/** * Background poller for pi-c2c auto-delivery. * * Drains three message sources (per-repo broker, sessions broker, public * relay), deduplicates, spools to disk, and injects into the pi transcript. * Failures on any single source are isolated so a relay hiccup doesn't lose * local messages. */ import type { C2cMessage } from "./c2c-cli.ts"; import { DeliveryDedup, deliveryOptionsFor, filterNovel, formatEnvelope, markDelivered, } from "./delivery.ts"; import { clearSpool, readSpool, writeSpool } from "./spool.ts"; import { extractStatusMessages } from "./peer-status.ts"; import { SPOOL_DIR, relayToC2c } from "./common.ts"; import type { PeerStatusStore } from "./peer-status.ts"; import type { LiveTelemetry, MessageSource } from "./telemetry.ts"; import type { C2cDeliveryDetails } from "./ui/compact-message.ts"; /** Dependencies injected by the extension factory. */ export interface PollerDeps { getCli(): import("./c2c-cli.ts").C2cCli | null; getIdentity(): import("./identity.ts").Identity | null; isShuttingDown(): boolean; sessionsBrokerRoot: string | undefined; isRelayRegistered(): boolean; getRelayAddress(): string | undefined; dedup: DeliveryDedup; peerStatusStore: PeerStatusStore; telemetry: LiveTelemetry; sendMessage(args: { customType: string; content: string; display: boolean; details: unknown }, opts: { triggerTurn?: boolean; deliverAs?: "steer" | "followUp" | "nextTurn" }): void; getQueuedSinceMs(): number | undefined; setQueuedSinceMs(v: number | undefined): void; } /** * Inject already-filtered messages into the transcript via `pi.sendMessage`. * Returns true if injection succeeded; false if it threw (stale runtime). */ export function inject( novel: C2cMessage[], deps: PollerDeps, ): boolean { if (novel.length === 0) return true; const allNonurgent = novel.every((m) => m.nonurgent === true); const identity = deps.getIdentity(); const body = novel.map((m) => formatEnvelope(m, identity?.alias)).join("\n\n"); const details: C2cDeliveryDetails = { count: novel.length, senders: [...new Set(novel.map((m) => m.from_alias || "unknown"))], selfAlias: identity?.alias, source: novel.every((m) => m.brokerSource === novel[0]?.brokerSource) ? novel[0]?.brokerSource : undefined, }; if (allNonurgent) { deps.setQueuedSinceMs(Date.now()); } try { deps.sendMessage( { customType: "c2c", content: body, display: true, details }, deliveryOptionsFor({ nonurgent: allNonurgent }), ); } catch { return false; } deps.telemetry.recordInjected(novel.length); return true; } /** * One tick of the background poller. Drains all sources, deduplicates, * spools to disk, and injects novel messages. Best-effort and loss-resistant: * - drained messages are spooled BEFORE injection (crash-safe); * - dedup is marked only AFTER successful injection; * - during shutdown, nothing is drained/injected. */ export async function pollTick(deps: PollerDeps): Promise { const cli = deps.getCli(); const identity = deps.getIdentity(); if (!cli || !identity || deps.isShuttingDown()) return; const sid = identity.sessionId; deps.telemetry.beginPoll(); if (deps.isShuttingDown()) { deps.telemetry.endPoll(); return; } const drained: C2cMessage[] = []; // Drain per-repo broker try { const localMsgs = await cli.pollInbox(); for (const m of localMsgs) m.brokerSource = "local"; drained.push(...localMsgs); deps.telemetry.recordBrokerOk("local"); } catch (e) { deps.telemetry.recordBrokerError("local", e); } // Drain sessions broker (cross-repo) if (deps.sessionsBrokerRoot) { try { const sessionMsgs = await cli.pollInbox({ brokerRoot: deps.sessionsBrokerRoot }); for (const m of sessionMsgs) m.brokerSource = "sessions"; drained.push(...sessionMsgs); deps.telemetry.recordBrokerOk("sessions"); } catch (e) { deps.telemetry.recordBrokerError("sessions", e); } } // Drain public relay if (deps.isRelayRegistered() && deps.getRelayAddress()) { try { const relayMsgs = await cli.relayDmPoll(deps.getRelayAddress()!); for (const c of relayToC2c(relayMsgs)) drained.push(c); deps.telemetry.recordRelayOk(); } catch (e) { deps.telemetry.recordRelayError(e); } } // Record telemetry for (const m of drained) { const source: MessageSource = m.source === "relay" ? "relay" : m.brokerSource ?? "local"; deps.telemetry.recordReceived({ from: m.from_alias, content: m.content, source }); } // Replay the spool, then filter statuses from the complete consumed batch so // a status persisted by an older extension version can never reach chat. const spooled = readSpool(SPOOL_DIR, sid); deps.telemetry.recordSpoolCount(spooled.length); const combined = [...spooled, ...drained]; const { messages: deliverable } = extractStatusMessages(combined, deps.peerStatusStore); deps.telemetry.recordPeerStatusCount(deps.peerStatusStore.live().length); const novel = filterNovel(deliverable, deps.dedup); if (novel.length === 0) { if (combined.length > 0) { clearSpool(SPOOL_DIR, sid); deps.telemetry.recordSpoolCount(0); } deps.telemetry.endPoll(); return; } // Persist before injecting (crash-safe) writeSpool(SPOOL_DIR, sid, novel); deps.telemetry.recordSpoolCount(novel.length); if (deps.isShuttingDown()) { deps.telemetry.endPoll(); return; } if (inject(novel, deps)) { markDelivered(novel, deps.dedup); clearSpool(SPOOL_DIR, sid); deps.telemetry.recordSpoolCount(0); } // else: spool persists, dedup unmarked → retried next tick deps.telemetry.endPoll(); }