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/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"
|
||||||
|
|||||||
@@ -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 =
|
||||||
|
|||||||
@@ -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": {
|
"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
@@ -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"]
|
||||||
|
|||||||
Reference in New Issue
Block a user