From 0d02df760f40f1c323b1c4bd5421faa9ef9e511a Mon Sep 17 00:00:00 2001 From: Cristina Poncela Cubeiro <140309543+cristinaponcela@users.noreply.github.com> Date: Thu, 18 Jun 2026 13:25:50 +0200 Subject: [PATCH] feat: supervisor --- packages/orchestrator/package.json | 4 +- packages/orchestrator/src/handler.ts | 19 ++--- packages/orchestrator/src/index.ts | 1 + packages/orchestrator/src/supervisor.ts | 99 +++++++++++++++++++++++++ 4 files changed, 109 insertions(+), 14 deletions(-) create mode 100644 packages/orchestrator/src/supervisor.ts diff --git a/packages/orchestrator/package.json b/packages/orchestrator/package.json index a614b8da..555169bc 100644 --- a/packages/orchestrator/package.json +++ b/packages/orchestrator/package.json @@ -39,7 +39,9 @@ "engines": { "node": ">=22.19.0" }, - "dependencies": {}, + "dependencies": { + "@earendil-works/pi-coding-agent": "0.79.6" + }, "devDependencies": { "shx": "0.4.0" } diff --git a/packages/orchestrator/src/handler.ts b/packages/orchestrator/src/handler.ts index 0e894ba0..317a9247 100644 --- a/packages/orchestrator/src/handler.ts +++ b/packages/orchestrator/src/handler.ts @@ -1,4 +1,3 @@ -import { randomUUID } from "node:crypto"; import type { ErrorResponse, InstanceSummary, @@ -13,7 +12,7 @@ import type { StopRequest, StopResponse, } from "./ipc/protocol.ts"; -import { getInstance, loadInstances, removeInstance, upsertInstance } from "./storage.ts"; +import { supervisor } from "./supervisor.ts"; import type { InstanceRecord } from "./types.ts"; function toInstanceSummary(instance: InstanceRecord): InstanceSummary { @@ -43,15 +42,10 @@ export async function handleIpcRequest(request: OrchestratorRequest): Promise { switch (request.type) { case "spawn": { - const instance: InstanceRecord = { - id: randomUUID(), - status: "starting", + const instance = await supervisor.spawnInstance({ cwd: request.cwd, - createdAt: new Date().toISOString(), - lastSeenAt: new Date().toISOString(), label: request.label, - }; - upsertInstance(instance); + }); return { type: "spawn_result", ok: true, @@ -63,12 +57,12 @@ export async function handleIpcRequest(request: OrchestratorRequest): Promise { + const agentDir = getAgentDir(); + const sessionManager = SessionManager.inMemory(cwd); + const runtimeFactory: CreateAgentSessionRuntimeFactory = async ({ + cwd, + agentDir, + sessionManager, + sessionStartEvent, + }) => { + const services = await createAgentSessionServices({ cwd, agentDir }); + const created = await createAgentSessionFromServices({ + services, + sessionManager, + sessionStartEvent, + }); + return { + ...created, + services, + diagnostics: services.diagnostics, + }; + }; + + return createAgentSessionRuntime(runtimeFactory, { + cwd, + agentDir, + sessionManager, + }); +} + +export class OrchestratorSupervisor { + private readonly liveInstances = new Map(); + + listInstances(): InstanceRecord[] { + return loadInstances().map(cloneInstance); + } + + getInstance(instanceId: string): InstanceRecord | undefined { + const live = this.liveInstances.get(instanceId); + if (live) { + return cloneInstance(live.record); + } + const stored = getInstance(instanceId); + return stored ? cloneInstance(stored) : undefined; + } + + async spawnInstance(options: { cwd: string; label?: string }): Promise { + const runtime = await createRuntime(options.cwd); + const now = new Date().toISOString(); + const record: InstanceRecord = { + id: randomUUID(), + status: "online", + cwd: options.cwd, + createdAt: now, + lastSeenAt: now, + label: options.label, + sessionId: runtime.session.sessionId, + }; + + this.liveInstances.set(record.id, { runtime, record }); + upsertInstance(record); + return cloneInstance(record); + } + + async stopInstance(instanceId: string): Promise { + const live = this.liveInstances.get(instanceId); + if (!live) { + return undefined; + } + + await live.runtime.dispose(); + this.liveInstances.delete(instanceId); + removeInstance(instanceId); + return cloneInstance(live.record); + } +} + +export const supervisor = new OrchestratorSupervisor();