import { homedir } from "node:os"; import { dirname, join } from "node:path"; import { fileURLToPath } from "node:url"; import type { ExtensionAPI, ExtensionContext, ToolDefinition, } from "@earendil-works/pi-coding-agent"; import type { RunningA2aServer } from "./a2a/server.ts"; import { loadWorkflowOutputsFromEnv, type WorkflowOutputDeclarations, } from "./a2a/workflow-outputs.ts"; import { type A2aToolAccess, createA2aRuntimeListener, registerA2aTools } from "./a2a-extension.ts"; import channelEventsExtension, { locatorChannel, type SourceRegistration } from "./index.ts"; import relayExtension from "./relay-extension.ts"; import { errorMessage } from "./sources/util.ts"; import type { TaskPlane } from "./task-plane/plane.ts"; import { type RunningChannelsRuntime, startChannelsRuntime } from "./task-plane/runtime.ts"; import { TaskSessionHost, type TaskSessionHostOptions, type TaskTurnRunner, } from "./task-plane/task-sessions.ts"; import type { SourceTaskActivationSink } from "./task-plane/types.ts"; type TaskSessionOwner = TaskTurnRunner & { close(): Promise }; export interface ChannelsRuntimeExtensionDependencies { readonly startRuntime?: typeof startChannelsRuntime; readonly sources?: Readonly>; readonly log?: (record: Readonly>) => void; readonly createTaskSessionHost?: (options: TaskSessionHostOptions) => TaskSessionOwner; readonly loadWorkflowOutputs?: () => Promise; } const CHANNELS_PACKAGE_ROOT = dirname(dirname(fileURLToPath(import.meta.url))); /** * The package's sole Pi entrypoint. The relay remains a separate session plane; * every work-producing source is composed with the task plane. The entrypoint * is intentionally stateless so Jiti reloads with moduleCache:false * create an independent runtime rather than consulting a module-global owner. */ export default function channelsRuntimeExtension( pi: ExtensionAPI, dependencies: ChannelsRuntimeExtensionDependencies = {}, ): void { let runtime: RunningChannelsRuntime | undefined; let starting: Promise | undefined; let lifecycleRequest = 0; let stopped = false; let taskPlaneHealthy = false; let a2aServer: RunningA2aServer | undefined; let taskSessions: TaskSessionOwner | undefined; let startingTaskPlane: TaskPlane | undefined; let startingWakeQueue: RunningChannelsRuntime["wakeQueue"] | undefined; let startingSourceSink: SourceTaskActivationSink | undefined; let closing: Promise | undefined; let shutdownRequest = 0; let workflowOutputs: WorkflowOutputDeclarations | undefined; const taskTools: ToolDefinition[] = []; const taskToolPi = captureTools(pi, taskTools, async (locator) => { const sourceSink = runtime?.sourceSink ?? startingSourceSink; if (!sourceSink?.taskForLocator) throw new Error("task plane is not running"); await sourceSink.taskForLocator(locatorChannel(locator), locator); }); const selection = process.env.OUTFITTER_CHANNELS?.trim(); const taskPlaneEnabled = selection !== "off" && selection !== "none"; const startRuntime = dependencies.startRuntime ?? startChannelsRuntime; const log = dependencies.log ?? ((record: Readonly>): void => console.error(JSON.stringify(record))); const listener = createA2aRuntimeListener( { log }, (server) => { a2aServer = server; }, (taskId) => (runtime?.wakeQueue ?? startingWakeQueue)?.cancelTask(taskId), ); const taskAccess = (): A2aToolAccess | undefined => { if (a2aServer) return a2aServer; const store = runtime?.taskPlane.taskStore ?? startingTaskPlane?.taskStore; if (!store) return undefined; return { readTask: async (taskId) => (await store.lookup(taskId))?.task, controllerForTask: async (taskId) => { const stored = await store.lookup(taskId); if (!stored) return undefined; let current = stored.task; return { get task() { return current; }, async status(state, message) { current = await store.updateStatus(stored.principal, taskId, { state, ...(message ? { message } : {}), }); return current; }, async artifact(artifact) { current = await store.addArtifact(stored.principal, taskId, artifact); return current; }, }; }, }; }; if (taskPlaneEnabled || listener) { registerA2aTools( taskToolPi, taskAccess, async (taskId) => (runtime?.wakeQueue ?? startingWakeQueue)?.hasAuthority(taskId) ?? false, (taskId) => { const queue = runtime?.wakeQueue ?? startingWakeQueue; return queue !== undefined && queue.sourceForTask(taskId) === "a2a"; }, () => workflowOutputs, ); } const closeTaskPlane = (): Promise => { if (closing) return closing; const loaded = runtime; const sessions = taskSessions; runtime = undefined; taskSessions = undefined; startingTaskPlane = undefined; startingWakeQueue = undefined; startingSourceSink = undefined; const operation = (async () => { try { await loaded?.close(); } finally { await sessions?.close(); } })(); closing = operation; void operation.then( () => { if (closing === operation) closing = undefined; }, () => { if (closing === operation) closing = undefined; }, ); return operation; }; const clearStartingReferences = (): void => { startingTaskPlane = undefined; startingWakeQueue = undefined; startingSourceSink = undefined; }; const closeStartingGeneration = async ( loaded: RunningChannelsRuntime, sessionOwner: TaskSessionOwner, ): Promise => { await loaded.close(); if (taskSessions === sessionOwner) taskSessions = undefined; clearStartingReferences(); await sessionOwner.close().catch(() => {}); }; const createSessionOwner = ( context: ExtensionContext | undefined, taskPlaneRoot: string, ): TaskSessionOwner => { const getModel = (): ExtensionContext["model"] | undefined => { if (!context) return undefined; try { return context.model; } catch { return undefined; } }; const getThinkingLevel = (): ReturnType | undefined => { if (!context || typeof pi.getThinkingLevel !== "function") return undefined; try { return pi.getThinkingLevel(); } catch { return undefined; } }; const options = { cwd: context?.cwd ?? process.cwd(), sessionDir: join(taskPlaneRoot, "pi-sessions"), projectTrusted: context?.isProjectTrusted() ?? false, // Pi's event context resolves this getter against the active session on every read. model: getModel, thinkingLevel: getThinkingLevel, customTools: taskTools, excludedExtensionRoot: CHANNELS_PACKAGE_ROOT, log, }; return dependencies.createTaskSessionHost ? dependencies.createTaskSessionHost(options) : new TaskSessionHost(options); }; const startTaskRuntime = ( taskPlaneRoot: string, sessionOwner: TaskSessionOwner, ): Promise => startRuntime(pi, { storePath: join(taskPlaneRoot, "tasks.json"), originStorePath: join(taskPlaneRoot, "origins.json"), agentInterface: process.env.A2A_PUBLIC_URL?.trim() || `http://${process.env.A2A_HOST?.trim() || "127.0.0.1"}:${process.env.A2A_PORT?.trim() || "8788"}`, // Channel sources receive the guarded sink from this runtime after it opens. sources: [], taskTurnRunner: sessionOwner, taskPlaneReady: (taskPlane, wakeQueue, sourceSink) => { startingTaskPlane = taskPlane; startingWakeQueue = wakeQueue; startingSourceSink = sourceSink; }, ...(listener ? { listener } : {}), log, }); const launchTaskPlane = async (context: ExtensionContext | undefined): Promise => { await closing; if (stopped) return; const taskPlaneRoot = process.env.CHANNELS_TASK_STORE_PATH?.trim() || join( process.env.XDG_DATA_HOME?.trim() || join(homedir(), ".local", "share"), "outfitter", "channels", "task-plane", ); const sessionOwner = createSessionOwner(context, taskPlaneRoot); taskSessions = sessionOwner; try { const loaded = await startTaskRuntime(taskPlaneRoot, sessionOwner); if (stopped) { await closeStartingGeneration(loaded, sessionOwner); } else { runtime = loaded; clearStartingReferences(); taskPlaneHealthy = loaded.healthy; } } catch (error) { if (taskSessions === sessionOwner) taskSessions = undefined; clearStartingReferences(); await sessionOwner.close().catch(() => {}); throw error; } }; const joinStartingGeneration = async (): Promise => { const prior = starting; if (!prior) return false; if (!stopped) { await prior; return true; } lifecycleRequest += 1; stopped = false; await prior.catch(() => {}); if (runtime || stopped) return true; if (starting === prior) starting = undefined; return false; }; const beginTaskPlaneGeneration = async (context: ExtensionContext | undefined): Promise => { lifecycleRequest += 1; stopped = false; taskPlaneHealthy = false; starting = (async () => { try { workflowOutputs = await (dependencies.loadWorkflowOutputs ?? loadWorkflowOutputsFromEnv)(); } catch (error) { workflowOutputs = undefined; log({ event: "a2a_workflow_manifest_load_failed", error: errorMessage(error) }); } await launchTaskPlane(context); })(); try { await starting; } finally { starting = undefined; if (!runtime) startingTaskPlane = undefined; } }; // Register the task plane first. Pi dispatches lifecycle hooks in // registration order, so the durable local plane is ready before any // source or optional listener can accept work. if (taskPlaneEnabled || listener) { pi.on("session_start", async (_event, context) => { if (runtime) { taskPlaneHealthy = runtime.healthy; return; } while (starting) { if (await joinStartingGeneration()) return; } await beginTaskPlaneGeneration(context); }); // Mark shutdown before channel-source cleanup begins so a concurrently // requested restart can supersede this specific stop request. pi.on("session_shutdown", () => { shutdownRequest = ++lifecycleRequest; stopped = true; taskPlaneHealthy = false; }); } const channels = channelEventsExtension( taskToolPi, () => { const sourceSink = runtime?.sourceSink ?? startingSourceSink; if (!sourceSink) throw new Error("task plane is not running"); return sourceSink; }, dependencies.sources, async () => { taskPlaneHealthy = false; await closeTaskPlane(); }, () => !stopped && runtime !== undefined, ); // Channel shutdown was registered immediately above. Registering the plane // shutdown afterward guarantees providers stop before intake closes. pi.on("session_shutdown", async () => { const request = shutdownRequest; await starting?.catch(() => {}); if (request !== lifecycleRequest) return; await closeTaskPlane(); }); relayExtension(pi); if (taskPlaneEnabled || listener) { pi.on("session_start", () => { const selectedChannelsHealthy = !taskPlaneEnabled || channels?.startupSucceeded() === true; if (taskPlaneHealthy && selectedChannelsHealthy) log({ event: "channels_ready" }); else log({ event: "channels_unhealthy" }); }); } } function captureTools( pi: ExtensionAPI, captured: ToolDefinition[], authorizeChannelLocator: (locator: string) => Promise, ): ExtensionAPI { return new Proxy(pi, { get(target, property) { if (property === "registerTool") { return (tool: ToolDefinition): void => { if (isTaskTool(tool.name)) { captured.push(scopeTaskTool(tool, authorizeChannelLocator)); } target.registerTool(tool); }; } const value = Reflect.get(target, property, target); return typeof value === "function" ? value.bind(target) : value; }, }); } function scopeTaskTool( tool: ToolDefinition, authorizeChannelLocator: (locator: string) => Promise, ): ToolDefinition { if (tool.name !== "channel_read" && tool.name !== "channel_respond") return tool; return { ...tool, async execute(toolCallId, params, signal, onUpdate, context) { const locator = (params as { locator?: unknown }).locator; if (typeof locator !== "string") throw new Error("channel locator is required"); await authorizeChannelLocator(locator); return tool.execute(toolCallId, params, signal, onUpdate, context); }, }; } function isTaskTool(name: string): boolean { return name.startsWith("a2a_") || name === "channel_read" || name === "channel_respond"; }