feat(daemon): port Fase 5 daemon — IPC agent-driver loop (Submit/Close) + streaming tokens; wire CLI --daemon/--attach
This commit is contained in:
@@ -15,7 +15,8 @@
|
||||
"@zesdex/api": "workspace:*",
|
||||
"@zesdex/ws": "workspace:*",
|
||||
"@zesdex/grpc": "workspace:*",
|
||||
"@zesdex/web": "workspace:*"
|
||||
"@zesdex/web": "workspace:*",
|
||||
"@zesdex/daemon": "workspace:*"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/bun": "^1.2.0"
|
||||
|
||||
@@ -9,6 +9,7 @@
|
||||
import * as fs from "node:fs";
|
||||
import * as os from "node:os";
|
||||
import * as path from "node:path";
|
||||
import { newStore } from "@zesdex/domain";
|
||||
import { parseCli, modeCount } from "./args.ts";
|
||||
import { runSingleProcess } from "./compose.ts";
|
||||
|
||||
@@ -127,11 +128,19 @@ async function main(): Promise<void> {
|
||||
}
|
||||
|
||||
if (opts.daemon) {
|
||||
console.info("daemon mode not yet wired — falling back to REPL");
|
||||
await runRepl();
|
||||
const { runDaemon } = await import("@zesdex/daemon");
|
||||
console.info("daemon mode — creating per-session IPC socket");
|
||||
const handle = await runDaemon();
|
||||
const stop = () => void handle.stop().then(() => process.exit(0));
|
||||
process.on("SIGINT", stop);
|
||||
process.on("SIGTERM", stop);
|
||||
await new Promise(() => {});
|
||||
} else if (opts.attach !== null) {
|
||||
console.info(`attach mode for session '${opts.attach}' not yet wired`);
|
||||
process.exit(1);
|
||||
const { runAttach } = await import("@zesdex/daemon");
|
||||
console.info(`attach mode for daemon session '${opts.attach}'`);
|
||||
const store = newStore();
|
||||
const socketPath = path.join(store.base_dir, "run", `${opts.attach}.sock`);
|
||||
await runAttach(socketPath);
|
||||
} else if (opts.api) {
|
||||
const { newApiState, startApiServer } = await import("@zesdex/api");
|
||||
const baseDir =
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
{
|
||||
"name": "@zesdex/daemon",
|
||||
"version": "1.21.2",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"scripts": {
|
||||
"daemon": "bun src/main.ts",
|
||||
"build": "true"
|
||||
},
|
||||
"dependencies": {
|
||||
"@zesdex/domain": "workspace:*",
|
||||
"@zesdex/application": "workspace:*",
|
||||
"@zesdex/infrastructure": "workspace:*"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/bun": "^1.2.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
/**
|
||||
* Daemon client — connects to a running daemon over its Unix socket and
|
||||
* exchanges framed IPC messages. Mirrors `apps/interfaces/daemon/src/client.rs`.
|
||||
*/
|
||||
import { IpcClient } from "@zesdex/infrastructure";
|
||||
import type { DaemonFrame, ClientRequest } from "@zesdex/infrastructure";
|
||||
|
||||
/**
|
||||
* Connect to a running daemon socket and send one or more requests,
|
||||
* printing the frames received back. Returns when the daemon closes or
|
||||
* after a timeout.
|
||||
*/
|
||||
export async function runAttach(
|
||||
socketPath: string,
|
||||
initialText?: string,
|
||||
timeoutMs = 15_000,
|
||||
): Promise<void> {
|
||||
const client = await IpcClient.connectUnix(socketPath);
|
||||
|
||||
// Send a submit if provided, so the daemon starts a turn.
|
||||
if (initialText) {
|
||||
const req: ClientRequest = { kind: "Submit", text: initialText };
|
||||
client.send(req);
|
||||
}
|
||||
|
||||
// Read frames until closed or timeout.
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (Date.now() < deadline) {
|
||||
const frame = await client.receive<DaemonFrame>();
|
||||
if (frame === null) break; // daemon closed connection
|
||||
switch (frame.kind) {
|
||||
case "StreamToken":
|
||||
process.stdout.write(frame.token);
|
||||
break;
|
||||
case "StateUpdate":
|
||||
process.stdout.write(`\n[state: ${frame.payload.message_count} msgs]\n`);
|
||||
break;
|
||||
case "Closed":
|
||||
process.stdout.write("\n[closed]\n");
|
||||
client.close();
|
||||
return;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
client.close();
|
||||
}
|
||||
|
||||
/** Connect and send a Close request, for shutdown testing. */
|
||||
export async function sendClose(socketPath: string): Promise<void> {
|
||||
const client = await IpcClient.connectUnix(socketPath);
|
||||
const req: ClientRequest = { kind: "Close" };
|
||||
client.send(req);
|
||||
// Give the daemon a moment to reply Closed.
|
||||
const frame = await client.receive<DaemonFrame>();
|
||||
if (frame && frame.kind === "Closed") {
|
||||
process.stdout.write("[closed]\n");
|
||||
}
|
||||
client.close();
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
/**
|
||||
* Zesdex daemon interface package.
|
||||
* Mirrors `apps/interfaces/daemon/`.
|
||||
*/
|
||||
export { runDaemon } from "./server.ts";
|
||||
export type { DaemonHandle } from "./server.ts";
|
||||
export { runAttach, sendClose } from "./client.ts";
|
||||
@@ -0,0 +1,20 @@
|
||||
#!/usr/bin/env bun
|
||||
/**
|
||||
* Zesdex daemon — standalone entry point.
|
||||
* Mirrors the `--daemon` mode of the Rust gateway.
|
||||
*/
|
||||
import { runDaemon } from "./index.ts";
|
||||
|
||||
const handle = await runDaemon();
|
||||
console.log("daemon running — Ctrl+C to stop");
|
||||
|
||||
// Keep alive until SIGINT/SIGTERM.
|
||||
let stopping = false;
|
||||
async function shutdown(): Promise<void> {
|
||||
if (stopping) return;
|
||||
stopping = true;
|
||||
await handle.stop();
|
||||
process.exit(0);
|
||||
}
|
||||
process.on("SIGINT", () => void shutdown());
|
||||
process.on("SIGTERM", () => void shutdown());
|
||||
@@ -0,0 +1,249 @@
|
||||
/**
|
||||
* Daemon server — owns the agent state, listens on a per-session Unix socket,
|
||||
* and drives requests from an attached client.
|
||||
* Mirrors `apps/interfaces/daemon/src/server.rs` (simplified: focuses on the
|
||||
* IPC agent-driver loop; full TUI key-handling lives in the Fase 4 TUI).
|
||||
*
|
||||
* Flow: `runDaemon()` creates a session → binds a Unix socket under
|
||||
* `<store>/run/<session_id>.sock` → accepts a client → loops reading
|
||||
* `ClientRequest`s. On `Submit` it runs an agent turn and streams
|
||||
* `StreamToken`/`StateUpdate` frames back. On `Close`/disconnect it cleans
|
||||
* up the socket file.
|
||||
*/
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { IpcServer } from "@zesdex/infrastructure";
|
||||
import {
|
||||
JsonSettingsRepository,
|
||||
JsonAppConfigRepository,
|
||||
LlmClient,
|
||||
resolveApiKey,
|
||||
InfrastructureToolExecutor,
|
||||
allTools,
|
||||
toolDefs,
|
||||
} from "@zesdex/infrastructure";
|
||||
import type { ToolCtx } from "@zesdex/infrastructure";
|
||||
import { AgentTurnServiceImpl } from "@zesdex/application";
|
||||
import {
|
||||
newStore,
|
||||
ensureStoreDirs,
|
||||
newSessionId,
|
||||
resolveEffectiveModel,
|
||||
type TurnEvent,
|
||||
type TurnEventSink,
|
||||
type AgentTurnParams,
|
||||
userMessage,
|
||||
} from "@zesdex/domain";
|
||||
import type { ClientRequest, DaemonFrame, MessageEntry, StatePayload } from "@zesdex/infrastructure";
|
||||
|
||||
/** A running daemon handle. */
|
||||
export interface DaemonHandle {
|
||||
stop(): Promise<void>;
|
||||
}
|
||||
|
||||
/** Array-backed turn event sink (push + drain). */
|
||||
export class ArraySink implements TurnEventSink {
|
||||
constructor(public events: TurnEvent[] = []) {}
|
||||
push(event: TurnEvent): void {
|
||||
this.events.push(event);
|
||||
}
|
||||
drain(): TurnEvent[] {
|
||||
const out = this.events;
|
||||
this.events = [];
|
||||
return out;
|
||||
}
|
||||
}
|
||||
|
||||
/** Wire a fresh agent runtime (mirrors the CLI composition root). */
|
||||
async function buildRuntime() {
|
||||
const store = newStore();
|
||||
await ensureStoreDirs(store);
|
||||
|
||||
const settingsRepo = new JsonSettingsRepository();
|
||||
const appConfigRepo = new JsonAppConfigRepository();
|
||||
const settings = await settingsRepo.load(store.base_dir);
|
||||
const appConfig = await appConfigRepo.load(store.base_dir);
|
||||
|
||||
const provider = settings.provider;
|
||||
const model = resolveEffectiveModel(settings, appConfig);
|
||||
const baseUrl = appConfig.providers[provider]?.api_base ?? undefined;
|
||||
const apiKey = resolveApiKey(settings, appConfig);
|
||||
|
||||
const llmClient = new LlmClient(apiKey, model, baseUrl);
|
||||
const toolCtx = {
|
||||
sessionDir: store.base_dir,
|
||||
workspaces: [process.cwd()],
|
||||
turnEvents: new ArraySink(),
|
||||
workflowFindings: [],
|
||||
} as unknown as ToolCtx;
|
||||
const executor = new InfrastructureToolExecutor(toolCtx);
|
||||
const turnService = new AgentTurnServiceImpl(llmClient, executor, toolDefs(allTools()) as never);
|
||||
|
||||
const buildParams = (message: string, sink: TurnEventSink): AgentTurnParams => ({
|
||||
messages: [userMessage(message)],
|
||||
session_dir: store.base_dir,
|
||||
workspace_roots: [process.cwd()],
|
||||
turn_events: sink,
|
||||
in_flight: { value: false },
|
||||
abort: new AbortController(),
|
||||
api_key: apiKey,
|
||||
model,
|
||||
api_base: baseUrl,
|
||||
});
|
||||
|
||||
return { store, turnService, buildParams };
|
||||
}
|
||||
|
||||
/** Build a StatePayload snapshot from the accumulated messages. */
|
||||
function snapshot(sessionId: string, messages: MessageEntry[], dirty: boolean): StatePayload {
|
||||
return {
|
||||
session_id: sessionId,
|
||||
messages,
|
||||
edit_count: 0,
|
||||
message_count: messages.length,
|
||||
overlay: null,
|
||||
toasts: [],
|
||||
dirty,
|
||||
input_buffer: "",
|
||||
input_cursor: 0,
|
||||
};
|
||||
}
|
||||
|
||||
/** Handle one attached client: process its requests, driving agent turns. */
|
||||
async function handleDaemonClient(
|
||||
conn: { send<T>(m: T): void; receiveJson<T>(): Promise<T | null>; close(): void },
|
||||
sessionId: string,
|
||||
): Promise<void> {
|
||||
const runtime = await buildRuntime();
|
||||
const messages: MessageEntry[] = [];
|
||||
|
||||
const pushState = (dirty: boolean) =>
|
||||
conn.send<DaemonFrame>({
|
||||
kind: "StateUpdate",
|
||||
payload: snapshot(sessionId, messages, dirty),
|
||||
});
|
||||
|
||||
pushState(false);
|
||||
|
||||
let running = true;
|
||||
while (running) {
|
||||
let req: ClientRequest | null = null;
|
||||
try {
|
||||
req = await conn.receiveJson<ClientRequest>();
|
||||
} catch {
|
||||
break;
|
||||
}
|
||||
if (req === null) break; // client closed
|
||||
|
||||
switch (req.kind) {
|
||||
case "Submit": {
|
||||
const text = req.text;
|
||||
if (text.trim() === "") break;
|
||||
messages.push({ role: "user", content: text, timestamp: Date.now() });
|
||||
pushState(true);
|
||||
|
||||
// Run the turn and stream tokens as they arrive.
|
||||
const sink = new ArraySink();
|
||||
const params = runtime.buildParams(text, sink);
|
||||
void runTurnAndStream(conn, sink, runtime.turnService, params).then(() => {
|
||||
messages.push({
|
||||
role: "assistant",
|
||||
content: sink.events
|
||||
.filter((e) => e.kind === "stream_token")
|
||||
.map((e) => (e as { content: string }).content)
|
||||
.join(""),
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
pushState(false);
|
||||
});
|
||||
break;
|
||||
}
|
||||
case "Close":
|
||||
running = false;
|
||||
conn.send<DaemonFrame>({ kind: "Closed" });
|
||||
break;
|
||||
default:
|
||||
// KeyPress/Tick/etc. no-op in this simplified daemon.
|
||||
break;
|
||||
}
|
||||
}
|
||||
conn.close();
|
||||
}
|
||||
|
||||
/** Run a turn, streaming token frames to the client as they are produced. */
|
||||
async function runTurnAndStream(
|
||||
conn: { send<T>(m: T): void },
|
||||
sink: ArraySink,
|
||||
turnService: AgentTurnServiceImpl,
|
||||
params: AgentTurnParams,
|
||||
): Promise<void> {
|
||||
const poll = setInterval(() => {
|
||||
for (const ev of sink.drain()) {
|
||||
if (ev.kind === "stream_token") {
|
||||
conn.send<DaemonFrame>({ kind: "StreamToken", token: (ev as { content: string }).content });
|
||||
}
|
||||
}
|
||||
}, 50);
|
||||
|
||||
try {
|
||||
await turnService.runTurn(params);
|
||||
} finally {
|
||||
clearInterval(poll);
|
||||
// Drain any remaining tokens after the turn completes.
|
||||
for (const ev of sink.drain()) {
|
||||
if (ev.kind === "stream_token") {
|
||||
conn.send<DaemonFrame>({ kind: "StreamToken", token: (ev as { content: string }).content });
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Run zesdex as a background daemon owning agent state over a Unix socket. */
|
||||
export async function runDaemon(): Promise<DaemonHandle> {
|
||||
const store = newStore();
|
||||
await ensureStoreDirs(store);
|
||||
|
||||
const idRes = newSessionId(randomUUID());
|
||||
const sessionId = idRes.ok ? idRes.value : `sess-${Date.now()}`;
|
||||
fs.mkdirSync(path.join(store.base_dir, "sessions", sessionId), { recursive: true });
|
||||
|
||||
const runDir = path.join(store.base_dir, "run");
|
||||
fs.mkdirSync(runDir, { recursive: true });
|
||||
const socketPath = path.join(runDir, `${sessionId}.sock`);
|
||||
|
||||
const server = await IpcServer.bindUnix(socketPath);
|
||||
console.log(`zesdex-daemon listening on ${socketPath}`);
|
||||
|
||||
// Accept clients serially.
|
||||
void (async () => {
|
||||
for (;;) {
|
||||
try {
|
||||
const conn = await server.accept();
|
||||
console.log("daemon client connected");
|
||||
try {
|
||||
await handleDaemonClient(
|
||||
{ send: (m) => conn.send(m), receiveJson: () => conn.receiveJson(), close: () => conn.close() },
|
||||
sessionId,
|
||||
);
|
||||
} catch (e) {
|
||||
console.error("daemon client error:", e);
|
||||
}
|
||||
console.log("daemon client disconnected");
|
||||
} catch {
|
||||
break; // server closed
|
||||
}
|
||||
}
|
||||
})();
|
||||
|
||||
return {
|
||||
async stop() {
|
||||
await server.close();
|
||||
try {
|
||||
fs.unlinkSync(socketPath);
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -30,6 +30,7 @@
|
||||
"dependencies": {
|
||||
"@zesdex/api": "workspace:*",
|
||||
"@zesdex/application": "workspace:*",
|
||||
"@zesdex/daemon": "workspace:*",
|
||||
"@zesdex/domain": "workspace:*",
|
||||
"@zesdex/grpc": "workspace:*",
|
||||
"@zesdex/infrastructure": "workspace:*",
|
||||
@@ -40,6 +41,18 @@
|
||||
"@types/bun": "^1.2.0",
|
||||
},
|
||||
},
|
||||
"apps/interfaces/daemon": {
|
||||
"name": "@zesdex/daemon",
|
||||
"version": "1.21.2",
|
||||
"dependencies": {
|
||||
"@zesdex/application": "workspace:*",
|
||||
"@zesdex/domain": "workspace:*",
|
||||
"@zesdex/infrastructure": "workspace:*",
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/bun": "^1.2.0",
|
||||
},
|
||||
},
|
||||
"apps/interfaces/grpc": {
|
||||
"name": "@zesdex/grpc",
|
||||
"version": "1.21.2",
|
||||
@@ -137,6 +150,8 @@
|
||||
|
||||
"@zesdex/cli": ["@zesdex/cli@workspace:apps/interfaces/cli"],
|
||||
|
||||
"@zesdex/daemon": ["@zesdex/daemon@workspace:apps/interfaces/daemon"],
|
||||
|
||||
"@zesdex/domain": ["@zesdex/domain@workspace:apps/packages/domain"],
|
||||
|
||||
"@zesdex/grpc": ["@zesdex/grpc@workspace:apps/interfaces/grpc"],
|
||||
|
||||
+2
-1
@@ -26,7 +26,8 @@
|
||||
"@zesdex/api": ["./apps/interfaces/api/src/index.ts"],
|
||||
"@zesdex/ws": ["./apps/interfaces/ws/src/index.ts"],
|
||||
"@zesdex/grpc": ["./apps/interfaces/grpc/src/index.ts"],
|
||||
"@zesdex/web": ["./apps/interfaces/web/src/index.ts"]
|
||||
"@zesdex/web": ["./apps/interfaces/web/src/index.ts"],
|
||||
"@zesdex/daemon": ["./apps/interfaces/daemon/src/index.ts"]
|
||||
}
|
||||
},
|
||||
"include": ["apps/packages/*/src", "apps/interfaces/*/src", "apps/interfaces/*/bin"]
|
||||
|
||||
Reference in New Issue
Block a user