refactor(agents): split work

This commit is contained in:
Nicolas Meienberger 2026-04-09 18:21:13 +02:00
parent 6cda41df7c
commit 9b0d307620
No known key found for this signature in database
6 changed files with 332 additions and 321 deletions

View file

@ -2,7 +2,7 @@ import { Effect, Exit, Scope } from "effect";
import { expect, test, vi } from "vitest"; import { expect, test, vi } from "vitest";
import { fromPartial } from "@total-typescript/shoehorn"; import { fromPartial } from "@total-typescript/shoehorn";
import { createAgentMessage } from "@zerobyte/contracts/agent-protocol"; import { createAgentMessage } from "@zerobyte/contracts/agent-protocol";
import { createControllerAgentSession } from "../controller-agent-session"; import { createControllerAgentSession } from "../controller/session";
const createSocket = () => { const createSocket = () => {
return fromPartial<Parameters<typeof createControllerAgentSession>[0]>({ return fromPartial<Parameters<typeof createControllerAgentSession>[0]>({
@ -50,8 +50,7 @@ test("close emits a synthetic backup.cancelled for a started backup", () => {
expect(onBackupCancelled).toHaveBeenCalledWith({ expect(onBackupCancelled).toHaveBeenCalledWith({
jobId: "job-1", jobId: "job-1",
scheduleId: "schedule-1", scheduleId: "schedule-1",
message: message: "The connection to the backup agent was lost. Restart the backup to ensure it completes.",
"The connection to the backup agent was lost while this backup was running. Restart the backup to ensure it completes.",
}); });
}); });
@ -142,7 +141,6 @@ test("close emits a synthetic backup.cancelled for a queued backup", () => {
expect(onBackupCancelled).toHaveBeenCalledWith({ expect(onBackupCancelled).toHaveBeenCalledWith({
jobId: "job-queued", jobId: "job-queued",
scheduleId: "schedule-queued", scheduleId: "schedule-queued",
message: message: "The connection to the backup agent was lost. Restart the backup to ensure it completes.",
"The connection to the backup agent was lost before this backup started. Restart the backup to ensure it completes.",
}); });
}); });

View file

@ -1,56 +1,14 @@
import { type ChildProcess, spawn } from "node:child_process"; import type { ChildProcess } from "node:child_process";
import { existsSync } from "node:fs"; import { createAgentManagerRuntime, type AgentManagerRuntime } from "./controller/server";
import path from "node:path"; import { spawnLocalAgentProcess, stopLocalAgentProcess } from "./local/process";
import { Effect, Exit, Fiber, Scope } from "effect";
import { logger } from "@zerobyte/core/node";
import type {
BackupCancelPayload,
BackupCancelledPayload,
BackupCompletedPayload,
BackupFailedPayload,
BackupProgressPayload,
BackupRunPayload,
BackupStartedPayload,
} from "@zerobyte/contracts/agent-protocol";
import { config } from "../../core/config";
import { validateAgentToken, deriveLocalAgentToken } from "./agent-tokens";
import {
createControllerAgentSession,
type AgentConnectionData,
type ControllerAgentSession,
} from "./controller-agent-session";
type AgentBackupEventContext = { export type { AgentBackupEventHandlers } from "./controller/server";
agentId: string;
agentName: string;
payload:
| BackupStartedPayload
| BackupProgressPayload
| BackupCompletedPayload
| BackupFailedPayload
| BackupCancelledPayload;
};
export type AgentBackupEventHandlers = {
onBackupStarted?: (context: AgentBackupEventContext & { payload: BackupStartedPayload }) => void;
onBackupProgress?: (context: AgentBackupEventContext & { payload: BackupProgressPayload }) => void;
onBackupCompleted?: (context: AgentBackupEventContext & { payload: BackupCompletedPayload }) => void;
onBackupFailed?: (context: AgentBackupEventContext & { payload: BackupFailedPayload }) => void;
onBackupCancelled?: (context: AgentBackupEventContext & { payload: BackupCancelledPayload }) => void;
};
type AgentManagerRuntime = ReturnType<typeof createAgentManagerRuntime>;
type AgentRuntimeState = { type AgentRuntimeState = {
agentManager: AgentManagerRuntime; agentManager: AgentManagerRuntime;
localAgent: ChildProcess | null; localAgent: ChildProcess | null;
}; };
type ControllerAgentSessionHandle = {
session: ControllerAgentSession;
runFiber: Fiber.RuntimeFiber<void, never>;
scope: Scope.CloseableScope;
};
type ProcessWithAgentRuntime = NodeJS.Process & { type ProcessWithAgentRuntime = NodeJS.Process & {
__zerobyteAgentRuntime?: AgentRuntimeState; __zerobyteAgentRuntime?: AgentRuntimeState;
}; };
@ -75,278 +33,11 @@ const getAgentRuntimeState = () => {
const getAgentManagerRuntime = () => getAgentRuntimeState().agentManager; const getAgentManagerRuntime = () => getAgentRuntimeState().agentManager;
export const spawnLocalAgent = async () => { export const spawnLocalAgent = async () => {
await stopLocalAgent(); await spawnLocalAgentProcess(getAgentRuntimeState());
const sourceEntryPoint = path.join(process.cwd(), "apps", "agent", "src", "index.ts");
const productionEntryPoint = path.join(process.cwd(), ".output", "agent", "index.mjs");
if (config.__prod__ && !existsSync(productionEntryPoint)) {
throw new Error(`Local agent entrypoint not found at ${productionEntryPoint}`);
}
const agentEntryPoint = config.__prod__ ? productionEntryPoint : sourceEntryPoint;
const agentToken = await deriveLocalAgentToken();
const args = config.__prod__ ? ["run", agentEntryPoint] : ["run", "--watch", agentEntryPoint];
const runtime = getAgentRuntimeState();
const agentProcess = spawn("bun", args, {
env: {
PATH: process.env.PATH,
ZEROBYTE_CONTROLLER_URL: "ws://localhost:3001",
ZEROBYTE_AGENT_TOKEN: agentToken,
},
stdio: ["ignore", "pipe", "pipe"],
});
runtime.localAgent = agentProcess;
agentProcess.stdout?.on("data", (data: Buffer) => {
const line = data.toString().trim();
if (line) logger.info(`[agent] ${line}`);
});
agentProcess.stderr?.on("data", (data: Buffer) => {
const line = data.toString().trim();
if (line) logger.error(`[agent] ${line}`);
});
agentProcess.on("exit", (code, signal) => {
if (runtime.localAgent === agentProcess) {
runtime.localAgent = null;
}
logger.info(`Agent process exited with code ${code} and signal ${signal}`);
});
}; };
export const stopLocalAgent = async () => { export const stopLocalAgent = async () => {
const runtime = getAgentRuntimeState(); await stopLocalAgentProcess(getAgentRuntimeState());
if (!runtime.localAgent) {
return;
}
const agentProcess = runtime.localAgent;
runtime.localAgent = null;
if (agentProcess.exitCode !== null || agentProcess.signalCode !== null) {
return;
}
const exited = new Promise<void>((resolve) => {
agentProcess.once("exit", () => {
resolve();
});
});
agentProcess.kill();
await exited;
};
const createAgentManagerRuntime = () => {
let sessions = new Map<string, ControllerAgentSessionHandle>();
let backupHandlers: AgentBackupEventHandlers = {};
let runtimeScope: Scope.CloseableScope | null = null;
const closeSession = (sessionHandle: ControllerAgentSessionHandle) => {
Effect.runSync(Fiber.interrupt(sessionHandle.runFiber));
Effect.runSync(Scope.close(sessionHandle.scope, Exit.succeed(undefined)));
};
const closeAllSessions = () => {
const currentSessions = sessions;
sessions = new Map();
for (const sessionHandle of currentSessions.values()) {
closeSession(sessionHandle);
}
};
const getSessionHandle = (agentId: string) => sessions.get(agentId);
const getSession = (agentId: string) => getSessionHandle(agentId)?.session;
const createSession = (ws: Bun.ServerWebSocket<AgentConnectionData>) => {
// Manual scope management because we are out of Effect
const scope = Effect.runSync(Scope.make());
try {
const session = Effect.runSync(Scope.extend(createControllerAgentSession(ws, createSessionHandlers(ws)), scope));
const runFiber = Effect.runFork(Scope.extend(session.run, scope));
return { session, runFiber, scope };
} catch (error) {
Effect.runSync(Scope.close(scope, Exit.fail(error)));
throw error;
}
};
const setSession = (agentId: string, sessionHandle: ControllerAgentSessionHandle) => {
const existingSession = getSessionHandle(agentId);
if (existingSession) {
closeSession(existingSession);
}
sessions.set(agentId, sessionHandle);
};
const removeSession = (agentId: string, connectionId: string) => {
const sessionHandle = getSessionHandle(agentId);
if (!sessionHandle || sessionHandle.session.connectionId !== connectionId) {
return;
}
sessions.delete(agentId);
closeSession(sessionHandle);
};
const createSessionHandlers = (ws: Bun.ServerWebSocket<AgentConnectionData>) => {
const agentId = ws.data.agentId;
const agentName = ws.data.agentName;
return {
onBackupStarted: (payload: BackupStartedPayload) => {
backupHandlers.onBackupStarted?.({ agentId, agentName, payload });
},
onBackupProgress: (payload: BackupProgressPayload) => {
backupHandlers.onBackupProgress?.({ agentId, agentName, payload });
},
onBackupCompleted: (payload: BackupCompletedPayload) => {
backupHandlers.onBackupCompleted?.({ agentId, agentName, payload });
},
onBackupFailed: (payload: BackupFailedPayload) => {
backupHandlers.onBackupFailed?.({ agentId, agentName, payload });
},
onBackupCancelled: (payload: BackupCancelledPayload) => {
backupHandlers.onBackupCancelled?.({ agentId, agentName, payload });
},
};
};
const acquireServer = Effect.acquireRelease(
Effect.sync(() =>
Bun.serve<AgentConnectionData>({
port: 3001,
async fetch(req, srv) {
const url = new URL(req.url);
const token = url.searchParams.get("token");
if (!token) {
return new Response("Missing token", { status: 401 });
}
const result = await validateAgentToken(token);
if (!result) {
return new Response("Invalid or revoked token", { status: 401 });
}
const upgraded = srv.upgrade(req, {
data: {
id: Bun.randomUUIDv7(),
agentId: result.agentId,
organizationId: result.organizationId,
agentName: result.agentName,
},
});
if (upgraded) return undefined;
return new Response("WebSocket upgrade failed", { status: 400 });
},
websocket: {
open: (ws) => {
setSession(ws.data.agentId, createSession(ws));
logger.info(`Agent "${ws.data.agentName}" (${ws.data.agentId}) connected on ${ws.data.id}`);
},
message: (ws, data) => {
if (typeof data !== "string") {
logger.warn(`Ignoring non-text message from agent ${ws.data.agentId}`);
return;
}
const session = getSession(ws.data.agentId);
if (!session || session.connectionId !== ws.data.id) {
logger.warn(`No active session for agent ${ws.data.agentId} on ${ws.data.id}`);
return;
}
Effect.runSync(session.handleMessage(data));
},
close: (ws) => {
removeSession(ws.data.agentId, ws.data.id);
logger.info(`Agent "${ws.data.agentName}" (${ws.data.agentId}) disconnected`);
},
},
}),
),
(server) =>
Effect.sync(() => {
closeAllSessions();
void server.stop(true);
}),
);
const stop = () => {
if (!runtimeScope) {
return;
}
logger.info("Stopping Agent Manager...");
const scope = runtimeScope;
runtimeScope = null;
Effect.runSync(Scope.close(scope, Exit.succeed(undefined)));
};
const start = () => {
if (runtimeScope) {
stop();
}
logger.info("Starting Agent Manager...");
const scope = Effect.runSync(Scope.make());
try {
const server = Effect.runSync(Scope.extend(acquireServer, scope));
runtimeScope = scope;
logger.info(`Agent Manager listening on port ${server.port}`);
} catch (error) {
Effect.runSync(Scope.close(scope, Exit.fail(error)));
throw error;
}
};
return {
start,
sendBackup: (agentId: string, payload: BackupRunPayload) => {
const session = getSession(agentId);
if (!session) {
logger.warn(`Cannot send backup command. Agent ${agentId} is not connected.`);
return false;
}
if (!Effect.runSync(session.isReady())) {
logger.warn(`Cannot send backup command. Agent ${agentId} is not ready.`);
return false;
}
Effect.runSync(session.sendBackup(payload));
logger.info(`Sent backup command ${payload.jobId} to agent ${agentId} for schedule ${payload.scheduleId}`);
return true;
},
cancelBackup: (agentId: string, payload: BackupCancelPayload) => {
const session = getSession(agentId);
if (!session) {
logger.warn(`Cannot cancel backup command. Agent ${agentId} is not connected.`);
return false;
}
Effect.runSync(session.sendBackupCancel(payload));
logger.info(`Sent backup cancel for command ${payload.jobId} to agent ${agentId}`);
return true;
},
setBackupEventHandlers: (handlers: AgentBackupEventHandlers) => {
backupHandlers = handlers;
},
getBackupEventHandlers: () => backupHandlers,
stop,
};
}; };
export const agentManager = getAgentManagerRuntime(); export const agentManager = getAgentManagerRuntime();

View file

@ -0,0 +1,248 @@
import { Effect, Exit, Fiber, Scope } from "effect";
import { logger } from "@zerobyte/core/node";
import type {
BackupCancelPayload,
BackupCancelledPayload,
BackupCompletedPayload,
BackupFailedPayload,
BackupProgressPayload,
BackupRunPayload,
BackupStartedPayload,
} from "@zerobyte/contracts/agent-protocol";
import { createControllerAgentSession, type AgentConnectionData, type ControllerAgentSession } from "./session";
import { validateAgentToken } from "../helpers/tokens";
type AgentBackupEventContext = {
agentId: string;
agentName: string;
payload:
| BackupStartedPayload
| BackupProgressPayload
| BackupCompletedPayload
| BackupFailedPayload
| BackupCancelledPayload;
};
export type AgentBackupEventHandlers = {
onBackupStarted?: (context: AgentBackupEventContext & { payload: BackupStartedPayload }) => void;
onBackupProgress?: (context: AgentBackupEventContext & { payload: BackupProgressPayload }) => void;
onBackupCompleted?: (context: AgentBackupEventContext & { payload: BackupCompletedPayload }) => void;
onBackupFailed?: (context: AgentBackupEventContext & { payload: BackupFailedPayload }) => void;
onBackupCancelled?: (context: AgentBackupEventContext & { payload: BackupCancelledPayload }) => void;
};
type ControllerAgentSessionHandle = {
session: ControllerAgentSession;
runFiber: Fiber.RuntimeFiber<void, never>;
scope: Scope.CloseableScope;
};
export function createAgentManagerRuntime() {
let sessions = new Map<string, ControllerAgentSessionHandle>();
let backupHandlers: AgentBackupEventHandlers = {};
let runtimeScope: Scope.CloseableScope | null = null;
const closeSession = (sessionHandle: ControllerAgentSessionHandle) => {
Effect.runSync(Fiber.interrupt(sessionHandle.runFiber));
Effect.runSync(Scope.close(sessionHandle.scope, Exit.succeed(undefined)));
};
const closeAllSessions = () => {
const currentSessions = sessions;
sessions = new Map();
for (const sessionHandle of currentSessions.values()) {
closeSession(sessionHandle);
}
};
const getSessionHandle = (agentId: string) => sessions.get(agentId);
const getSession = (agentId: string) => getSessionHandle(agentId)?.session;
const createSessionHandlers = (ws: Bun.ServerWebSocket<AgentConnectionData>) => {
const agentId = ws.data.agentId;
const agentName = ws.data.agentName;
return {
onBackupStarted: (payload: BackupStartedPayload) => {
backupHandlers.onBackupStarted?.({ agentId, agentName, payload });
},
onBackupProgress: (payload: BackupProgressPayload) => {
backupHandlers.onBackupProgress?.({ agentId, agentName, payload });
},
onBackupCompleted: (payload: BackupCompletedPayload) => {
backupHandlers.onBackupCompleted?.({ agentId, agentName, payload });
},
onBackupFailed: (payload: BackupFailedPayload) => {
backupHandlers.onBackupFailed?.({ agentId, agentName, payload });
},
onBackupCancelled: (payload: BackupCancelledPayload) => {
backupHandlers.onBackupCancelled?.({ agentId, agentName, payload });
},
};
};
const createSession = (ws: Bun.ServerWebSocket<AgentConnectionData>) => {
// Manual scope management because we are out of Effect
const scope = Effect.runSync(Scope.make());
try {
const session = Effect.runSync(Scope.extend(createControllerAgentSession(ws, createSessionHandlers(ws)), scope));
const runFiber = Effect.runFork(Scope.extend(session.run, scope));
return { session, runFiber, scope };
} catch (error) {
Effect.runSync(Scope.close(scope, Exit.fail(error)));
throw error;
}
};
const setSession = (agentId: string, sessionHandle: ControllerAgentSessionHandle) => {
const existingSession = getSessionHandle(agentId);
if (existingSession) {
closeSession(existingSession);
}
sessions.set(agentId, sessionHandle);
};
const removeSession = (agentId: string, connectionId: string) => {
const sessionHandle = getSessionHandle(agentId);
if (!sessionHandle || sessionHandle.session.connectionId !== connectionId) {
return;
}
sessions.delete(agentId);
closeSession(sessionHandle);
};
const acquireServer = Effect.acquireRelease(
Effect.sync(() =>
Bun.serve<AgentConnectionData>({
port: 3001,
async fetch(req, srv) {
const url = new URL(req.url);
const token = url.searchParams.get("token");
if (!token) {
return new Response("Missing token", { status: 401 });
}
const result = await validateAgentToken(token);
if (!result) {
return new Response("Invalid or revoked token", { status: 401 });
}
const upgraded = srv.upgrade(req, {
data: {
id: Bun.randomUUIDv7(),
agentId: result.agentId,
organizationId: result.organizationId,
agentName: result.agentName,
},
});
if (upgraded) return undefined;
return new Response("WebSocket upgrade failed", { status: 400 });
},
websocket: {
open: (ws) => {
setSession(ws.data.agentId, createSession(ws));
logger.info(`Agent "${ws.data.agentName}" (${ws.data.agentId}) connected on ${ws.data.id}`);
},
message: (ws, data) => {
if (typeof data !== "string") {
logger.warn(`Ignoring non-text message from agent ${ws.data.agentId}`);
return;
}
const session = getSession(ws.data.agentId);
if (!session || session.connectionId !== ws.data.id) {
logger.warn(`No active session for agent ${ws.data.agentId} on ${ws.data.id}`);
return;
}
Effect.runSync(session.handleMessage(data));
},
close: (ws) => {
removeSession(ws.data.agentId, ws.data.id);
logger.info(`Agent "${ws.data.agentName}" (${ws.data.agentId}) disconnected`);
},
},
}),
),
(server) =>
Effect.sync(() => {
closeAllSessions();
void server.stop(true);
}),
);
const stop = () => {
if (!runtimeScope) {
return;
}
logger.info("Stopping Agent Manager...");
const scope = runtimeScope;
runtimeScope = null;
Effect.runSync(Scope.close(scope, Exit.succeed(undefined)));
};
const start = () => {
if (runtimeScope) {
stop();
}
logger.info("Starting Agent Manager...");
const scope = Effect.runSync(Scope.make());
try {
const server = Effect.runSync(Scope.extend(acquireServer, scope));
runtimeScope = scope;
logger.info(`Agent Manager listening on port ${server.port}`);
} catch (error) {
Effect.runSync(Scope.close(scope, Exit.fail(error)));
throw error;
}
};
return {
start,
sendBackup: (agentId: string, payload: BackupRunPayload) => {
const session = getSession(agentId);
if (!session) {
logger.warn(`Cannot send backup command. Agent ${agentId} is not connected.`);
return false;
}
if (!Effect.runSync(session.isReady())) {
logger.warn(`Cannot send backup command. Agent ${agentId} is not ready.`);
return false;
}
Effect.runSync(session.sendBackup(payload));
logger.info(`Sent backup command ${payload.jobId} to agent ${agentId} for schedule ${payload.scheduleId}`);
return true;
},
cancelBackup: (agentId: string, payload: BackupCancelPayload) => {
const session = getSession(agentId);
if (!session) {
logger.warn(`Cannot cancel backup command. Agent ${agentId} is not connected.`);
return false;
}
Effect.runSync(session.sendBackupCancel(payload));
logger.info(`Sent backup cancel for command ${payload.jobId} to agent ${agentId}`);
return true;
},
setBackupEventHandlers: (handlers: AgentBackupEventHandlers) => {
backupHandlers = handlers;
},
getBackupEventHandlers: () => backupHandlers,
stop,
};
}
export type AgentManagerRuntime = ReturnType<typeof createAgentManagerRuntime>;

View file

@ -3,12 +3,12 @@ import {
createControllerMessage, createControllerMessage,
parseAgentMessage, parseAgentMessage,
type AgentMessage, type AgentMessage,
type BackupCancelPayload,
type BackupCancelledPayload, type BackupCancelledPayload,
type BackupCompletedPayload, type BackupCompletedPayload,
type BackupFailedPayload, type BackupFailedPayload,
type BackupProgressPayload, type BackupProgressPayload,
type BackupRunPayload, type BackupRunPayload,
type BackupCancelPayload,
type BackupStartedPayload, type BackupStartedPayload,
type ControllerWireMessage, type ControllerWireMessage,
} from "@zerobyte/contracts/agent-protocol"; } from "@zerobyte/contracts/agent-protocol";

View file

@ -0,0 +1,74 @@
import { type ChildProcess, spawn } from "node:child_process";
import { existsSync } from "node:fs";
import path from "node:path";
import { logger } from "@zerobyte/core/node";
import { config } from "../../../core/config";
import { deriveLocalAgentToken } from "../helpers/tokens";
type LocalAgentState = {
localAgent: ChildProcess | null;
};
export async function spawnLocalAgentProcess(runtime: LocalAgentState) {
await stopLocalAgentProcess(runtime);
const sourceEntryPoint = path.join(process.cwd(), "apps", "agent", "src", "index.ts");
const productionEntryPoint = path.join(process.cwd(), ".output", "agent", "index.mjs");
if (config.__prod__ && !existsSync(productionEntryPoint)) {
throw new Error(`Local agent entrypoint not found at ${productionEntryPoint}`);
}
const agentEntryPoint = config.__prod__ ? productionEntryPoint : sourceEntryPoint;
const agentToken = await deriveLocalAgentToken();
const args = config.__prod__ ? ["run", agentEntryPoint] : ["run", "--watch", agentEntryPoint];
const agentProcess = spawn("bun", args, {
env: {
PATH: process.env.PATH,
ZEROBYTE_CONTROLLER_URL: "ws://localhost:3001",
ZEROBYTE_AGENT_TOKEN: agentToken,
},
stdio: ["ignore", "pipe", "pipe"],
});
runtime.localAgent = agentProcess;
agentProcess.stdout?.on("data", (data: Buffer) => {
const line = data.toString().trim();
if (line) logger.info(`[agent] ${line}`);
});
agentProcess.stderr?.on("data", (data: Buffer) => {
const line = data.toString().trim();
if (line) logger.error(`[agent] ${line}`);
});
agentProcess.on("exit", (code, signal) => {
if (runtime.localAgent === agentProcess) {
runtime.localAgent = null;
}
logger.info(`Agent process exited with code ${code} and signal ${signal}`);
});
}
export async function stopLocalAgentProcess(runtime: LocalAgentState) {
if (!runtime.localAgent) {
return;
}
const agentProcess = runtime.localAgent;
runtime.localAgent = null;
if (agentProcess.exitCode !== null || agentProcess.signalCode !== null) {
return;
}
const exited = new Promise<void>((resolve) => {
agentProcess.once("exit", () => {
resolve();
});
});
agentProcess.kill();
await exited;
}