fix(buzz): reconnect silently skips retained inbound messages (#116925)

* fix(buzz): page reconnect history instead of dropping older messages

On reconnect the Buzz gateway asked the relay for a single capped page of
room history and treated EOSE as a complete recovery. Any backlog larger
than the per-room replay limit was never requested again, so those messages
never reached handleBuzzInbound, the agent, or the session transcript, and
nothing reported the loss.

The room subscription now records how much history its first page returned
and, when that page filled the limit, pages older history through the relay
until it is exhausted. Each page is dispatched through the existing bounded
replay queue and the next page waits for queue capacity, so recovery stays
memory bounded. A room whose backlog cannot be paged past a single
timestamp is now reported instead of dropped silently.

Membership tracking moves to room-membership-tracker.ts because buzz-bus.ts
was already at the 700 line ceiling.

* fix(buzz): hold dispatch capacity while a history page is in flight

Catch-up asked the replay queue whether capacity existed and then issued the
relay query, but nothing held those slots. Live room events arriving while
the query was in flight could consume them, so admitting the page afterwards
could report overflow, which closes the dispatch queue and drops every
queued message before forcing a reconnect. A busy room with slow inbound
handlers could repeat that instead of finishing recovery.

Capacity is now reserved rather than sampled. reserveCapacity resolves with a
reservation that owns its slots until released, page events are admitted
through that reservation, and the reservation is released once the page is
enqueued. Live enqueues keep the full pending limit and are never rejected
earlier because a reservation is outstanding. An overflow reported through a
reservation means the relay exceeded the page it was asked for, so it is
reported as a history error instead of tearing down the session.

* fix(buzz): bound reconnect history pages

* fix(buzz): preserve replay reservations

* test(buzz): cover catch-up settlement paths

* fix(buzz): drain saturated history ranges

---------

Co-authored-by: Vincent Koc <vincentkoc@ieee.org>
This commit is contained in:
Yuval Dinodia
2026-07-31 20:28:03 -04:00
committed by GitHub
parent 4f7fbcb34d
commit e8d546c8da
7 changed files with 1366 additions and 358 deletions

View File

@@ -0,0 +1,447 @@
import { mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os";
import path from "node:path";
import { finalizeEvent, getPublicKey, type Event, type Filter } from "nostr-tools";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
const relayMocks = vi.hoisted(() => ({
connect: vi.fn<() => Promise<void>>(),
auth: vi.fn<() => Promise<string>>(),
publish: vi.fn<(event: Event) => Promise<string>>(),
send: vi.fn<(message: string) => Promise<void>>(),
close: vi.fn(),
connected: true,
storedEvents: [] as Event[],
historyRequests: [] as Filter[],
historySubscriptionCloses: 0,
closeHistoryPagesReason: undefined as string | undefined,
overReturnHistoryPages: false,
stallHistoryPages: false,
}));
function matchesRelayFilter(event: Event, filter: Filter): boolean {
if (filter.kinds && !filter.kinds.includes(event.kind)) {
return false;
}
if (filter.authors && !filter.authors.includes(event.pubkey)) {
return false;
}
for (const [key, values] of Object.entries(filter)) {
if (!key.startsWith("#") || !Array.isArray(values)) {
continue;
}
const tagName = key.slice(1);
const tagValues = event.tags.filter((tag) => tag[0] === tagName).map((tag) => tag[1] ?? "");
if (!tagValues.some((value) => (values as string[]).includes(value))) {
return false;
}
}
if (filter.since !== undefined && event.created_at < filter.since) {
return false;
}
if (filter.until !== undefined && event.created_at > filter.until) {
return false;
}
return true;
}
function selectRelayEvents(filter: Filter): Event[] {
const matched = relayMocks.storedEvents
.filter((event) => matchesRelayFilter(event, filter))
.toSorted((left, right) => right.created_at - left.created_at);
return filter.limit === undefined ||
(relayMocks.overReturnHistoryPages && filter.until !== undefined)
? matched
: matched.slice(0, filter.limit);
}
vi.mock("nostr-tools", async (importOriginal) => {
const actual = await importOriginal<typeof import("nostr-tools")>();
return {
...actual,
Relay: class {
onauth?: (template: unknown) => Promise<unknown>;
idleSince: number | undefined;
ongoingOperations = 0;
get connected() {
return relayMocks.connected;
}
connect = relayMocks.connect;
auth = relayMocks.auth;
publish = relayMocks.publish;
send = relayMocks.send;
close = relayMocks.close;
scheduleIdleClose = vi.fn();
prepareSubscription(
filters: Filter[],
handlers: {
onevent: (event: Event) => void;
oneose?: () => void;
onclose: (reason: string) => void;
},
) {
let isHistoryPage = false;
for (const filter of filters) {
if (filter.kinds?.includes(9)) {
relayMocks.historyRequests.push(filter);
isHistoryPage ||= filter.until !== undefined;
}
for (const event of selectRelayEvents(filter)) {
handlers.onevent(event);
}
if (relayMocks.closeHistoryPagesReason && isHistoryPage) {
queueMicrotask(() => {
handlers.onclose(relayMocks.closeHistoryPagesReason ?? "relay closed");
});
return {
id: `sub:${relayMocks.historyRequests.length}`,
close: vi.fn(),
closed: false,
};
}
if (relayMocks.stallHistoryPages && isHistoryPage) {
return {
id: `sub:${relayMocks.historyRequests.length}`,
close: vi.fn(),
closed: false,
};
}
}
handlers.oneose?.();
return {
id: `sub:${relayMocks.historyRequests.length}`,
close: vi.fn(() => {
if (isHistoryPage) {
relayMocks.historySubscriptionCloses += 1;
}
}),
closed: false,
};
}
},
};
});
import { startBuzzBus } from "./buzz-bus.js";
const BUZZ_NORMAL_MESSAGE_KIND = 9;
const BUZZ_ROOM_MEMBERSHIP_KIND = 39_002;
const PRIVATE_KEY = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
const SENDER_PRIVATE_KEY = "0102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f20";
const SENDER_SECRET = Uint8Array.from(Buffer.from(SENDER_PRIVATE_KEY, "hex"));
const ACCOUNT_ID = "default";
const CHANNEL_ID = "7c4a6d2a-2ed9-4b4e-a5e2-4d705ee9b34c";
const BOT_PUBLIC_KEY = getPublicKey(Uint8Array.from(Buffer.from(PRIVATE_KEY, "hex")));
const SENDER_PUBLIC_KEY = getPublicKey(SENDER_SECRET);
const RELAY_PUBLIC_KEY = "f".repeat(64);
const HISTORY_LIMIT = 100;
const BASE_TIMESTAMP = 1_700_000_000;
const tempDirs = new Set<string>();
let previousStateDir: string | undefined;
function buildMessageEvent(index: number, createdAt: number): Event {
return finalizeEvent(
{
kind: BUZZ_NORMAL_MESSAGE_KIND,
content: `offline-message-${String(index).padStart(3, "0")}`,
created_at: createdAt,
tags: [["h", CHANNEL_ID]],
},
SENDER_SECRET,
);
}
function seedOfflineBacklog(count: number, createdAt: (index: number) => number): void {
for (let index = 0; index < count; index += 1) {
relayMocks.storedEvents.push(buildMessageEvent(index, createdAt(index)));
}
}
async function waitForSettled(predicate: () => boolean): Promise<void> {
for (let attempt = 0; attempt < 300; attempt += 1) {
if (predicate()) {
return;
}
await new Promise((resolve) => {
setTimeout(resolve, 10);
});
}
}
describe("Buzz reconnect history catch-up", () => {
beforeEach(() => {
previousStateDir = process.env.OPENCLAW_STATE_DIR;
// openclaw-temp-dir: allow extension tests cannot import root test helpers.
const stateDir = mkdtempSync(path.join(tmpdir(), "openclaw-buzz-catchup-"));
tempDirs.add(stateDir);
process.env.OPENCLAW_STATE_DIR = stateDir;
vi.clearAllMocks();
relayMocks.historyRequests.length = 0;
relayMocks.historySubscriptionCloses = 0;
relayMocks.closeHistoryPagesReason = undefined;
relayMocks.overReturnHistoryPages = false;
relayMocks.stallHistoryPages = false;
relayMocks.storedEvents = [
{
id: "membership-1",
kind: BUZZ_ROOM_MEMBERSHIP_KIND,
pubkey: RELAY_PUBLIC_KEY,
created_at: BASE_TIMESTAMP - 3_600,
content: "",
sig: "e".repeat(128),
tags: [
["d", CHANNEL_ID],
["p", BOT_PUBLIC_KEY, "", "bot"],
["p", SENDER_PUBLIC_KEY, "", "member"],
],
},
];
relayMocks.connect.mockResolvedValue();
relayMocks.auth.mockResolvedValue("ok");
relayMocks.publish.mockResolvedValue("");
relayMocks.send.mockResolvedValue();
relayMocks.connected = true;
vi.stubGlobal(
"fetch",
vi.fn(async () => ({
ok: true,
json: async () => ({
self: RELAY_PUBLIC_KEY,
software: "https://github.com/block/buzz",
}),
})),
);
});
afterEach(() => {
if (previousStateDir === undefined) {
delete process.env.OPENCLAW_STATE_DIR;
} else {
process.env.OPENCLAW_STATE_DIR = previousStateDir;
}
for (const tempDir of tempDirs) {
rmSync(tempDir, { recursive: true, force: true });
}
tempDirs.clear();
vi.useRealTimers();
});
it("delivers backlog older than the per-room history limit", async () => {
seedOfflineBacklog(HISTORY_LIMIT + 1, (index) => BASE_TIMESTAMP + index);
const received: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async (message) => {
received.push(message.text);
},
});
await waitForSettled(() => received.length >= HISTORY_LIMIT + 1);
await bus.close();
expect(new Set(received).size).toBe(HISTORY_LIMIT + 1);
expect(received).toContain("offline-message-000");
expect(received.length).toBe(HISTORY_LIMIT + 1);
expect(relayMocks.historySubscriptionCloses).toBe(1);
});
it("pages a backlog spanning several history windows", async () => {
const backlogSize = 250;
seedOfflineBacklog(backlogSize, (index) => BASE_TIMESTAMP + index);
const received: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async (message) => {
received.push(message.text);
},
});
await waitForSettled(() => received.length >= backlogSize);
await bus.close();
expect(new Set(received).size).toBe(backlogSize);
expect(received).toContain("offline-message-000");
expect(relayMocks.historyRequests.length).toBeGreaterThan(1);
});
it("drains a backlog that exceeds one page at the same timestamp", async () => {
seedOfflineBacklog(HISTORY_LIMIT + 1, () => BASE_TIMESTAMP);
const historyErrors: string[] = [];
const received: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async (message) => {
received.push(message.text);
},
onHistoryError: (error) => {
historyErrors.push(error.message);
},
});
await waitForSettled(() => received.length >= HISTORY_LIMIT + 1);
await bus.close();
expect(historyErrors).toEqual([]);
expect(new Set(received).size).toBe(HISTORY_LIMIT + 1);
expect(received.length).toBe(HISTORY_LIMIT + 1);
expect(relayMocks.historyRequests.some((filter) => filter.limit === undefined)).toBe(true);
});
it("bounds a catch-up page when the relay ignores its history limit", async () => {
seedOfflineBacklog(250, (index) => BASE_TIMESTAMP + index);
relayMocks.overReturnHistoryPages = true;
const historyErrors: string[] = [];
const received: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async (message) => {
received.push(message.text);
},
onHistoryError: (error) => {
historyErrors.push(error.message);
},
});
await waitForSettled(() => received.length >= 250);
await bus.close();
expect(historyErrors).toEqual([]);
expect(new Set(received).size).toBe(250);
expect(received.length).toBe(250);
expect(relayMocks.historySubscriptionCloses).toBe(2);
});
it("bisects an overfull relay range until every bounded page fits", async () => {
const backlogSize = 1_300;
seedOfflineBacklog(backlogSize, (index) => BASE_TIMESTAMP + index);
relayMocks.overReturnHistoryPages = true;
const historyErrors: string[] = [];
const received: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async (message) => {
received.push(message.text);
},
onHistoryError: (error) => {
historyErrors.push(error.message);
},
});
await waitForSettled(() => received.length >= backlogSize);
await bus.close();
expect(historyErrors).toEqual([]);
expect(new Set(received).size).toBe(backlogSize);
expect(received.length).toBe(backlogSize);
expect(
relayMocks.historyRequests.filter((filter) => filter.limit === undefined).length,
).toBeGreaterThan(2);
});
it("fails the bus when a catch-up subscription never reaches EOSE", async () => {
vi.useFakeTimers();
seedOfflineBacklog(HISTORY_LIMIT + 1, (index) => BASE_TIMESTAMP + index);
relayMocks.stallHistoryPages = true;
const fatalErrors: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async () => {},
onFatalError: (error) => {
fatalErrors.push(error.message);
},
});
await vi.advanceTimersByTimeAsync(10_000);
await bus.close();
expect(fatalErrors).toEqual([`Timed out loading Buzz room history for ${CHANNEL_ID}`]);
expect(relayMocks.close).toHaveBeenCalled();
expect(relayMocks.historySubscriptionCloses).toBe(0);
});
it("fails the bus when a catch-up subscription closes unexpectedly", async () => {
seedOfflineBacklog(HISTORY_LIMIT + 1, (index) => BASE_TIMESTAMP + index);
relayMocks.closeHistoryPagesReason = "relay rejected subscription";
const fatalErrors: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async () => {},
onFatalError: (error) => {
fatalErrors.push(error.message);
},
});
await waitForSettled(() => fatalErrors.length > 0);
await bus.close();
expect(fatalErrors).toEqual([
`Buzz room history query closed for ${CHANNEL_ID}: relay rejected subscription`,
]);
expect(relayMocks.close).toHaveBeenCalled();
});
it("stops an active history query quietly when the bus closes", async () => {
seedOfflineBacklog(250, (index) => BASE_TIMESTAMP + index);
relayMocks.stallHistoryPages = true;
const fatalErrors: string[] = [];
const historyErrors: string[] = [];
const received: string[] = [];
const bus = await startBuzzBus({
accountId: ACCOUNT_ID,
relayUrl: "wss://buzz.example.com",
privateKey: PRIVATE_KEY,
channelIds: [CHANNEL_ID],
since: BASE_TIMESTAMP - 60,
onMessage: async (message) => {
received.push(message.text);
},
onFatalError: (error) => {
fatalErrors.push(error.message);
},
onHistoryError: (error) => {
historyErrors.push(error.message);
},
});
await waitForSettled(() => relayMocks.historyRequests.length > 1);
await bus.close();
const requestsAtClose = relayMocks.historyRequests.length;
await new Promise((resolve) => {
setTimeout(resolve, 100);
});
expect(relayMocks.historyRequests.length).toBe(requestsAtClose);
expect(received.length).toBeLessThan(250);
expect(fatalErrors).toEqual([]);
expect(historyErrors).toEqual([]);
});
});

