Controller, claim store, live-session guard, engine link and seal, turn tracker, cohort force stop and recovery, client library, transcript and mediated terminal, with the fake engine and tests. Fixtures only; no live cutover. Dewey built it. Darkwing (comment 26690) and Filbert (comment 26694) approved round 2. Manifest I1-r2-manifest.sha256 (2b48e333, 27 files). Suites on an export: conversation 152/152, control-board 124, webui 14, seat 19, chat-00/01/01c checks, and all nine scripts/test-*.sh green. Follow-ups for I3 are in DEFERRED. Gate E stays with Jason. Co-Authored-By: Claude Opus 5.5 <[email protected]>
235 lines
8.4 KiB
JavaScript
235 lines
8.4 KiB
JavaScript
// 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);
|
||
});
|
||
}
|
||
}
|