feat(ipc): port 3g IPC + bgbash — framed Unix-socket server/client + background job registry, wire bash tools
This commit is contained in:
@@ -0,0 +1,44 @@
|
||||
/**
|
||||
* Background bash control — list, cancel, and inspect background processes.
|
||||
* Mirrors `apps/infrastructure/src/bgbash/control.rs`.
|
||||
*/
|
||||
import { type BashJob } from "./job.ts";
|
||||
|
||||
/** Central registry of all running background bash jobs. */
|
||||
export class BashControl {
|
||||
private jobs = new Map<string, BashJob>();
|
||||
|
||||
/** Register a new background job. */
|
||||
register(job: BashJob): void {
|
||||
this.jobs.set(job.id, job);
|
||||
}
|
||||
|
||||
/** Cancel a job by ID. Returns true if found and cancelled. */
|
||||
cancel(id: string): boolean {
|
||||
const job = this.jobs.get(id);
|
||||
if (!job) return false;
|
||||
job.cancel();
|
||||
this.jobs.delete(id);
|
||||
return true;
|
||||
}
|
||||
|
||||
/** List all active jobs as [id, command, isRunning] tuples. */
|
||||
list(): Array<[string, string, boolean]> {
|
||||
this.prune();
|
||||
return [...this.jobs.entries()].map(([id, j]) => [id, j.command, j.isRunning()]);
|
||||
}
|
||||
|
||||
/** Clean up completed jobs. */
|
||||
prune(): void {
|
||||
for (const [id, job] of this.jobs) {
|
||||
if (!job.isRunning()) this.jobs.delete(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Global accessor for the shared BashControl singleton. */
|
||||
let _bashControl: BashControl | null = null;
|
||||
export function bashControl(): BashControl {
|
||||
if (!_bashControl) _bashControl = new BashControl();
|
||||
return _bashControl;
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
/**
|
||||
* Background bash module — job + control registry.
|
||||
* Mirrors `apps/infrastructure/src/bgbash/mod.rs`.
|
||||
*/
|
||||
export { BashJob, spawnBashJob } from "./job.ts";
|
||||
export { BashControl, bashControl } from "./control.ts";
|
||||
@@ -0,0 +1,65 @@
|
||||
/**
|
||||
* Background bash job — spawns a `bash -c` subprocess and tracks its life.
|
||||
* Mirrors `apps/infrastructure/src/bgbash/job.rs`.
|
||||
*/
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import { randomUUID } from "node:crypto";
|
||||
|
||||
/** A handle to a spawned background bash job. */
|
||||
export class BashJob {
|
||||
readonly id: string;
|
||||
readonly command: string;
|
||||
private process: ChildProcess | null = null;
|
||||
private cancelled = false;
|
||||
private finished = false;
|
||||
|
||||
constructor(command: string) {
|
||||
this.id = randomUUID();
|
||||
this.command = command;
|
||||
}
|
||||
|
||||
/** Spawn the `bash -c` child process and begin tracking it. */
|
||||
start(): void {
|
||||
const child = spawn("bash", ["-c", this.command], {
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
});
|
||||
this.process = child;
|
||||
|
||||
// Drain stdout/stderr so the process doesn't block on a full pipe, and
|
||||
// mark finished when it exits.
|
||||
child.stdout?.resume();
|
||||
child.stderr?.resume();
|
||||
child.on("exit", () => {
|
||||
this.finished = true;
|
||||
});
|
||||
child.on("error", () => {
|
||||
this.finished = true;
|
||||
});
|
||||
}
|
||||
|
||||
/** Cancel the job: set the cancelled flag and kill the child process. */
|
||||
cancel(): void {
|
||||
this.cancelled = true;
|
||||
if (this.process && !this.finished) {
|
||||
try {
|
||||
this.process.kill("SIGKILL");
|
||||
} catch {
|
||||
/* already gone */
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Whether the job is currently running. */
|
||||
isRunning(): boolean {
|
||||
if (this.cancelled || this.finished) return false;
|
||||
if (!this.process) return false;
|
||||
return this.process.exitCode === null && !this.finished;
|
||||
}
|
||||
}
|
||||
|
||||
/** Spawn a background bash job and return its handle (node-side). */
|
||||
export function spawnBashJob(command: string): BashJob {
|
||||
const job = new BashJob(command);
|
||||
job.start();
|
||||
return job;
|
||||
}
|
||||
@@ -0,0 +1,43 @@
|
||||
/**
|
||||
* IPC client — connects to the daemon's Unix socket and sends/receives
|
||||
* framed JSON messages.
|
||||
* Mirrors `apps/infrastructure/src/ipc/client.rs`.
|
||||
*/
|
||||
import * as net from "node:net";
|
||||
import { Connection } from "./conn.ts";
|
||||
|
||||
/** A thread-safe IPC client connected to a Zesdex daemon over a Unix socket. */
|
||||
export class IpcClient {
|
||||
private conn: Connection;
|
||||
|
||||
private constructor(conn: Connection) {
|
||||
this.conn = conn;
|
||||
}
|
||||
|
||||
/** Connect to a Unix socket at `path`. */
|
||||
static connectUnix(path: string): Promise<IpcClient> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const socket = net.createConnection(path);
|
||||
socket.once("error", reject);
|
||||
socket.once("connect", () => {
|
||||
socket.removeListener("error", reject);
|
||||
resolve(new IpcClient(new Connection(socket)));
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
/** Send a JSON-serializable message as one framed frame. */
|
||||
send<T>(msg: T): void {
|
||||
this.conn.send(msg);
|
||||
}
|
||||
|
||||
/** Receive the next message and JSON-decode it (null on close). */
|
||||
async receive<T>(): Promise<T | null> {
|
||||
return this.conn.receiveJson<T>();
|
||||
}
|
||||
|
||||
/** Close the connection. */
|
||||
close(): void {
|
||||
this.conn.close();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
/**
|
||||
* IPC connection — a framed, JSON-message socket over node:net.
|
||||
* Mirrors `apps/infrastructure/src/ipc/conn.rs`.
|
||||
*
|
||||
* Wraps a `net.Socket` (or `net.Server` connection) with the length-prefixed
|
||||
* framing from `frame.ts`. Provides `send` (write one framed message) and
|
||||
* an async-style receive via an event emitter / queue.
|
||||
*/
|
||||
import * as net from "node:net";
|
||||
import { FrameReassembler, encodeFrame } from "./frame.ts";
|
||||
|
||||
/** A connection over which framed JSON messages are exchanged. */
|
||||
export class Connection {
|
||||
private socket: net.Socket;
|
||||
private reassembler = new FrameReassembler();
|
||||
private messageQueue: Buffer[] = [];
|
||||
private waiters: Array<(buf: Buffer | null) => void> = [];
|
||||
|
||||
constructor(socket: net.Socket) {
|
||||
this.socket = socket;
|
||||
socket.on("data", (chunk: Buffer) => {
|
||||
let frames: Buffer[];
|
||||
try {
|
||||
frames = this.reassembler.feed(chunk);
|
||||
} catch (e) {
|
||||
// oversized frame — fail the connection
|
||||
const err = e as Error;
|
||||
this._wakeWaiters(null);
|
||||
this.socket.destroy(new Error(err.message));
|
||||
return;
|
||||
}
|
||||
for (const frame of frames) {
|
||||
this._enqueue(frame);
|
||||
}
|
||||
});
|
||||
socket.on("end", () => {
|
||||
this.messageQueue = [];
|
||||
this._wakeWaiters(null);
|
||||
});
|
||||
socket.on("error", () => {
|
||||
this._wakeWaiters(null);
|
||||
});
|
||||
}
|
||||
|
||||
/** Serialize and send a message as one framed frame. */
|
||||
send<T>(msg: T): void {
|
||||
this.socket.write(encodeFrame(JSON.stringify(msg)));
|
||||
}
|
||||
|
||||
/**
|
||||
* Receive the next complete message as a Buffer (parsed JSON).
|
||||
* Returns `null` when the connection closes.
|
||||
*/
|
||||
receive(): Promise<Buffer | null> {
|
||||
if (this.messageQueue.length > 0) {
|
||||
return Promise.resolve(this.messageQueue.shift()!);
|
||||
}
|
||||
return new Promise((resolve) => {
|
||||
this.waiters.push(resolve);
|
||||
});
|
||||
}
|
||||
|
||||
/** Receive the next message and JSON-decode it. */
|
||||
async receiveJson<T>(): Promise<T | null> {
|
||||
const buf = await this.receive();
|
||||
if (buf === null) return null;
|
||||
return JSON.parse(buf.toString("utf8")) as T;
|
||||
}
|
||||
|
||||
/** Close the underlying socket. */
|
||||
close(): void {
|
||||
this.socket.end();
|
||||
}
|
||||
|
||||
private _enqueue(buf: Buffer): void {
|
||||
const waiter = this.waiters.shift();
|
||||
if (waiter) {
|
||||
waiter(buf);
|
||||
} else {
|
||||
this.messageQueue.push(buf);
|
||||
}
|
||||
}
|
||||
|
||||
private _wakeWaiters(buf: Buffer | null): void {
|
||||
while (this.waiters.length > 0) {
|
||||
const w = this.waiters.shift()!;
|
||||
w(buf);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
/**
|
||||
* Length-prefixed framing for Unix-socket IPC.
|
||||
* Mirrors `apps/infrastructure/src/ipc/frame.rs`.
|
||||
*
|
||||
* Every message on the wire is encoded as:
|
||||
* ```text
|
||||
* [ 4-byte big-endian payload length ][ payload bytes (JSON) ]
|
||||
* ```
|
||||
*
|
||||
* In TS this operates over `node:net` Socket (byte stream). A helper tracks
|
||||
* a partial buffer so a single `data` event may contain multiple frames or a
|
||||
* partial frame.
|
||||
*/
|
||||
|
||||
/** Maximum payload size accepted (64 MiB). */
|
||||
export const MAX_PAYLOAD = 64 * 1024 * 1024;
|
||||
|
||||
/** Buffer for accumulating a frame across multiple `data` events. */
|
||||
export class FrameReassembler {
|
||||
private chunks: Buffer[] = [];
|
||||
private buffered = 0;
|
||||
private pendingLength: number | null = null;
|
||||
|
||||
/**
|
||||
* Feed raw bytes and return any complete frames that were buffered.
|
||||
* Each returned element is one complete JSON payload.
|
||||
*/
|
||||
feed(chunk: Buffer): Buffer[] {
|
||||
this.chunks.push(chunk);
|
||||
this.buffered += chunk.length;
|
||||
return this.tryFrames();
|
||||
}
|
||||
|
||||
private tryFrames(): Buffer[] {
|
||||
const frames: Buffer[] = [];
|
||||
for (;;) {
|
||||
// Ensure we have at least the 4-byte length prefix.
|
||||
if (this.buffered < 4) break;
|
||||
if (this.pendingLength === null) {
|
||||
this.pendingLength = this.peekLength();
|
||||
}
|
||||
const len = this.pendingLength;
|
||||
if (len === null) break;
|
||||
if (len > MAX_PAYLOAD) {
|
||||
throw new Error(`frame payload too large: ${len} bytes (max ${MAX_PAYLOAD})`);
|
||||
}
|
||||
if (this.buffered < 4 + len) break; // wait for full payload
|
||||
// Consume the full frame.
|
||||
const frame = this.take(4 + len);
|
||||
const payload = frame.subarray(4, 4 + len);
|
||||
frames.push(Buffer.from(payload));
|
||||
this.pendingLength = null;
|
||||
}
|
||||
return frames;
|
||||
}
|
||||
|
||||
private peekLength(): number | null {
|
||||
if (this.buffered < 4) return null;
|
||||
const first = this.take(4);
|
||||
// Put the 4 bytes back.
|
||||
this.unshift(first);
|
||||
return first.readUInt32BE(0);
|
||||
}
|
||||
|
||||
/** Take `n` bytes from the front of the buffer. */
|
||||
private take(n: number): Buffer {
|
||||
let remaining = n;
|
||||
const parts: Buffer[] = [];
|
||||
while (remaining > 0 && this.chunks.length > 0) {
|
||||
const head = this.chunks[0]!;
|
||||
if (head.length <= remaining) {
|
||||
parts.push(head);
|
||||
this.chunks.shift();
|
||||
remaining -= head.length;
|
||||
} else {
|
||||
parts.push(head.subarray(0, remaining));
|
||||
this.chunks[0] = head.subarray(remaining);
|
||||
remaining = 0;
|
||||
}
|
||||
}
|
||||
this.buffered -= n;
|
||||
return Buffer.concat(parts);
|
||||
}
|
||||
|
||||
/** Push `n` bytes back to the front of the buffer. */
|
||||
private unshift(buf: Buffer): void {
|
||||
this.chunks.unshift(buf);
|
||||
this.buffered += buf.length;
|
||||
}
|
||||
}
|
||||
|
||||
/** Encode a JSON payload into a single length-prefixed frame buffer. */
|
||||
export function encodeFrame(plainText: string | Buffer): Buffer {
|
||||
const payload = typeof plainText === "string" ? Buffer.from(plainText, "utf8") : plainText;
|
||||
if (payload.length > MAX_PAYLOAD) {
|
||||
throw new Error(`frame payload too large: ${payload.length} bytes (max ${MAX_PAYLOAD})`);
|
||||
}
|
||||
const lenBuf = Buffer.alloc(4);
|
||||
lenBuf.writeUInt32BE(payload.length, 0);
|
||||
return Buffer.concat([lenBuf, payload]);
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
/**
|
||||
* IPC module — frame, protocol, connection, server, client.
|
||||
* Mirrors `apps/infrastructure/src/ipc/mod.rs`.
|
||||
*/
|
||||
export { FrameReassembler, encodeFrame, MAX_PAYLOAD } from "./frame.ts";
|
||||
export { Connection } from "./conn.ts";
|
||||
export { IpcServer } from "./server.ts";
|
||||
export { IpcClient } from "./client.ts";
|
||||
export type {
|
||||
KeyAction,
|
||||
KeyActionName,
|
||||
ClientRequest,
|
||||
MessageEntry,
|
||||
ToastEntry,
|
||||
StatePayload,
|
||||
DaemonFrame,
|
||||
} from "./protocol.ts";
|
||||
@@ -0,0 +1,75 @@
|
||||
/**
|
||||
* Wire types for the Zesdex IPC protocol.
|
||||
* Mirrors `apps/infrastructure/src/ipc/protocol.rs`.
|
||||
*/
|
||||
|
||||
/** A resolved key press sent from the daemon to the client. */
|
||||
export type KeyActionName =
|
||||
| "Char"
|
||||
| "Enter"
|
||||
| "Escape"
|
||||
| "Backspace"
|
||||
| "Delete"
|
||||
| "Tab"
|
||||
| "Up"
|
||||
| "Down"
|
||||
| "Left"
|
||||
| "Right"
|
||||
| "Home"
|
||||
| "End"
|
||||
| "PageUp"
|
||||
| "PageDown"
|
||||
| "Function";
|
||||
|
||||
export interface KeyAction {
|
||||
action: KeyActionName;
|
||||
char?: string;
|
||||
n?: number;
|
||||
}
|
||||
|
||||
/** A message sent from the TUI client to the daemon over the IPC socket. */
|
||||
export type ClientRequest =
|
||||
| { kind: "Tick" }
|
||||
| { kind: "KeyPress"; key: KeyAction; ctrl: boolean; alt: boolean; shift: boolean }
|
||||
| { kind: "Submit"; text: string }
|
||||
| { kind: "Paste"; text: string }
|
||||
| { kind: "Resize"; cols: number; rows: number }
|
||||
| { kind: "Close" }
|
||||
| { kind: "ScrollUp" }
|
||||
| { kind: "ScrollDown" };
|
||||
|
||||
/** A single chat message within a session. */
|
||||
export interface MessageEntry {
|
||||
role: string;
|
||||
content: string;
|
||||
timestamp: number;
|
||||
}
|
||||
|
||||
/** A transient toast notification sent to the client. */
|
||||
export interface ToastEntry {
|
||||
kind: string;
|
||||
message: string;
|
||||
created_at: number;
|
||||
lifetime_ms: number;
|
||||
}
|
||||
|
||||
/** Full UI state snapshot pushed from the daemon to the client. */
|
||||
export interface StatePayload {
|
||||
session_id: string;
|
||||
messages: MessageEntry[];
|
||||
edit_count: number;
|
||||
message_count: number;
|
||||
overlay: string | null;
|
||||
toasts: ToastEntry[];
|
||||
dirty: boolean;
|
||||
input_buffer: string;
|
||||
input_cursor: number;
|
||||
}
|
||||
|
||||
/** A frame sent from the daemon to the client. */
|
||||
export type DaemonFrame =
|
||||
| { kind: "StateUpdate"; payload: StatePayload }
|
||||
| { kind: "StreamToken"; token: string }
|
||||
| { kind: "SystemNote"; noteKind: string; message: string }
|
||||
| { kind: "ClipboardCopy"; text: string }
|
||||
| { kind: "Closed" };
|
||||
@@ -0,0 +1,60 @@
|
||||
/**
|
||||
* IPC server — binds a Unix socket and accepts incoming client connections.
|
||||
* Mirrors `apps/infrastructure/src/ipc/server.rs`.
|
||||
*/
|
||||
import * as net from "node:net";
|
||||
import * as fs from "node:fs";
|
||||
import { Connection } from "./conn.ts";
|
||||
|
||||
/** A Unix-socket IPC server. */
|
||||
export class IpcServer {
|
||||
private server: net.Server;
|
||||
private path: string;
|
||||
|
||||
private constructor(path: string, server: net.Server) {
|
||||
this.path = path;
|
||||
this.server = server;
|
||||
}
|
||||
|
||||
/** Bind a Unix-socket server at `path`, removing any stale socket file. */
|
||||
static bindUnix(path: string): Promise<IpcServer> {
|
||||
return new Promise((resolve, reject) => {
|
||||
if (fs.existsSync(path)) {
|
||||
fs.unlinkSync(path);
|
||||
}
|
||||
const server = net.createServer();
|
||||
server.once("error", reject);
|
||||
server.listen(path, () => {
|
||||
resolve(new IpcServer(path, server));
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
/** Accept the next incoming connection as a `Connection`. */
|
||||
accept(): Promise<Connection> {
|
||||
return new Promise((resolve) => {
|
||||
const onConnection = (socket: net.Socket) => {
|
||||
serverCleanup();
|
||||
resolve(new Connection(socket));
|
||||
};
|
||||
const serverCleanup = () => {
|
||||
this.server.removeListener("connection", onConnection);
|
||||
};
|
||||
this.server.once("connection", onConnection);
|
||||
});
|
||||
}
|
||||
|
||||
/** Close the server and remove the socket file. */
|
||||
close(): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
this.server.close(() => {
|
||||
try {
|
||||
fs.unlinkSync(this.path);
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -7,6 +7,7 @@ import * as path from "node:path";
|
||||
import type { JsonValue } from "@zesdex/domain";
|
||||
import { type Tool, type ToolCtx } from "./mod.ts";
|
||||
import { argStr } from "./util.ts";
|
||||
import { bashControl } from "../bgbash/control.ts";
|
||||
|
||||
/** Get the output of a background bash job by ID. */
|
||||
export class BashOutput implements Tool {
|
||||
@@ -59,7 +60,11 @@ export class BashKill implements Tool {
|
||||
|
||||
run(_ctx: ToolCtx, args: JsonValue): string {
|
||||
const jobId = argStr(args, "job_id");
|
||||
// Background job control is a global registry — not yet ported (bgbash 3g).
|
||||
throw new Error(`no active background job found with ID '${jobId}'`);
|
||||
// Look up the job in the global bgbash registry.
|
||||
const cancelled = bashControl().cancel(jobId);
|
||||
if (!cancelled) {
|
||||
throw new Error(`no active background job found with ID '${jobId}'`);
|
||||
}
|
||||
return `Background job '${jobId}' killed.`;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,7 +26,7 @@ export class Bash implements Tool {
|
||||
required: ["command"],
|
||||
};
|
||||
|
||||
run(_ctx: ToolCtx, args: JsonValue): string {
|
||||
async run(_ctx: ToolCtx, args: JsonValue): Promise<string> {
|
||||
const cmd = argStr(args, "command");
|
||||
let timeoutMs = optInt(args, "timeout", DEFAULT_TIMEOUT_MS);
|
||||
timeoutMs = Math.min(timeoutMs, MAX_TIMEOUT_MS);
|
||||
@@ -38,28 +38,12 @@ export class Bash implements Tool {
|
||||
const runInBackground = optBool(args, "run_in_background", false);
|
||||
|
||||
if (runInBackground) {
|
||||
// Background jobs handled via bgbash
|
||||
const jobId = `${Date.now()}-${Math.random().toString(36).slice(2, 10)}`;
|
||||
try {
|
||||
const outDir = _ctx.sessionDir ? `${_ctx.sessionDir}/bash-outputs` : "";
|
||||
if (outDir) {
|
||||
import("node:fs").then((fs) => {
|
||||
fs.mkdirSync(outDir, { recursive: true });
|
||||
const child = Bun.spawn({ cmd: ["bash", "-c", cmd], stdout: "pipe", stderr: "pipe" });
|
||||
const read = async () => {
|
||||
const out = child.stdout ? await new Response(child.stdout).text() : "";
|
||||
const err = child.stderr ? await new Response(child.stderr).text() : "";
|
||||
const code = await child.exited;
|
||||
const combined = err ? `${out}\n${err}` : out;
|
||||
fs.writeFileSync(`${outDir}/${jobId}`, `Exit code: ${code}\n\n${combined}`);
|
||||
};
|
||||
read().catch(() => {});
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
// Background spawn failed — still report the job id
|
||||
}
|
||||
return `Background job: ${jobId}`;
|
||||
// Background jobs handled via bgbash registry
|
||||
const { spawnBashJob } = await import("../../bgbash/job.ts");
|
||||
const { bashControl } = await import("../../bgbash/control.ts");
|
||||
const job = spawnBashJob(cmd);
|
||||
bashControl().register(job);
|
||||
return `Background job: ${job.id}`;
|
||||
}
|
||||
|
||||
const proc = Bun.spawnSync({ cmd: ["bash", "-c", cmd], stdout: "pipe", stderr: "pipe", timeout: timeoutMs });
|
||||
|
||||
Reference in New Issue
Block a user