// The CHAT-03 client library (#1507, CHAT-03 §1, deviation V-1). // // It speaks the controller's line protocol over the Unix socket: `hello` // with a grant, then CHAT-01 v2 client requests, each carrying the // controller's incarnation token. The library never resends. A request still // pending when the connection drops, or one refused `stale-incarnation`, is // shown as "outcome unknown, check the transcript" (lead decision 25, H21– // H23). An actor may copy its text into a new request; the library doesn't // present that as a safe retry. // // Drafts are client-local in CHAT-03: the draft ID and revision name the // client's own buffer, and the text travels in the request envelope. CHAT-04 // owns server-side drafts. import { connect as netConnect } from "node:net"; import { LineSplitter, encodeLine, parseLine } from "./framing.mjs"; import { newId, targetOf } from "./records.mjs"; export const OUTCOME_UNKNOWN = "outcome unknown, check the transcript"; const STALE = "stale-incarnation"; const FINAL = new Set(["finished", "failed", "dispatch-refused", "delivery-unknown"]); export class ConversationClient { constructor({ socketPath, grant = "grant-local", connect = netConnect }) { this.socketPath = socketPath; this.grant = grant; this.netConnect = connect; this.sock = null; this.incarnation = null; this.connection = null; this.binding = null; this.streamEpoch = null; this.pending = new Map(); this.unknown = []; this.receipts = new Map(); this.events = []; this.pushes = []; this.listeners = new Set(); this.waiters = new Set(); this.closed = true; } get target() { return this.binding ? targetOf(this.binding) : null; } get isController() { return !!this.connection && this.binding?.controllerConnection === this.connection.id; } // `fn(message)` for every push, reply and status change. on(fn) { this.listeners.add(fn); return () => this.listeners.delete(fn); } #notify(msg) { for (const fn of this.listeners) fn(msg); for (const w of [...this.waiters]) w(); } // Resolves with the welcome. A new incarnation turns every request still // pending under the old one into outcome-unknown; none is resent (H23). async connect() { this.#dropPending("disconnected"); const sock = this.netConnect(this.socketPath); this.sock = sock; await new Promise((resolve, reject) => { sock.once("connect", resolve); sock.once("error", reject); }); this.closed = false; const welcome = new Promise((resolve, reject) => { this.welcomeWaiter = { resolve, reject }; }); const splitter = new LineSplitter((line) => this.#onLine(line)); sock.on("data", (c) => splitter.push(c)); sock.on("error", () => {}); sock.on("close", () => { if (this.sock !== sock) return; this.closed = true; this.welcomeWaiter?.reject(new Error("the controller closed the connection")); this.welcomeWaiter = null; this.#dropPending("disconnected"); this.#notify({ type: "status", status: "disconnected" }); }); sock.write(encodeLine({ type: "hello", grant: this.grant })); return welcome; } close() { this.sock?.destroy(); } #dropPending(reason) { for (const [id, p] of this.pending) { const entry = { request: id, operation: p.envelope.command.operation, incarnation: p.incarnation, reason, display: OUTCOME_UNKNOWN }; this.unknown.push(entry); p.resolve({ outcome: "outcome-unknown", refusal: null, display: OUTCOME_UNKNOWN, reason, request: null, receipt: null }); this.#notify({ type: "outcome-unknown", ...entry }); } this.pending.clear(); } // A receipt last seen short of a final state belongs to the old controller; // the new one doesn't know it. Its outcome is unknown (H20, H23). #unsettled() { for (const r of this.receipts.values()) { if (FINAL.has(r.state) || this.unknown.some((u) => u.receipt === r.id)) continue; const entry = { request: null, receipt: r.id, operation: "prompt", incarnation: this.incarnation, reason: STALE, display: OUTCOME_UNKNOWN }; this.unknown.push(entry); this.#notify({ type: "outcome-unknown", ...entry }); } } #onLine(line) { const parsed = parseLine(line); if (parsed.error) return; const msg = parsed.value; if (msg.type === "welcome") { if (this.incarnation !== null && this.incarnation !== msg.incarnation) { this.#dropPending(STALE); this.#unsettled(); } this.incarnation = msg.incarnation; this.connection = msg.connection; this.binding = msg.binding; this.streamEpoch = msg.streamEpoch; this.welcomeWaiter?.resolve(msg); this.welcomeWaiter = null; this.#notify(msg); return; } if (msg.type === "refused" && this.welcomeWaiter) { this.welcomeWaiter.reject(Object.assign(new Error(`refused: ${msg.refusal}`), { refusal: msg.refusal })); this.welcomeWaiter = null; return; } if (msg.type === "reply") { const p = this.pending.get(msg.id); if (!p) return; this.pending.delete(msg.id); const reply = { ...msg }; if (msg.refusal === STALE) { reply.display = OUTCOME_UNKNOWN; this.unknown.push({ request: msg.id, operation: p.envelope.command.operation, incarnation: p.incarnation, reason: STALE, display: OUTCOME_UNKNOWN }); } if (msg.receipt) this.receipts.set(msg.receipt.id, msg.receipt); p.resolve(reply); this.#notify(reply); return; } if (msg.type === "push") { this.pushes.push(msg); if (msg.kind === "binding") this.binding = msg.binding; else if (msg.kind === "connection" && msg.connection.id === this.connection?.id) this.connection = msg.connection; else if (msg.kind === "receipt") this.receipts.set(msg.receipt.id, msg.receipt); else if (msg.kind === "event") this.events.push(msg.event); this.#notify(msg); } } // The envelope for one operation on the current target. envelope(operation, fields = {}) { if (!this.connection || !this.binding) throw new Error("not connected"); return { version: 2, kind: "clientRequest", id: newId("req"), connection: this.connection.id, target: this.target, command: { operation, ...fields } }; } // Sends one request and resolves with the controller's reply. request(operation, fields = {}, text) { return this.send(this.envelope(operation, fields), text); } // Sends an envelope under the token this client holds now. An envelope is // sent once; the library has no retry path. send(envelope, text, { incarnation = this.incarnation } = {}) { if (this.closed) return Promise.resolve({ outcome: "refused:channel", refusal: "channel", display: "not connected" }); return new Promise((resolve) => { this.pending.set(envelope.id, { envelope, incarnation, resolve }); const msg = { type: "request", incarnation, request: envelope }; if (text !== undefined) msg.text = text; this.sock.write(encodeLine(msg)); }); } prompt(text, draft = { id: newId("draft"), revision: 1 }) { return this.request("prompt", { draft: draft.id, draftRevision: draft.revision }, text); } observe({ cursor = null, limit = 50 } = {}) { return this.request("observe", { cursor, limit }); } takeover() { return this.request("takeover"); } interrupt() { return this.request("interrupt"); } // Issue, confirm and use a confirmation for a force stop, a recovery or // recovery control. async confirmed(operation, fields = {}) { const issued = await this.request("issue-confirmation", { operationToConfirm: operation }); if (issued.outcome !== "confirmation-issued") return issued; const id = issued.data.confirmation.id; const answered = await this.request("answer-confirmation", { confirmation: id, answer: "confirm" }); if (answered.outcome !== "confirmation-confirmed") return answered; return this.request(operation, { ...fields, confirmation: id }); } receipt(id) { return this.receipts.get(id) ?? null; } // Resolves when `pred()` holds, or false after `ms`. waitFor(pred, ms = 5000) { if (pred()) return Promise.resolve(true); return new Promise((resolve) => { const w = () => { if (!pred()) return; this.waiters.delete(w); clearTimeout(t); resolve(true); }; const t = setTimeout(() => { this.waiters.delete(w); resolve(false); }, ms); this.waiters.add(w); }); } }