// The engine pipe (#1507, CHAT-03 §1). // // One EngineLink per execution owns the engine's stdin and reads its stdout. // Output is split on LF only (framing.mjs); every line reaches the controller // in order, tagged with the link's execution, so a replaced engine's late // output can be dropped by incarnation (H14). // // Write outcomes (§1): // written the whole line was accepted by the pipe (the write callback ran // without an error). That is not native consumption. // unknown the write returned an error (EPIPE), or its callback did not run // within the bound. A partial line may be in the pipe, so the link // is poisoned and never written again. // `acknowledged` is the response to a request, which `request()` reports // separately. A request whose response does not come within its bound is // reported as a timeout; the controller decides what that means. import { LineSplitter, encodeLine, parseLine } from "./framing.mjs"; export const WRITE_TIMEOUT_MS = 5000; export class EngineLink { // `onLine(link, value, bytes)` gets every parsed non-response line; // `onResponse(link, value, entry)` every response, before its promise // resolves; `onGap(link, kind, detail)` an unparseable line, an unmatched // response, an overlong line or EOF with a partial line. constructor({ execution, stdin, stdout, onLine, onResponse = () => {}, onGap, onEnd = () => {}, writeTimeoutMs = WRITE_TIMEOUT_MS }) { this.execution = execution; this.stdin = stdin; this.onLine = onLine; this.onResponse = onResponse; this.onGap = onGap; this.onEnd = onEnd; this.writeTimeoutMs = writeTimeoutMs; this.poisoned = null; this.pending = new Map(); this.serial = 0; this.bytesWritten = 0; this.ended = false; this.writes = []; stdin.on("error", (err) => { this.stdinError = err.code ?? err.message; }); this.splitter = new LineSplitter((line) => this.#line(line), { onOverflow: () => this.onGap(this, "overlong-line", null) }); stdout.on("data", (chunk) => this.splitter.push(chunk)); stdout.on("end", () => { this.ended = true; if (this.splitter.pending() > 0) this.onGap(this, "partial-line-at-eof", { bytes: this.splitter.pending() }); this.onEnd(this); }); stdout.on("error", () => {}); } #line(line) { const bytes = Buffer.byteLength(line, "utf8"); const parsed = parseLine(line); if (parsed.error) return this.onGap(this, "unparseable-line", { bytes, error: parsed.error }); const value = parsed.value; if (value.type === "response") { const entry = typeof value.id === "string" ? this.pending.get(value.id) : undefined; if (!entry) return this.onGap(this, "unmatched-response", { bytes, command: typeof value.command === "string" ? value.command.slice(0, 40) : null }); this.pending.delete(value.id); entry.late = entry.timedOut; this.onResponse(this, value, entry); entry.resolve({ response: value, late: entry.late }); return undefined; } return this.onLine(this, value, bytes); } poison(reason) { if (!this.poisoned) this.poisoned = reason; } // Writes one record. Never writes to a poisoned link. write(value) { if (this.poisoned) return Promise.resolve({ outcome: "refused", reason: "poisoned" }); const line = encodeLine(value); const bytes = Buffer.byteLength(line, "utf8"); return new Promise((resolve) => { let done = false; const finish = (o) => { if (done) return; done = true; clearTimeout(timer); if (o.outcome === "unknown") this.poison(o.reason); else this.bytesWritten += bytes; this.writes.push({ type: value.type, id: value.id, outcome: o.outcome }); resolve(o); }; const timer = setTimeout(() => finish({ outcome: "unknown", reason: "write-timeout" }), this.writeTimeoutMs); try { this.stdin.write(line, (err) => finish(err ? { outcome: "unknown", reason: err.code ?? "write-error" } : { outcome: "written" })); } catch (err) { finish({ outcome: "unknown", reason: err.code ?? "write-error" }); } }); } // Sends a command. Returns { id, written, response }: `written` settles // with the write outcome, `response` with { response } or { timeout: true } // or { unsent: outcome }. The response is registered before the write, so // it can never be unmatched. request(type, fields = {}, { timeoutMs = 5000, beforeWrite = null } = {}) { const id = `${type}-${++this.serial}`; let resolve; const response = new Promise((r) => (resolve = r)); const entry = { id, type, resolve, timedOut: false, late: false }; this.pending.set(id, entry); beforeWrite?.(id); const written = this.write({ id, type, ...fields }); written.then((w) => { if (w.outcome !== "written") { this.pending.delete(id); resolve({ unsent: w }); return; } setTimeout(() => { if (!this.pending.has(id)) return; entry.timedOut = true; // kept, so a late answer is recognised and not a gap resolve({ timeout: true }); }, timeoutMs); }); return { id, written, response }; } }