// Slice 1 S5 (#1522): the Console's read of the bus. It runs the S4 human // transport (packages/bus/src/human-cli.mjs, lead decision 70) as a child // with spawn, not spawnSync, so one slow read never stalls the HTTP server. // Only the four verbs of the Q1 module (lead decision 56) are reachable: // the WebUI writes nothing, not even seen (Q4). The bus does the proof; a // server started inside an agent run gets `human-required` from the broker. import { spawn } from 'node:child_process'; import { HUMAN_CLI } from '../../cli/src/transport.mjs'; import { readHostState } from '../../cli/src/host.mjs'; import { views } from '../../bus/src/views.mjs'; export class BusReadError extends Error { constructor(code, message = code) { super(message); this.code = code; } } // The broker's identifier rule (packages/bus/src/broker.mjs `id`), checked // here so a bad subject is a 400 and never reaches the bus. export const subjectOk = s => typeof s === 'string' && s.length > 0 && s.length <= 160 && /^[a-zA-Z0-9][a-zA-Z0-9_.:/-]*$/.test(s); // Same protocol as humanTransport in packages/cli/src/transport.mjs: the // request on stdin, JSON on stdout, or a bus error code as the last // `^[a-z-]{1,64}$` line of stderr. export function humanCall({ socket, business, env = process.env, cli = HUMAN_CLI, timeoutMs = 30000, maxBytes = 8 * 1024 * 1024 }) { return (verb, args = {}) => new Promise((resolve, reject) => { let child; try { child = spawn(process.execPath, [cli, socket], { env, stdio: ['pipe', 'pipe', 'pipe'] }); } catch { return reject(new BusReadError('outcome-unknown', 'the bus transport did not start')); } const out = [], err = []; let size = 0, settled = false; const done = (fn, value) => { if (!settled) { settled = true; clearTimeout(timer); fn(value); } }; const timer = setTimeout(() => { child.kill('SIGKILL'); done(reject, new BusReadError('outcome-unknown', 'the bus transport did not finish in time')); }, timeoutMs); child.stdout.on('data', chunk => { size += chunk.length; if (size > maxBytes) { child.kill('SIGKILL'); return done(reject, new BusReadError('response-too-large')); } out.push(chunk); }); child.stderr.on('data', chunk => { if (err.length < 64) err.push(chunk); }); child.on('error', () => done(reject, new BusReadError('outcome-unknown', 'the bus transport did not start'))); child.on('close', status => { if (status === null) return done(reject, new BusReadError('outcome-unknown', 'the bus transport did not finish')); if (status !== 0) { const code = Buffer.concat(err).toString('utf8').trim().split('\n').reverse().find(l => /^[a-z-]{1,64}$/.test(l)) ?? 'invalid-response'; return done(reject, new BusReadError(code)); } try { done(resolve, JSON.parse(Buffer.concat(out).toString('utf8'))); } catch { done(reject, new BusReadError('invalid-response')); } }); child.stdin.on('error', () => {}); // A child that exits early reports through 'close'. child.stdin.end(`${JSON.stringify({ business, verb, args })}\n`); }); } // The business is the --business flag, or else the live bus host's, read // on every request so a host started after the Console is picked up. export function busReader({ dataRoot, socket, business = null, env = process.env, cli = HUMAN_CLI, timeoutMs = 30000 }) { const call = async (verb, args) => { let id = business; if (!id) { let state; try { state = readHostState(dataRoot); } catch { throw new BusReadError('no-bus-host', 'the bus host state file is unreadable'); } if (!state?.live) throw new BusReadError('no-bus-host', 'no bus host is running and no --business was given'); id = state.business; } return humanCall({ socket, business: id, env, cli, timeoutMs })(verb, args); }; return views({ call }); } // HTTP status by bus error code. 403: the bus refused this caller; 503: no // bus to ask, or no answer; 502: an answer that can't be used. const STATUS = { 'human-required': 403, unauthenticated: 403, 'unknown-business': 403, 'read-only': 403, 'outcome-unknown': 503, 'no-bus-host': 503, 'not-configured': 503, 'invalid-request': 400, }; export const busStatus = code => STATUS[code] ?? 502;