View File

@@ -3,11 +3,9 @@ import { createChannelReplayGuard } from "openclaw/plugin-sdk/persistent-dedupe"
import { queryBuzzDirectoryRooms, startBuzzDirectoryRelay } from "./directory-relay.js";
import { BuzzDirectoryState } from "./directory-state.js";
import {
BUZZ_INBOUND_MESSAGE_KINDS,
BUZZ_NORMAL_MESSAGE_KIND,
BUZZ_TYPING_INDICATOR_KIND,
buildBuzzMessageTags,
isBuzzInboundMessageKind,
parseBuzzMessageEvent,
type BuzzInboundMessage,
} from "./message-event.js";
@@ -17,20 +15,13 @@ import {
connectAuthenticatedBuzzRelaySession,
parseBuzzAuthTag,
} from "./relay-auth.js";
import { openBuzzRelaySubscription } from "./relay-subscription.js";
import {
BUZZ_REPLAY_DISPATCH_MAX_PENDING,
createBuzzReplayDispatchQueue,
resolveBuzzRoomHistoryLimit,
} from "./replay-dispatch.js";
import { startBuzzRoomMembershipNotifications } from "./room-membership-notification.js";
import { queryBuzzRoomMemberships } from "./room-membership-query.js";
import {
BUZZ_ROOM_SYSTEM_KIND,
isNewerBuzzRoomMembership,
parseBuzzRoomMembershipChangeEvent,
type BuzzRoomMembership,
} from "./room-membership.js";
import { createBuzzRoomMembershipTracker } from "./room-membership-tracker.js";
import { resolveBuzzSubscriptionBudget } from "./subscription-budget.js";
import { decodeBuzzPrivateKey, resolveBuzzPublicKey } from "./types.js";
@@ -40,11 +31,6 @@ const REPLAY_TTL_MS = 30 * 24 * 60 * 60 * 1000;
const REPLAY_MAX_ENTRIES = 10_000;
const REPLAY_STATE_MAX_ENTRIES = 50_000;
const REPLAY_NAMESPACE_PREFIX = "buzz.inbound-dedupe";
const MEMBERSHIP_READY_TIMEOUT_MS = 10_000;
const MEMBERSHIP_TRACKER_SETUP_CLOSE_REASON = "membership tracker setup failed";
const BUZZ_ROOM_METADATA_EDIT_KIND = 9_002;
const MEMBERSHIP_REFRESH_DELAYS_MS = [100, 500, 1_500, 3_000] as const;
const MEMBERSHIP_EVENT_CACHE_MAX_ENTRIES = 10_000;
export interface BuzzBus {
publicKey: string;
@@ -154,332 +140,6 @@ function startBuzzPresenceHeartbeat(params: {
};
}
async function sleepWithSignal(delayMs: number, signal?: AbortSignal): Promise<void> {
signal?.throwIfAborted();
await new Promise<void>((resolve, reject) => {
let settled = false;
const finish = (error?: unknown) => {
if (settled) {
return;
}
settled = true;
clearTimeout(timer);
signal?.removeEventListener("abort", onAbort);
if (error === undefined) {
resolve();
} else {
reject(
error instanceof Error
? error
: new Error("Buzz room membership refresh failed", { cause: error }),
);
}
};
const onAbort = () =>
finish(signal?.reason ?? new Error("Buzz room membership refresh aborted"));
const timer = setTimeout(() => finish(), delayMs);
signal?.addEventListener("abort", onAbort, { once: true });
if (signal?.aborted) {
onAbort();
}
});
}
async function createBuzzRoomMembershipTracker(params: {
relay: Relay;
relayPublicKey: string;
channelIds: string[];
botPublicKey: string;
since: number;
messageSince: number;
messageLimit: number;
onMessageEvent: (
event: Event,
isMember: (channelId: string, publicKey: string) => boolean,
) => void;
onFatalError?: (error: Error) => void;
onMembershipsChanged?: (memberships: ReadonlyMap<string, BuzzRoomMembership>) => void;
onRoomMetadataChanged?: (channelId: string) => void;
signal?: AbortSignal;
}): Promise<{
memberships: () => ReadonlyMap<string, BuzzRoomMembership>;
}> {
type ExpectedMembership = "present" | "absent";
type RefreshState = {
generation: number;
lastAttemptedGeneration: number;
promise: Promise<void>;
};
const historicalRooms = new Set<string>();
const seenEventIds = new Map<string, true>();
const blockedRooms = new Set<string>();
const deniedMembers = new Map<string, Set<string>>();
const pendingMemberships = new Map<string, Map<string, ExpectedMembership>>();
const refreshes = new Map<string, RefreshState>();
let membershipQueryTail = Promise.resolve();
const memberships = await queryBuzzRoomMemberships(params);
const isMember = (channelId: string, publicKey: string) =>
!blockedRooms.has(channelId) &&
!deniedMembers.get(channelId)?.has(publicKey.trim().toLowerCase()) &&
memberships.get(channelId)?.members.has(publicKey.trim().toLowerCase()) === true;
const markSystemEventSeen = (eventId: string): boolean => {
if (seenEventIds.has(eventId)) {
return false;
}
seenEventIds.set(eventId, true);
if (seenEventIds.size > MEMBERSHIP_EVENT_CACHE_MAX_ENTRIES) {
const oldestEventId = seenEventIds.keys().next().value;
if (oldestEventId) {
seenEventIds.delete(oldestEventId);
}
}
return true;
};
const reportSystemEventError = (error: unknown) => {
if (params.signal?.aborted) {
return;
}
params.onFatalError?.(error instanceof Error ? error : new Error(String(error)));
params.relay.close();
};
const queryMembership = (channelId: string): Promise<BuzzRoomMembership | undefined> => {
const query = membershipQueryTail.then(async () =>
(
await queryBuzzRoomMemberships({
relay: params.relay,
relayPublicKey: params.relayPublicKey,
channelIds: [channelId],
signal: params.signal,
})
).get(channelId),
);
membershipQueryTail = query.then(
() => undefined,
() => undefined,
);
return query;
};
const refreshMembership = async (channelId: string, state: RefreshState): Promise<void> => {
const baseline = memberships.get(channelId);
if (!baseline) {
throw new Error(`Missing Buzz room membership for ${channelId}`);
}
for (const delayMs of MEMBERSHIP_REFRESH_DELAYS_MS) {
const generation = state.generation;
state.lastAttemptedGeneration = generation;
await sleepWithSignal(delayMs, params.signal);
if (state.generation !== generation) {
continue;
}
let refreshed: BuzzRoomMembership | undefined;
try {
refreshed = await queryMembership(channelId);
} catch (error) {
if (params.signal?.aborted) {
throw error;
}
continue;
}
if (state.generation !== generation || !refreshed) {
continue;
}
const pending = pendingMemberships.get(channelId);
const pendingMatches =
!pending ||
[...pending].every(
([publicKey, expected]) => refreshed.members.has(publicKey) === (expected === "present"),
);
const botMembershipChanged = pending?.has(params.botPublicKey) === true;
if (
!pendingMatches ||
(botMembershipChanged && !isNewerBuzzRoomMembership(refreshed, baseline))
) {
continue;
}
if (
refreshed.roles.get(params.botPublicKey) !== "bot" ||
!refreshed.members.has(params.botPublicKey)
) {
blockedRooms.add(channelId);
throw new Error(`Buzz bot no longer has the Bot role in room ${channelId}`);
}
memberships.set(channelId, refreshed);
pendingMemberships.delete(channelId);
deniedMembers.delete(channelId);
blockedRooms.delete(channelId);
params.onMembershipsChanged?.(memberships);
return;
}
if (state.generation !== state.lastAttemptedGeneration) {
return;
}
blockedRooms.add(channelId);
throw new Error(`Could not refresh Buzz room membership for ${channelId}`);
};
const refreshMembershipOnce = (channelId: string): Promise<void> => {
const current = refreshes.get(channelId);
if (current) {
current.generation += 1;
return current.promise;
}
const state = {
generation: 1,
lastAttemptedGeneration: 0,
promise: Promise.resolve(),
} satisfies RefreshState;
state.promise = refreshMembership(channelId, state).finally(() => {
if (refreshes.get(channelId) === state) {
refreshes.delete(channelId);
}
if (
state.generation !== state.lastAttemptedGeneration &&
pendingMemberships.has(channelId) &&
!params.signal?.aborted
) {
void refreshMembershipOnce(channelId).catch(reportSystemEventError);
}
});
refreshes.set(channelId, state);
return state.promise;
};
const handleSystemEvent = (event: Event): Promise<void> | undefined => {
if (!markSystemEventSeen(event.id)) {
return undefined;
}
const channelId = event.tags
.find((tag) => tag[0] === "h")?.[1]
?.trim()
.toLowerCase();
if (!channelId) {
return undefined;
}
if (event.kind === BUZZ_ROOM_METADATA_EDIT_KIND) {
params.onRoomMetadataChanged?.(channelId);
return undefined;
}
const membership = memberships.get(channelId);
if (!membership) {
return undefined;
}
const change = parseBuzzRoomMembershipChangeEvent(event, membership);
if (!change) {
return undefined;
}
// System events invalidate membership; the relay-signed roster decides the
// final state. Removals deny immediately, while joins wait for confirmation.
const expected = change.type === "member_joined" ? "present" : "absent";
const pending = pendingMemberships.get(channelId) ?? new Map<string, ExpectedMembership>();
pending.set(change.targetPublicKey, expected);
pendingMemberships.set(channelId, pending);
if (expected === "absent") {
const denied = deniedMembers.get(channelId) ?? new Set<string>();
denied.add(change.targetPublicKey);
deniedMembers.set(channelId, denied);
}
if (change.targetPublicKey === params.botPublicKey) {
blockedRooms.add(channelId);
}
return refreshMembershipOnce(channelId);
};
const handleRoomEvent = (event: Event) => {
if (isBuzzInboundMessageKind(event.kind)) {
params.onMessageEvent(event, isMember);
return;
}
void handleSystemEvent(event)?.catch(reportSystemEventError);
};
for (const channelId of params.channelIds) {
if (memberships.get(channelId)?.roles.get(params.botPublicKey) !== "bot") {
throw new Error(`Buzz bot does not have the Bot role in configured room ${channelId}`);
}
}
let resolveHistorical: (() => void) | undefined;
let rejectHistorical: ((error: Error) => void) | undefined;
const historicalReady = new Promise<void>((resolve, reject) => {
resolveHistorical = resolve;
rejectHistorical = reject;
});
const historicalTimeout = setTimeout(() => {
const error = new Error("Timed out loading Buzz room membership changes");
rejectHistorical?.(error);
params.relay.close();
}, MEMBERSHIP_READY_TIMEOUT_MS);
const subscriptions: Array<ReturnType<Relay["prepareSubscription"]>> = [];
try {
// Snapshot membership before room history so startup memory stays bounded.
// Buzz emits these filters in order: system changes since session start
// update or deny membership before the following message history is handled.
for (const channelId of params.channelIds) {
subscriptions.push(
openBuzzRelaySubscription(
params.relay,
[
{
kinds: [BUZZ_ROOM_SYSTEM_KIND, BUZZ_ROOM_METADATA_EDIT_KIND],
"#h": [channelId],
since: params.since,
},
{
kinds: [...BUZZ_INBOUND_MESSAGE_KINDS],
"#h": [channelId],
since: params.messageSince,
limit: params.messageLimit,
},
],
{
onevent: handleRoomEvent,
oneose: () => {
historicalRooms.add(channelId);
if (historicalRooms.size === params.channelIds.length) {
resolveHistorical?.();
}
},
onclose: (reason) => {
if (!historicalRooms.has(channelId)) {
rejectHistorical?.(
new Error(`Buzz membership subscription closed for ${channelId}: ${reason}`),
);
} else if (
reason !== "shutdown" &&
reason !== "relay connection closed by us" &&
reason !== MEMBERSHIP_TRACKER_SETUP_CLOSE_REASON &&
!params.signal?.aborted
) {
params.onFatalError?.(
new Error(`Buzz membership subscription closed for ${channelId}: ${reason}`),
);
}
},
},
),
);
}
await historicalReady;
} catch (error) {
if (params.relay.connected) {
for (const subscription of subscriptions) {
if (!subscription.closed) {
subscription.close(MEMBERSHIP_TRACKER_SETUP_CLOSE_REASON);
}
}
}
throw error;
} finally {
clearTimeout(historicalTimeout);
}
return {
memberships: () => memberships,
};
}
export async function sendBuzzTextOneShot(params: {
relayUrl: string;
privateKey: string;
@@ -515,6 +175,7 @@ export async function startBuzzBus(options: {
onMessageError?: (error: Error) => void;
onFatalError?: (error: Error) => void;
onDedupeError?: (error: Error) => void;
onHistoryError?: (error: Error) => void;
onPresenceError?: (error: Error) => void;
profileName?: string;
onProfilePublished?: (eventId: string) => void;
@@ -642,7 +303,9 @@ export async function startBuzzBus(options: {
since: sessionStartedAt,
messageSince: options.since ?? sessionStartedAt,
messageLimit: resolveBuzzRoomHistoryLimit(activeChannelIds.length),
onMessageEvent: (event, isMember) => {
reserveDispatchCapacity: (slots) => dispatchQueue.reserveCapacity(slots),
onHistoryError: options.onHistoryError,
onMessageEvent: (event, isMember, reservation) => {
if (signal.aborted || event.pubkey === publicKey) {
return;
}
@@ -652,19 +315,28 @@ export async function startBuzzBus(options: {
}
// Admit only room members to bounded workers; claim replay dedupe inside
// each worker so queued history cannot create unbounded in-flight state.
const admission = dispatchQueue.enqueue(async () => {
const admission = (reservation ?? dispatchQueue).enqueue(async () => {
await replayGuard.processGuarded(event, async () => {
await options.onMessage(message, bus, signal);
});
});
if (admission === "overflow") {
void dispatchQueue.close();
reportFatalError(
if (admission !== "overflow") {
return;
}
if (reservation) {
options.onHistoryError?.(
new Error(
`Buzz inbound replay exceeded the ${BUZZ_REPLAY_DISPATCH_MAX_PENDING}-message pending limit`,
`Buzz room ${message.channelId} returned more history than the ${BUZZ_REPLAY_DISPATCH_MAX_PENDING}-message pending limit allows`,
),
);
return;
}
void dispatchQueue.close();
reportFatalError(
new Error(
`Buzz inbound replay exceeded the ${BUZZ_REPLAY_DISPATCH_MAX_PENDING}-message pending limit`,
),
);
},
onFatalError: reportFatalError,
onMembershipsChanged: (memberships) => {
@@ -688,6 +360,7 @@ export async function startBuzzBus(options: {
: undefined;
directory.replaceMemberships(membershipTracker?.memberships() ?? new Map());
directoryRelay.replaceProfilePublicKeys(directory.profilePublicKeys());
void membershipTracker?.catchUpHistory();
stopPresenceHeartbeat = startBuzzPresenceHeartbeat({
relay,
secretKey,

View File

@@ -114,6 +114,11 @@ export async function startBuzzGatewayAccount(ctx: ChannelGatewayContext<Resolve
onDedupeError: (error) => {
ctx.log?.error?.(`[${account.accountId}] Buzz replay state failed: ${error.message}`);
},
onHistoryError: (error) => {
ctx.log?.warn?.(
`[${account.accountId}] Buzz history recovery incomplete: ${error.message}`,
);
},
onPresenceError: (error) => {
ctx.log?.warn?.(
`[${account.accountId}] Buzz presence heartbeat failed: ${error.message}`,

View File

@@ -0,0 +1,265 @@
import type { Event, Relay } from "nostr-tools";
import { BUZZ_INBOUND_MESSAGE_KINDS } from "./message-event.js";
import { openBuzzRelaySubscription } from "./relay-subscription.js";
import {
BUZZ_REPLAY_DISPATCH_MAX_PENDING,
type BuzzReplayDispatchReservation,
} from "./replay-dispatch.js";
const HISTORY_PAGE_TIMEOUT_MS = 10_000;
const HISTORY_PAGE_COMPLETE_REASON = "buzz room history page loaded";
type BuzzRoomHistoryCatchUp = "complete" | "aborted" | "timestamp-over-limit";
type BuzzRoomHistoryPage = {
events: Event[];
overLimit: boolean;
};
async function queryBuzzRoomHistoryPage(params: {
relay: Relay;
channelId: string;
since: number;
until: number;
requestLimit?: number;
maxEvents: number;
skipEventIds?: ReadonlySet<string>;
signal?: AbortSignal;
}): Promise<BuzzRoomHistoryPage> {
const events: Event[] = [];
let overLimit = false;
return await new Promise<BuzzRoomHistoryPage>((resolve, reject) => {
let settled = false;
let receivedEose = false;
const timeout = setTimeout(() => {
const error = new Error(`Timed out loading Buzz room history for ${params.channelId}`);
finish(error);
params.relay.close();
}, HISTORY_PAGE_TIMEOUT_MS);
const subscriptionRef: { current?: ReturnType<Relay["prepareSubscription"]> } = {};
const finish = (error?: unknown) => {
if (settled) {
return;
}
settled = true;
clearTimeout(timeout);
params.signal?.removeEventListener("abort", onAbort);
if (receivedEose) {
subscriptionRef.current?.close(HISTORY_PAGE_COMPLETE_REASON);
}
if (error === undefined) {
resolve({ events, overLimit });
} else {
reject(
error instanceof Error
? error
: new Error("Buzz room history query failed", { cause: error }),
);
}
};
const onAbort = () =>
finish(params.signal?.reason ?? new Error("Buzz room history query aborted"));
params.signal?.addEventListener("abort", onAbort, { once: true });
try {
subscriptionRef.current = openBuzzRelaySubscription(
params.relay,
[
{
kinds: [...BUZZ_INBOUND_MESSAGE_KINDS],
"#h": [params.channelId],
since: params.since,
until: params.until,
...(params.requestLimit === undefined ? {} : { limit: params.requestLimit }),
},
],
{
onevent: (event) => {
if (params.skipEventIds?.has(event.id)) {
return;
}
if (events.length < params.maxEvents) {
events.push(event);
} else {
overLimit = true;
}
},
oneose: () => {
receivedEose = true;
if (settled) {
subscriptionRef.current?.close(HISTORY_PAGE_COMPLETE_REASON);
} else {
finish();
}
},
onclose: (reason) => {
if (reason !== HISTORY_PAGE_COMPLETE_REASON) {
finish(
new Error(`Buzz room history query closed for ${params.channelId}: ${reason}`),
);
}
},
},
);
} catch (error) {
finish(error);
return;
}
if (settled && receivedEose) {
subscriptionRef.current.close(HISTORY_PAGE_COMPLETE_REASON);
}
if (params.signal?.aborted) {
onAbort();
}
});
}
async function drainBuzzRoomHistoryRange(params: {
relay: Relay;
channelId: string;
since: number;
until: number;
skipEventIds: ReadonlySet<string>;
reserveCapacity: (slots: number) => Promise<BuzzReplayDispatchReservation | undefined>;
onEvent: (event: Event, reservation: BuzzReplayDispatchReservation) => void;
signal?: AbortSignal;
}): Promise<BuzzRoomHistoryCatchUp> {
if (params.signal?.aborted) {
return "aborted";
}
const page = await queryBuzzRoomHistoryPage({
relay: params.relay,
channelId: params.channelId,
since: params.since,
until: params.until,
maxEvents: BUZZ_REPLAY_DISPATCH_MAX_PENDING,
skipEventIds: params.skipEventIds,
signal: params.signal,
});
if (!page.overLimit) {
if (page.events.length === 0) {
return "complete";
}
const reservation = await params.reserveCapacity(page.events.length);
if (!reservation) {
return "aborted";
}
try {
for (const event of page.events) {
params.onEvent(event, reservation);
}
} finally {
reservation.release();
}
return "complete";
}
if (params.since === params.until) {
const reservation = await params.reserveCapacity(page.events.length);
if (!reservation) {
return "aborted";
}
try {
for (const event of page.events) {
params.onEvent(event, reservation);
}
} finally {
reservation.release();
}
return "timestamp-over-limit";
}
// NIP-01 has only a second-resolution time cursor. Split an overfull range
// until every query fits; only a single overfull second is irreducible.
const midpoint = Math.floor((params.since + params.until) / 2);
const newer = await drainBuzzRoomHistoryRange({
...params,
since: midpoint + 1,
});
if (newer !== "complete") {
return newer;
}
return await drainBuzzRoomHistoryRange({
...params,
until: midpoint,
});
}
export async function catchUpBuzzRoomHistory(params: {
relay: Relay;
channelId: string;
since: number;
until: number;
limit: number;
reserveCapacity: (slots: number) => Promise<BuzzReplayDispatchReservation | undefined>;
onEvent: (event: Event, reservation: BuzzReplayDispatchReservation) => void;
signal?: AbortSignal;
}): Promise<BuzzRoomHistoryCatchUp> {
let until = params.until;
while (!params.signal?.aborted) {
const reservation = await params.reserveCapacity(params.limit);
if (!reservation) {
return "aborted";
}
let page: BuzzRoomHistoryPage;
try {
page = await queryBuzzRoomHistoryPage({
relay: params.relay,
channelId: params.channelId,
since: params.since,
until,
requestLimit: params.limit,
maxEvents: params.limit,
signal: params.signal,
});
if (page.events.length === 0) {
return "complete";
}
for (const event of page.events) {
params.onEvent(event, reservation);
}
} finally {
reservation.release();
}
let oldest = until;
for (const event of page.events) {
oldest = Math.min(oldest, event.created_at);
}
const skipEventIds = new Set(page.events.map((event) => event.id));
if (page.overLimit) {
return await drainBuzzRoomHistoryRange({
relay: params.relay,
channelId: params.channelId,
since: params.since,
until,
skipEventIds,
reserveCapacity: params.reserveCapacity,
onEvent: params.onEvent,
signal: params.signal,
});
}
if (page.events.length < params.limit) {
return "complete";
}
if (oldest >= until) {
const outcome = await drainBuzzRoomHistoryRange({
relay: params.relay,
channelId: params.channelId,
since: until,
until,
skipEventIds,
reserveCapacity: params.reserveCapacity,
onEvent: params.onEvent,
signal: params.signal,
});
if (outcome !== "complete") {
return outcome;
}
if (until <= params.since) {
return "complete";
}
until -= 1;
continue;
}
until = oldest;
}
return "aborted";
}

View File

@@ -0,0 +1,139 @@
import { describe, expect, it } from "vitest";
import {
BUZZ_REPLAY_DISPATCH_MAX_PENDING,
createBuzzReplayDispatchQueue,
} from "./replay-dispatch.js";
const REPLAY_DISPATCH_CONCURRENCY = 8;
function createBlockedQueue() {
const releases: Array<() => void> = [];
const queue = createBuzzReplayDispatchQueue({ onTaskError: () => {} });
const blockTask = () => {
let release = () => {};
const gate = new Promise<void>((resolve) => {
release = resolve;
});
releases.push(release);
return async () => {
await gate;
};
};
return { queue, releases, blockTask };
}
async function flush(): Promise<void> {
for (let index = 0; index < 5; index += 1) {
await Promise.resolve();
}
}
describe("Buzz replay dispatch capacity reservations", () => {
it("withholds a reservation while queued work occupies the pending limit", async () => {
const { queue, blockTask } = createBlockedQueue();
for (
let index = 0;
index < BUZZ_REPLAY_DISPATCH_MAX_PENDING + REPLAY_DISPATCH_CONCURRENCY;
index += 1
) {
expect(queue.enqueue(blockTask())).toBe("accepted");
}
let granted: unknown = "pending";
void queue.reserveCapacity(10).then((reservation) => {
granted = reservation;
});
await flush();
expect(granted).toBe("pending");
expect(queue.enqueue(blockTask())).toBe("overflow");
});
it("withholds reserved slots from later live events", async () => {
const { queue, releases, blockTask } = createBlockedQueue();
const pageSize = 100;
for (
let index = 0;
index < BUZZ_REPLAY_DISPATCH_MAX_PENDING + REPLAY_DISPATCH_CONCURRENCY;
index += 1
) {
queue.enqueue(blockTask());
}
const reservationPromise = queue.reserveCapacity(pageSize);
for (let index = 0; index < pageSize; index += 1) {
releases[index]?.();
}
const reservation = await reservationPromise;
expect(reservation).toBeDefined();
for (let index = 0; index < pageSize; index += 1) {
expect(queue.enqueue(blockTask())).toBe("overflow");
}
const admissions = new Set<string>();
for (let index = 0; index < pageSize; index += 1) {
admissions.add(reservation?.enqueue(blockTask()) ?? "missing");
}
expect([...admissions]).toEqual(["accepted"]);
expect(reservation?.enqueue(blockTask())).toBe("overflow");
});
it("returns unused slots when a reservation is released", async () => {
const { queue, releases, blockTask } = createBlockedQueue();
for (
let index = 0;
index < BUZZ_REPLAY_DISPATCH_MAX_PENDING + REPLAY_DISPATCH_CONCURRENCY;
index += 1
) {
queue.enqueue(blockTask());
}
const firstPromise = queue.reserveCapacity(50);
for (let index = 0; index < 50; index += 1) {
releases[index]?.();
}
const first = await firstPromise;
expect(first).toBeDefined();
let second: unknown = "pending";
void queue.reserveCapacity(50).then((reservation) => {
second = reservation;
});
await flush();
expect(second).toBe("pending");
first?.release();
await flush();
expect(second).toBeDefined();
expect(second).not.toBe("pending");
});
it("abandons waiting reservations once the queue closes", async () => {
const { queue, blockTask } = createBlockedQueue();
for (
let index = 0;
index < BUZZ_REPLAY_DISPATCH_MAX_PENDING + REPLAY_DISPATCH_CONCURRENCY;
index += 1
) {
queue.enqueue(blockTask());
}
const reservationPromise = queue.reserveCapacity(10);
void queue.close();
expect(await reservationPromise).toBeUndefined();
expect(await queue.reserveCapacity(1)).toBeUndefined();
});
it("rejects work from a held reservation after the queue closes", async () => {
const { queue, blockTask } = createBlockedQueue();
const reservation = await queue.reserveCapacity(2);
await queue.close();
expect(reservation?.enqueue(blockTask())).toBe("closed");
reservation?.release();
expect(await queue.reserveCapacity(1)).toBeUndefined();
});
});

View File

@@ -2,8 +2,16 @@ const REPLAY_DISPATCH_CONCURRENCY = 8;
export const BUZZ_REPLAY_DISPATCH_MAX_PENDING = 1_024;
const REPLAY_HISTORY_MAX_PER_ROOM = 100;
type BuzzReplayDispatchAdmission = "accepted" | "closed" | "overflow";
export type BuzzReplayDispatchReservation = {
enqueue: (task: () => Promise<void>) => BuzzReplayDispatchAdmission;
release: () => void;
};
type BuzzReplayDispatchQueue = {
enqueue: (task: () => Promise<void>) => "accepted" | "closed" | "overflow";
enqueue: (task: () => Promise<void>) => BuzzReplayDispatchAdmission;
reserveCapacity: (slots: number) => Promise<BuzzReplayDispatchReservation | undefined>;
close: () => Promise<void>;
};
@@ -26,6 +34,14 @@ export function createBuzzReplayDispatchQueue(params: {
}
};
let reserved = 0;
const reservationWaiters: Array<{
slots: number;
resolve: (reservation: BuzzReplayDispatchReservation | undefined) => void;
}> = [];
const availableCapacity = () =>
BUZZ_REPLAY_DISPATCH_MAX_PENDING - (pending.length - pendingHead) - reserved;
const compactPending = () => {
if (pendingHead > 256 && pendingHead * 2 >= pending.length) {
pending.splice(0, pendingHead);
@@ -54,28 +70,88 @@ export function createBuzzReplayDispatchQueue(params: {
drain();
});
}
settleReservationWaiters();
};
return {
enqueue(task) {
if (closed) {
return "closed";
}
if (active < REPLAY_DISPATCH_CONCURRENCY) {
const enqueueTask = (task: () => Promise<void>): BuzzReplayDispatchAdmission => {
if (closed) {
return "closed";
}
if (active < REPLAY_DISPATCH_CONCURRENCY) {
pending.push(task);
drain();
return "accepted";
}
if (availableCapacity() <= 0) {
return "overflow";
}
pending.push(task);
return "accepted";
};
const createReservation = (slots: number): BuzzReplayDispatchReservation => {
let remaining = slots;
reserved += slots;
return {
enqueue(task) {
if (closed) {
return "closed";
}
if (remaining === 0) {
return "overflow";
}
remaining -= 1;
reserved -= 1;
pending.push(task);
drain();
return "accepted";
},
release() {
reserved -= remaining;
remaining = 0;
settleReservationWaiters();
},
};
};
const settleReservationWaiters = () => {
while (reservationWaiters.length > 0) {
const waiter = reservationWaiters[0];
if (!waiter) {
reservationWaiters.shift();
continue;
}
if (pending.length - pendingHead >= BUZZ_REPLAY_DISPATCH_MAX_PENDING) {
return "overflow";
if (closed) {
reservationWaiters.shift();
waiter.resolve(undefined);
continue;
}
pending.push(task);
return "accepted";
if (availableCapacity() < waiter.slots) {
return;
}
reservationWaiters.shift();
waiter.resolve(createReservation(waiter.slots));
}
};
return {
enqueue: enqueueTask,
async reserveCapacity(slots) {
if (closed) {
return undefined;
}
if (reservationWaiters.length === 0 && availableCapacity() >= slots) {
return createReservation(slots);
}
return await new Promise<BuzzReplayDispatchReservation | undefined>((resolve) => {
reservationWaiters.push({ slots, resolve });
});
},
async close() {
closed = true;
pending.length = 0;
pendingHead = 0;
settleReservationWaiters();
settleDrained();
await drained;
},

View File

@@ -0,0 +1,403 @@
import type { Event, Relay } from "nostr-tools";
import { catchUpBuzzRoomHistory } from "./history-catchup.js";
import { BUZZ_INBOUND_MESSAGE_KINDS, isBuzzInboundMessageKind } from "./message-event.js";
import { openBuzzRelaySubscription } from "./relay-subscription.js";
import {
BUZZ_REPLAY_DISPATCH_MAX_PENDING,
type BuzzReplayDispatchReservation,
} from "./replay-dispatch.js";
import { queryBuzzRoomMemberships } from "./room-membership-query.js";
import {
BUZZ_ROOM_SYSTEM_KIND,
isNewerBuzzRoomMembership,
parseBuzzRoomMembershipChangeEvent,
type BuzzRoomMembership,
} from "./room-membership.js";
const MEMBERSHIP_READY_TIMEOUT_MS = 10_000;
const MEMBERSHIP_TRACKER_SETUP_CLOSE_REASON = "membership tracker setup failed";
const BUZZ_ROOM_METADATA_EDIT_KIND = 9_002;
const MEMBERSHIP_REFRESH_DELAYS_MS = [100, 500, 1_500, 3_000] as const;
const MEMBERSHIP_EVENT_CACHE_MAX_ENTRIES = 10_000;
async function sleepWithSignal(delayMs: number, signal?: AbortSignal): Promise<void> {
signal?.throwIfAborted();
await new Promise<void>((resolve, reject) => {
let settled = false;
const finish = (error?: unknown) => {
if (settled) {
return;
}
settled = true;
clearTimeout(timer);
signal?.removeEventListener("abort", onAbort);
if (error === undefined) {
resolve();
} else {
reject(
error instanceof Error
? error
: new Error("Buzz room membership refresh failed", { cause: error }),
);
}
};
const onAbort = () =>
finish(signal?.reason ?? new Error("Buzz room membership refresh aborted"));
const timer = setTimeout(() => finish(), delayMs);
signal?.addEventListener("abort", onAbort, { once: true });
if (signal?.aborted) {
onAbort();
}
});
}
export async function createBuzzRoomMembershipTracker(params: {
relay: Relay;
relayPublicKey: string;
channelIds: string[];
botPublicKey: string;
since: number;
messageSince: number;
messageLimit: number;
reserveDispatchCapacity: (slots: number) => Promise<BuzzReplayDispatchReservation | undefined>;
onMessageEvent: (
event: Event,
isMember: (channelId: string, publicKey: string) => boolean,
reservation?: BuzzReplayDispatchReservation,
) => void;
onFatalError?: (error: Error) => void;
onHistoryError?: (error: Error) => void;
onMembershipsChanged?: (memberships: ReadonlyMap<string, BuzzRoomMembership>) => void;
onRoomMetadataChanged?: (channelId: string) => void;
signal?: AbortSignal;
}): Promise<{
memberships: () => ReadonlyMap<string, BuzzRoomMembership>;
catchUpHistory: () => Promise<void>;
}> {
type ExpectedMembership = "present" | "absent";
type RefreshState = {
generation: number;
lastAttemptedGeneration: number;
promise: Promise<void>;
};
const historicalRooms = new Set<string>();
const historyPages = new Map<string, { count: number; oldest: number }>();
const seenEventIds = new Map<string, true>();
const blockedRooms = new Set<string>();
const deniedMembers = new Map<string, Set<string>>();
const pendingMemberships = new Map<string, Map<string, ExpectedMembership>>();
const refreshes = new Map<string, RefreshState>();
let membershipQueryTail = Promise.resolve();
const memberships = await queryBuzzRoomMemberships(params);
const isMember = (channelId: string, publicKey: string) =>
!blockedRooms.has(channelId) &&
!deniedMembers.get(channelId)?.has(publicKey.trim().toLowerCase()) &&
memberships.get(channelId)?.members.has(publicKey.trim().toLowerCase()) === true;
const markSystemEventSeen = (eventId: string): boolean => {
if (seenEventIds.has(eventId)) {
return false;
}
seenEventIds.set(eventId, true);
if (seenEventIds.size > MEMBERSHIP_EVENT_CACHE_MAX_ENTRIES) {
const oldestEventId = seenEventIds.keys().next().value;
if (oldestEventId) {
seenEventIds.delete(oldestEventId);
}
}
return true;
};
const reportSystemEventError = (error: unknown) => {
if (params.signal?.aborted) {
return;
}
params.onFatalError?.(error instanceof Error ? error : new Error(String(error)));
params.relay.close();
};
const queryMembership = (channelId: string): Promise<BuzzRoomMembership | undefined> => {
const query = membershipQueryTail.then(async () =>
(
await queryBuzzRoomMemberships({
relay: params.relay,
relayPublicKey: params.relayPublicKey,
channelIds: [channelId],
signal: params.signal,
})
).get(channelId),
);
membershipQueryTail = query.then(
() => undefined,
() => undefined,
);
return query;
};
const refreshMembership = async (channelId: string, state: RefreshState): Promise<void> => {
const baseline = memberships.get(channelId);
if (!baseline) {
throw new Error(`Missing Buzz room membership for ${channelId}`);
}
for (const delayMs of MEMBERSHIP_REFRESH_DELAYS_MS) {
const generation = state.generation;
state.lastAttemptedGeneration = generation;
await sleepWithSignal(delayMs, params.signal);
if (state.generation !== generation) {
continue;
}
let refreshed: BuzzRoomMembership | undefined;
try {
refreshed = await queryMembership(channelId);
} catch (error) {
if (params.signal?.aborted) {
throw error;
}
continue;
}
if (state.generation !== generation || !refreshed) {
continue;
}
const pending = pendingMemberships.get(channelId);
const pendingMatches =
!pending ||
[...pending].every(
([publicKey, expected]) => refreshed.members.has(publicKey) === (expected === "present"),
);
const botMembershipChanged = pending?.has(params.botPublicKey) === true;
if (
!pendingMatches ||
(botMembershipChanged && !isNewerBuzzRoomMembership(refreshed, baseline))
) {
continue;
}
if (
refreshed.roles.get(params.botPublicKey) !== "bot" ||
!refreshed.members.has(params.botPublicKey)
) {
blockedRooms.add(channelId);
throw new Error(`Buzz bot no longer has the Bot role in room ${channelId}`);
}
memberships.set(channelId, refreshed);
pendingMemberships.delete(channelId);
deniedMembers.delete(channelId);
blockedRooms.delete(channelId);
params.onMembershipsChanged?.(memberships);
return;
}
if (state.generation !== state.lastAttemptedGeneration) {
return;
}
blockedRooms.add(channelId);
throw new Error(`Could not refresh Buzz room membership for ${channelId}`);
};
const refreshMembershipOnce = (channelId: string): Promise<void> => {
const current = refreshes.get(channelId);
if (current) {
current.generation += 1;
return current.promise;
}
const state = {
generation: 1,
lastAttemptedGeneration: 0,
promise: Promise.resolve(),
} satisfies RefreshState;
state.promise = refreshMembership(channelId, state).finally(() => {
if (refreshes.get(channelId) === state) {
refreshes.delete(channelId);
}
if (
state.generation !== state.lastAttemptedGeneration &&
pendingMemberships.has(channelId) &&
!params.signal?.aborted
) {
void refreshMembershipOnce(channelId).catch(reportSystemEventError);
}
});
refreshes.set(channelId, state);
return state.promise;
};
const handleSystemEvent = (event: Event): Promise<void> | undefined => {
if (!markSystemEventSeen(event.id)) {
return undefined;
}
const channelId = event.tags
.find((tag) => tag[0] === "h")?.[1]
?.trim()
.toLowerCase();
if (!channelId) {
return undefined;
}
if (event.kind === BUZZ_ROOM_METADATA_EDIT_KIND) {
params.onRoomMetadataChanged?.(channelId);
return undefined;
}
const membership = memberships.get(channelId);
if (!membership) {
return undefined;
}
const change = parseBuzzRoomMembershipChangeEvent(event, membership);
if (!change) {
return undefined;
}
// System events invalidate membership; the relay-signed roster decides the
// final state. Removals deny immediately, while joins wait for confirmation.
const expected = change.type === "member_joined" ? "present" : "absent";
const pending = pendingMemberships.get(channelId) ?? new Map<string, ExpectedMembership>();
pending.set(change.targetPublicKey, expected);
pendingMemberships.set(channelId, pending);
if (expected === "absent") {
const denied = deniedMembers.get(channelId) ?? new Set<string>();
denied.add(change.targetPublicKey);
deniedMembers.set(channelId, denied);
}
if (change.targetPublicKey === params.botPublicKey) {
blockedRooms.add(channelId);
}
return refreshMembershipOnce(channelId);
};
const handleRoomEvent = (event: Event, reservation?: BuzzReplayDispatchReservation) => {
if (isBuzzInboundMessageKind(event.kind)) {
params.onMessageEvent(event, isMember, reservation);
return;
}
void handleSystemEvent(event)?.catch(reportSystemEventError);
};
for (const channelId of params.channelIds) {
if (memberships.get(channelId)?.roles.get(params.botPublicKey) !== "bot") {
throw new Error(`Buzz bot does not have the Bot role in configured room ${channelId}`);
}
}
let resolveHistorical: (() => void) | undefined;
let rejectHistorical: ((error: Error) => void) | undefined;
const historicalReady = new Promise<void>((resolve, reject) => {
resolveHistorical = resolve;
rejectHistorical = reject;
});
const historicalTimeout = setTimeout(() => {
const error = new Error("Timed out loading Buzz room membership changes");
rejectHistorical?.(error);
params.relay.close();
}, MEMBERSHIP_READY_TIMEOUT_MS);
const subscriptions: Array<ReturnType<Relay["prepareSubscription"]>> = [];
try {
// Snapshot membership before room history so startup memory stays bounded.
// Buzz emits these filters in order: system changes since session start
// update or deny membership before the following message history is handled.
for (const channelId of params.channelIds) {
subscriptions.push(
openBuzzRelaySubscription(
params.relay,
[
{
kinds: [BUZZ_ROOM_SYSTEM_KIND, BUZZ_ROOM_METADATA_EDIT_KIND],
"#h": [channelId],
since: params.since,
},
{
kinds: [...BUZZ_INBOUND_MESSAGE_KINDS],
"#h": [channelId],
since: params.messageSince,
limit: params.messageLimit,
},
],
{
onevent: (event) => {
if (!historicalRooms.has(channelId) && isBuzzInboundMessageKind(event.kind)) {
const page = historyPages.get(channelId);
if (page) {
page.count += 1;
page.oldest = Math.min(page.oldest, event.created_at);
} else {
historyPages.set(channelId, { count: 1, oldest: event.created_at });
}
}
handleRoomEvent(event);
},
oneose: () => {
historicalRooms.add(channelId);
if (historicalRooms.size === params.channelIds.length) {
resolveHistorical?.();
}
},
onclose: (reason) => {
if (!historicalRooms.has(channelId)) {
rejectHistorical?.(
new Error(`Buzz membership subscription closed for ${channelId}: ${reason}`),
);
} else if (
reason !== "shutdown" &&
reason !== "relay connection closed by us" &&
reason !== MEMBERSHIP_TRACKER_SETUP_CLOSE_REASON &&
!params.signal?.aborted
) {
params.onFatalError?.(
new Error(`Buzz membership subscription closed for ${channelId}: ${reason}`),
);
}
},
},
),
);
}
await historicalReady;
} catch (error) {
if (params.relay.connected) {
for (const subscription of subscriptions) {
if (!subscription.closed) {
subscription.close(MEMBERSHIP_TRACKER_SETUP_CLOSE_REASON);
}
}
}
throw error;
} finally {
clearTimeout(historicalTimeout);
}
return {
memberships: () => memberships,
catchUpHistory: async () => {
for (const channelId of params.channelIds) {
const page = historyPages.get(channelId);
if (params.signal?.aborted) {
return;
}
if (!page || page.count < params.messageLimit) {
continue;
}
try {
const outcome = await catchUpBuzzRoomHistory({
relay: params.relay,
channelId,
since: params.messageSince,
until: page.oldest,
limit: params.messageLimit,
reserveCapacity: params.reserveDispatchCapacity,
onEvent: handleRoomEvent,
signal: params.signal,
});
if (outcome === "timestamp-over-limit") {
params.onHistoryError?.(
new Error(
`Buzz room ${channelId} kept more than ${BUZZ_REPLAY_DISPATCH_MAX_PENDING} additional messages at one timestamp; older history was not recovered`,
),
);
}
} catch (error) {
if (params.signal?.aborted) {
return;
}
reportSystemEventError(
error instanceof Error
? error
: new Error(`Buzz room history recovery failed for ${channelId}`, { cause: error }),
);
return;
}
}
},
};
}