feat(daemon): port Fase 5 daemon — IPC agent-driver loop (Submit/Close) + streaming tokens; wire CLI --daemon/--attach

This commit is contained in:
asepharyana
2026-09-02 22:50:38 +07:00
parent d134678530
commit 14d3556179
9 changed files with 387 additions and 6 deletions
+2 -1
View File
@@ -15,7 +15,8 @@
"@zesdex/api": "workspace:*", "@zesdex/api": "workspace:*",
"@zesdex/ws": "workspace:*", "@zesdex/ws": "workspace:*",
"@zesdex/grpc": "workspace:*", "@zesdex/grpc": "workspace:*",
"@zesdex/web": "workspace:*" "@zesdex/web": "workspace:*",
"@zesdex/daemon": "workspace:*"
}, },
"devDependencies": { "devDependencies": {
"@types/bun": "^1.2.0" "@types/bun": "^1.2.0"
+13 -4
View File
@@ -9,6 +9,7 @@
import * as fs from "node:fs"; import * as fs from "node:fs";
import * as os from "node:os"; import * as os from "node:os";
import * as path from "node:path"; import * as path from "node:path";
import { newStore } from "@zesdex/domain";
import { parseCli, modeCount } from "./args.ts"; import { parseCli, modeCount } from "./args.ts";
import { runSingleProcess } from "./compose.ts"; import { runSingleProcess } from "./compose.ts";
@@ -127,11 +128,19 @@ async function main(): Promise<void> {
} }
if (opts.daemon) { if (opts.daemon) {
console.info("daemon mode not yet wired — falling back to REPL"); const { runDaemon } = await import("@zesdex/daemon");
await runRepl(); 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) { } else if (opts.attach !== null) {
console.info(`attach mode for session '${opts.attach}' not yet wired`); const { runAttach } = await import("@zesdex/daemon");
process.exit(1); 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) { } else if (opts.api) {
const { newApiState, startApiServer } = await import("@zesdex/api"); const { newApiState, startApiServer } = await import("@zesdex/api");
const baseDir = const baseDir =
+18
View File
@@ -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"
}
}
+61
View File
@@ -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();
}
+7
View File
@@ -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";
+20
View File
@@ -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());
+249
View File
@@ -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 */
}
},
};
}
+15
View File
@@ -30,6 +30,7 @@
"dependencies": { "dependencies": {
"@zesdex/api": "workspace:*", "@zesdex/api": "workspace:*",
"@zesdex/application": "workspace:*", "@zesdex/application": "workspace:*",
"@zesdex/daemon": "workspace:*",
"@zesdex/domain": "workspace:*", "@zesdex/domain": "workspace:*",
"@zesdex/grpc": "workspace:*", "@zesdex/grpc": "workspace:*",
"@zesdex/infrastructure": "workspace:*", "@zesdex/infrastructure": "workspace:*",
@@ -40,6 +41,18 @@
"@types/bun": "^1.2.0", "@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": { "apps/interfaces/grpc": {
"name": "@zesdex/grpc", "name": "@zesdex/grpc",
"version": "1.21.2", "version": "1.21.2",
@@ -137,6 +150,8 @@
"@zesdex/cli": ["@zesdex/cli@workspace:apps/interfaces/cli"], "@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/domain": ["@zesdex/domain@workspace:apps/packages/domain"],
"@zesdex/grpc": ["@zesdex/grpc@workspace:apps/interfaces/grpc"], "@zesdex/grpc": ["@zesdex/grpc@workspace:apps/interfaces/grpc"],
+2 -1
View File
@@ -26,7 +26,8 @@
"@zesdex/api": ["./apps/interfaces/api/src/index.ts"], "@zesdex/api": ["./apps/interfaces/api/src/index.ts"],
"@zesdex/ws": ["./apps/interfaces/ws/src/index.ts"], "@zesdex/ws": ["./apps/interfaces/ws/src/index.ts"],
"@zesdex/grpc": ["./apps/interfaces/grpc/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"] "include": ["apps/packages/*/src", "apps/interfaces/*/src", "apps/interfaces/*/bin"]