feat(ws): port Fase 5 WS server — Bun native WebSocket, ZESDEX_WS_TOKEN guard, prompt turn proxy + streaming token/done/error
This commit is contained in:
@@ -0,0 +1,18 @@
|
|||||||
|
{
|
||||||
|
"name": "@zesdex/ws",
|
||||||
|
"version": "1.21.2",
|
||||||
|
"private": true,
|
||||||
|
"type": "module",
|
||||||
|
"scripts": {
|
||||||
|
"ws": "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,6 @@
|
|||||||
|
/**
|
||||||
|
* Zesdex WebSocket interface package.
|
||||||
|
* Mirrors `apps/interfaces/ws/`.
|
||||||
|
*/
|
||||||
|
export { startWsServer } from "./server.ts";
|
||||||
|
export type { WsState, SocketSend } from "./server.ts";
|
||||||
@@ -0,0 +1,9 @@
|
|||||||
|
#!/usr/bin/env bun
|
||||||
|
/**
|
||||||
|
* Zesdex WS server — standalone entry point.
|
||||||
|
* Mirrors the `--ws` mode of the Rust gateway.
|
||||||
|
*/
|
||||||
|
import { startWsServer } from "./index.ts";
|
||||||
|
|
||||||
|
const port = Number.parseInt(process.env.ZESDEX_WS_PORT ?? "8081", 10);
|
||||||
|
startWsServer(port);
|
||||||
@@ -0,0 +1,220 @@
|
|||||||
|
#!/usr/bin/env bun
|
||||||
|
/**
|
||||||
|
* Zesdex WebSocket interface — real-time bidirectional communication.
|
||||||
|
* Mirrors `apps/interfaces/ws/src/lib.rs`.
|
||||||
|
*
|
||||||
|
* Security: accepts an optional `?token=` query param. When `ZESDEX_WS_TOKEN`
|
||||||
|
* env is set, connections MUST present a matching token, else rejected —
|
||||||
|
* prevents the endpoint from being used as an open LLM proxy.
|
||||||
|
*
|
||||||
|
* Protocol (JSON text frames):
|
||||||
|
* client → { "type": "prompt", "message": "...", "model"?: "..." }
|
||||||
|
* server → { "type": "connected", "session": ..., "message": "..." } (on connect)
|
||||||
|
* server → { "type": "token", "content": "..." } (streaming)
|
||||||
|
* server → { "type": "done" } | { "type": "error", "message": "..." } (terminal)
|
||||||
|
* server → { "type": "echo", "data": "..." } (fallback)
|
||||||
|
*/
|
||||||
|
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,
|
||||||
|
type TurnEvent,
|
||||||
|
type TurnEventSink,
|
||||||
|
type AgentTurnParams,
|
||||||
|
userMessage,
|
||||||
|
resolveEffectiveModel,
|
||||||
|
} from "@zesdex/domain";
|
||||||
|
|
||||||
|
/** Shared application state for the WS server. */
|
||||||
|
export interface WsState {
|
||||||
|
store_base_dir: string;
|
||||||
|
session_id: string | null;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Minimal send interface used by streaming helpers. */
|
||||||
|
export interface SocketSend {
|
||||||
|
send(json: string): void;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Array-backed turn event sink (push + drain). */
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Validate the token against ZESDEX_WS_TOKEN (when configured). */
|
||||||
|
function tokenAllowed(url: URL): boolean {
|
||||||
|
const configured = process.env.ZESDEX_WS_TOKEN;
|
||||||
|
const t = url.searchParams.get("token");
|
||||||
|
if (configured) {
|
||||||
|
return t === configured;
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Wire a turn service for one prompt (matches the CLI composition root). */
|
||||||
|
async function buildSingleShot(
|
||||||
|
sessionDir: string,
|
||||||
|
eventSink: TurnEventSink,
|
||||||
|
): Promise<{ turnService: AgentTurnServiceImpl; params: AgentTurnParams }> {
|
||||||
|
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,
|
||||||
|
workspaces: [sessionDir],
|
||||||
|
turnEvents: eventSink,
|
||||||
|
workflowFindings: [],
|
||||||
|
} as unknown as ToolCtx;
|
||||||
|
const executor = new InfrastructureToolExecutor(toolCtx);
|
||||||
|
const turnService = new AgentTurnServiceImpl(llmClient, executor, toolDefs(allTools()) as never);
|
||||||
|
|
||||||
|
const params: AgentTurnParams = {
|
||||||
|
messages: [userMessage("")],
|
||||||
|
session_dir: sessionDir,
|
||||||
|
workspace_roots: [sessionDir],
|
||||||
|
turn_events: eventSink,
|
||||||
|
in_flight: { value: false },
|
||||||
|
abort: new AbortController(),
|
||||||
|
api_key: apiKey,
|
||||||
|
model,
|
||||||
|
api_base: baseUrl,
|
||||||
|
};
|
||||||
|
|
||||||
|
return { turnService, params };
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Start the WS server. `onPrompt` is invoked with the client's prompt payload,
|
||||||
|
* and `forward` sends raw strings to the client.
|
||||||
|
*/
|
||||||
|
export function startWsServer(port: number): { stop: () => void } {
|
||||||
|
const state: WsState = { store_base_dir: ".", session_id: null };
|
||||||
|
|
||||||
|
const server = Bun.serve({
|
||||||
|
port,
|
||||||
|
fetch(req, server) {
|
||||||
|
const url = new URL(req.url);
|
||||||
|
if (url.pathname === "/ws") {
|
||||||
|
if (!tokenAllowed(url)) {
|
||||||
|
return new Response("missing or invalid token", { status: 401 });
|
||||||
|
}
|
||||||
|
if (server.upgrade(req, {})) {
|
||||||
|
return undefined;
|
||||||
|
}
|
||||||
|
return new Response("upgrade failed", { status: 400 });
|
||||||
|
}
|
||||||
|
return new Response("Not found", { status: 404 });
|
||||||
|
},
|
||||||
|
websocket: {
|
||||||
|
open(ws) {
|
||||||
|
const sender: SocketSend = { send: (s) => ws.send(s) };
|
||||||
|
// Welcome message on connect.
|
||||||
|
sender.send(
|
||||||
|
JSON.stringify({
|
||||||
|
type: "connected",
|
||||||
|
session: state.session_id,
|
||||||
|
message: "Connected to Zesdex WebSocket server",
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
},
|
||||||
|
message(ws, msg) {
|
||||||
|
const text = typeof msg === "string" ? msg : JSON.stringify(msg);
|
||||||
|
if (!text.startsWith("{")) {
|
||||||
|
ws.send(JSON.stringify({ type: "echo", data: text }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let val: Record<string, unknown>;
|
||||||
|
try {
|
||||||
|
val = JSON.parse(text) as Record<string, unknown>;
|
||||||
|
} catch {
|
||||||
|
ws.send(JSON.stringify({ type: "echo", data: text }));
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (val.type === "prompt" && typeof val.message === "string") {
|
||||||
|
const sender: SocketSend = { send: (s) => ws.send(s) };
|
||||||
|
void runPrompt(sender, val.message, typeof val.model === "string" ? val.model : undefined);
|
||||||
|
} else {
|
||||||
|
ws.send(JSON.stringify({ type: "echo", data: text }));
|
||||||
|
}
|
||||||
|
},
|
||||||
|
close() {},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
console.log(`zesdex-ws listening on ws://0.0.0.0:${server.port}`);
|
||||||
|
return { stop: () => server.stop(true) };
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Run a single prompt and stream events to the client. */
|
||||||
|
async function runPrompt(ws: SocketSend, message: string, modelOverride?: string): Promise<void> {
|
||||||
|
const sessionDir = process.cwd();
|
||||||
|
const sink = new ArraySink();
|
||||||
|
try {
|
||||||
|
const { turnService, params } = await buildSingleShot(sessionDir, sink);
|
||||||
|
if (modelOverride) params.model = modelOverride;
|
||||||
|
params.messages = [userMessage(message)];
|
||||||
|
|
||||||
|
// Poll the sink and forward events to the client.
|
||||||
|
const poll = setInterval(() => {
|
||||||
|
const evs = sink.drain();
|
||||||
|
for (const ev of evs) {
|
||||||
|
if (ev.kind === "stream_token") {
|
||||||
|
ws.send(JSON.stringify({ type: "token", content: ev.content }));
|
||||||
|
} else if (ev.kind === "done") {
|
||||||
|
ws.send(JSON.stringify({ type: "done" }));
|
||||||
|
clearInterval(poll);
|
||||||
|
} else if (ev.kind === "error") {
|
||||||
|
ws.send(JSON.stringify({ type: "error", message: ev.message }));
|
||||||
|
clearInterval(poll);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}, 50);
|
||||||
|
|
||||||
|
try {
|
||||||
|
await turnService.runTurn(params);
|
||||||
|
// Ensure a terminal done/error is sent even if none auto-emitted
|
||||||
|
// while the poll stopped early.
|
||||||
|
const remaining = sink.drain();
|
||||||
|
for (const ev of remaining) {
|
||||||
|
if (ev.kind === "stream_token") ws.send(JSON.stringify({ type: "token", content: ev.content }));
|
||||||
|
else if (ev.kind === "done") ws.send(JSON.stringify({ type: "done" }));
|
||||||
|
else if (ev.kind === "error") ws.send(JSON.stringify({ type: "error", message: ev.message }));
|
||||||
|
}
|
||||||
|
if (!remaining.some((e) => e.kind === "done" || e.kind === "error")) {
|
||||||
|
ws.send(JSON.stringify({ type: "done" }));
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
clearInterval(poll);
|
||||||
|
}
|
||||||
|
} catch (e) {
|
||||||
|
ws.send(JSON.stringify({ type: "error", message: (e as Error).message }));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -36,6 +36,18 @@
|
|||||||
"@types/bun": "^1.2.0",
|
"@types/bun": "^1.2.0",
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
"apps/interfaces/ws": {
|
||||||
|
"name": "@zesdex/ws",
|
||||||
|
"version": "1.21.2",
|
||||||
|
"dependencies": {
|
||||||
|
"@zesdex/application": "workspace:*",
|
||||||
|
"@zesdex/domain": "workspace:*",
|
||||||
|
"@zesdex/infrastructure": "workspace:*",
|
||||||
|
},
|
||||||
|
"devDependencies": {
|
||||||
|
"@types/bun": "^1.2.0",
|
||||||
|
},
|
||||||
|
},
|
||||||
"apps/packages/application": {
|
"apps/packages/application": {
|
||||||
"name": "@zesdex/application",
|
"name": "@zesdex/application",
|
||||||
"version": "1.21.2",
|
"version": "1.21.2",
|
||||||
@@ -111,6 +123,8 @@
|
|||||||
|
|
||||||
"@zesdex/infrastructure": ["@zesdex/infrastructure@workspace:apps/packages/infrastructure"],
|
"@zesdex/infrastructure": ["@zesdex/infrastructure@workspace:apps/packages/infrastructure"],
|
||||||
|
|
||||||
|
"@zesdex/ws": ["@zesdex/ws@workspace:apps/interfaces/ws"],
|
||||||
|
|
||||||
"bun-types": ["bun-types@1.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-iIKw23BspnQQYd3prITOBxeUsxBHnwzX6YJfGMuNOZzeNcMmVqzIIVGRm1l69ogaPQmb4wB6BN8mA5bE9YuC5Q=="],
|
"bun-types": ["bun-types@1.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-iIKw23BspnQQYd3prITOBxeUsxBHnwzX6YJfGMuNOZzeNcMmVqzIIVGRm1l69ogaPQmb4wB6BN8mA5bE9YuC5Q=="],
|
||||||
|
|
||||||
"typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="],
|
"typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="],
|
||||||
|
|||||||
Reference in New Issue
Block a user