From 5555434064a8591d4e91f9306f43f4df44c5f3fd Mon Sep 17 00:00:00 2001 From: kjgbot Date: Tue, 25 Aug 2026 21:47:32 +0200 Subject: [PATCH 1/3] refactor(core): extract broker transport port --- .../src/__tests__/broker-transport.test.ts | 47 + .../core/src/__tests__/idle-nudge.test.ts | 2 +- .../src/__tests__/workflow-runner.test.ts | 29 +- packages/core/src/agent-handle.ts | 6 +- packages/core/src/broker-transport.ts | 736 ++++++++++++++++ packages/core/src/index.ts | 1 + packages/core/src/runner.ts | 809 +++--------------- 7 files changed, 921 insertions(+), 709 deletions(-) create mode 100644 packages/core/src/__tests__/broker-transport.test.ts create mode 100644 packages/core/src/broker-transport.ts diff --git a/packages/core/src/__tests__/broker-transport.test.ts b/packages/core/src/__tests__/broker-transport.test.ts new file mode 100644 index 0000000..d0ec5b4 --- /dev/null +++ b/packages/core/src/__tests__/broker-transport.test.ts @@ -0,0 +1,47 @@ +import { describe, expect, it, vi } from 'vitest'; +import { + HarnessBrokerTransport, + resolveBrokerTransportMode, + type BrokerTransportPort, +} from '../broker-transport.js'; +import { WorkflowRunner } from '../runner.js'; + +describe('broker transport selection', () => { + it('defaults to legacy mode', () => { + expect(resolveBrokerTransportMode(undefined, {})).toBe('legacy'); + const runner = new WorkflowRunner(); + expect((runner as any).brokerTransport).toBeInstanceOf(HarnessBrokerTransport); + expect((runner as any).brokerTransport.mode).toBe('legacy'); + }); + + it('prefers the per-run selector over the environment selector', () => { + const runner = new WorkflowRunner({ + brokerTransportMode: 'shadow', + relay: { env: { RELAYFLOWS_INTEGRATION_TRANSPORT: 'adapter' } }, + }); + expect((runner as any).brokerTransport.mode).toBe('shadow'); + }); + + it('uses the rollout environment selector when no per-run selector is supplied', () => { + const runner = new WorkflowRunner({ + relay: { env: { RELAYFLOWS_INTEGRATION_TRANSPORT: 'adapter' } }, + }); + expect((runner as any).brokerTransport.mode).toBe('adapter'); + }); + + it('prefers an explicitly injected port over both selectors', () => { + const explicit = { mode: 'legacy', shutdown: vi.fn() } as unknown as BrokerTransportPort; + const runner = new WorkflowRunner({ + brokerTransport: explicit, + brokerTransportMode: 'shadow', + relay: { env: { RELAYFLOWS_INTEGRATION_TRANSPORT: 'adapter' } }, + }); + expect((runner as any).brokerTransport).toBe(explicit); + }); + + it('rejects an invalid rollout selector', () => { + expect(() => + resolveBrokerTransportMode(undefined, { RELAYFLOWS_INTEGRATION_TRANSPORT: 'invalid' }) + ).toThrow('Invalid broker transport mode'); + }); +}); diff --git a/packages/core/src/__tests__/idle-nudge.test.ts b/packages/core/src/__tests__/idle-nudge.test.ts index aaaf35b..2f6db28 100644 --- a/packages/core/src/__tests__/idle-nudge.test.ts +++ b/packages/core/src/__tests__/idle-nudge.test.ts @@ -225,7 +225,7 @@ describe('Idle Nudge Detection', () => { const agentDef = { name: 'worker', cli: 'claude' }; (runner as any).currentConfig = config; - (runner as any).relay = { sendMessage: mockSendMessage }; + (runner as any).brokerTransport.relay = { sendMessage: mockSendMessage }; const result = await (runner as any).waitForExitWithIdleNudging( wrappedMockAgent(), agentDef, diff --git a/packages/core/src/__tests__/workflow-runner.test.ts b/packages/core/src/__tests__/workflow-runner.test.ts index 23d597e..b7ce46b 100644 --- a/packages/core/src/__tests__/workflow-runner.test.ts +++ b/packages/core/src/__tests__/workflow-runner.test.ts @@ -650,14 +650,28 @@ agents: await Promise.all( runners.map((candidate, index) => - (candidate as any).startOrReuseSharedBroker(`run-${index}`, `wf-${index}`, true) + (candidate as any).brokerTransport.start( + { + runId: `run-${index}`, + brokerName: `relayflows-run-${index}`, + channel: `wf-${index}`, + relaycastDisabled: true, + }, + { + onEvent: vi.fn(), + onLog: vi.fn(), + getActiveAgentNames: () => [], + } + ) ) ); expect(HarnessDriverClient.spawn).toHaveBeenCalledTimes(1); expect(HarnessDriverClient.connect).toHaveBeenCalledTimes(3); - const starter = runners.find((candidate) => (candidate as any).sharedBrokerLease?.startedBroker); + const starter = runners.find( + (candidate) => (candidate as any).brokerTransport.sharedBrokerLease?.startedBroker + ); expect(starter).toBeDefined(); const attached = runners.filter((candidate) => candidate !== starter); @@ -810,15 +824,20 @@ agents: it('refuses broker recovery while workflow agents are still active', async () => { const candidate = new WorkflowRunner({ db, workspaceId: 'ws-test' }); - (candidate as any).currentBrokerContext = { + const transport = (candidate as any).brokerTransport; + transport.context = { runId: 'run-active', brokerName: 'relayflows-run-active', channel: 'wf-active', relaycastDisabled: true, }; - (candidate as any).activeAgentHandles.set('active-agent', { name: 'active-agent' }); + transport.hooks = { + onEvent: vi.fn(), + onLog: vi.fn(), + getActiveAgentNames: () => ['active-agent'], + }; - await expect((candidate as any).recoverBroker('spawn failed')).rejects.toThrow( + await expect(transport.recoverBroker('spawn failed')).rejects.toThrow( 'Broker recovery is unsafe while 1 agent is still active: active-agent' ); expect(mockHarnessDriverSpawn).not.toHaveBeenCalled(); diff --git a/packages/core/src/agent-handle.ts b/packages/core/src/agent-handle.ts index 38a0458..39a46b5 100644 --- a/packages/core/src/agent-handle.ts +++ b/packages/core/src/agent-handle.ts @@ -8,16 +8,16 @@ * string contract the workflow runner has always used (`'exited'` / `'timeout'` * / `'idle'`), so `runner.ts` keeps consuming handles unchanged. */ -import type { SpawnedAgentHandle } from '@agent-relay/harness-driver'; +import type { BrokerAgentHandle } from './broker-transport.js'; export class WorkflowAgentHandle { - constructor(private readonly inner: SpawnedAgentHandle) {} + constructor(private readonly inner: BrokerAgentHandle) {} get name(): string { return this.inner.name; } - get runtime(): SpawnedAgentHandle['runtime'] { + get runtime(): BrokerAgentHandle['runtime'] { return this.inner.runtime; } diff --git a/packages/core/src/broker-transport.ts b/packages/core/src/broker-transport.ts new file mode 100644 index 0000000..32c188a --- /dev/null +++ b/packages/core/src/broker-transport.ts @@ -0,0 +1,736 @@ +import { randomBytes } from 'node:crypto'; +import { + mkdirSync, + readFileSync, + readdirSync, + rmSync, + statSync, + unlinkSync, + writeFileSync, +} from 'node:fs'; +import type { Dirent } from 'node:fs'; +import path from 'node:path'; +import chalk from 'chalk'; +import { + HarnessDriverClient, + type BrokerEvent, + type ListAgent, + type RuntimeSpawnOptions, + type SendMessageInput, + type SpawnedAgentHandle, + type SpawnPtyInput, +} from '@agent-relay/harness-driver'; +import { RelayCast, RelayError, type AgentClient } from '@relaycast/sdk'; + +const BROKER_CONNECTION_FILENAME = 'connection.json'; +const SHARED_BROKER_LOCK_DIRNAME = '.relayflows-start.lock'; +const SHARED_BROKER_LEASE_DIRNAME = 'relayflows-runs'; +const SHARED_BROKER_OWNER_FILENAME = 'relayflows-owner.json'; +const SHARED_BROKER_LOCK_POLL_MS = 200; +const SHARED_BROKER_DEFAULT_STARTUP_TIMEOUT_MS = 45_000; +const BROKER_OPERATION_MAX_ATTEMPTS = 3; +const BROKER_OPERATION_RETRY_DELAY_MS = 1_000; + +export type BrokerTransportMode = 'legacy' | 'shadow' | 'adapter'; + +export interface BrokerRunContext { + runId: string; + brokerName: string; + channel: string; + relaycastDisabled: boolean; +} + +export interface BrokerTransportHooks { + onEvent: (event: BrokerEvent) => void; + onLog: (message: string) => void; + getActiveAgentNames: () => string[]; +} + +/** Broker-owned lifecycle handle returned to the workflow runtime. */ +export interface BrokerAgentHandle { + readonly name: string; + readonly runtime: SpawnedAgentHandle['runtime']; + readonly exitCode: number | undefined; + readonly exitSignal: string | undefined; + waitForExit(timeoutMs?: number): ReturnType; + waitForIdle(timeoutMs?: number): ReturnType; + release(reason?: string): Promise<{ name: string }>; +} + +/** + * Broker transport boundary used by WorkflowRunner. + * + * Implementations own broker connection/recovery, event subscription, worker + * lifecycle operations, Relaycast messaging, and shared lock/lease files. + */ +export interface BrokerTransportPort { + readonly mode: BrokerTransportMode; + readonly apiKey: string | undefined; + readonly apiKeyAutoCreated: boolean; + readonly connected: boolean; + ensureApiKey(channel: string): Promise; + start(context: BrokerRunContext, hooks: BrokerTransportHooks): Promise; + spawnPty(input: SpawnPtyInput, operation: string): Promise; + listAgents(operation: string): Promise; + release(name: string, reason: string | undefined, operation: string): Promise<{ name: string }>; + sendMessage(input: SendMessageInput, operation: string): Promise<{ event_id: string; targets: string[] }>; + /** Returns `pty` when written to stdin and `message` for the compatibility fallback. */ + sendInput(name: string, text: string, operation: string): Promise<'pty' | 'message'>; + createAndJoinChannel(channel: string, topic?: string): Promise; + startExternalAgentHeartbeat(name: string, persona?: string): Promise<(() => void) | undefined>; + inviteAgent(channel: string, name: string): Promise; + postToChannel(channel: string, text: string): Promise; + shutdown(): Promise; +} + +export interface HarnessBrokerTransportOptions { + mode?: BrokerTransportMode; + cwd: string; + relay?: RuntimeSpawnOptions; + resolveRelayEnv?: () => NodeJS.ProcessEnv | undefined; +} + +interface BrokerConnectionFile { + url: string; + api_key: string; + pid: number; + port?: number; +} + +interface SharedBrokerLease { + stateDir: string; + connectionPath: string; + ownerPath: string; + leasePath: string; + startedBroker: boolean; +} + +function parseBrokerConnectionFile(raw: string): BrokerConnectionFile | null { + try { + const conn = JSON.parse(raw); + if ( + typeof conn.url === 'string' && + typeof conn.api_key === 'string' && + typeof conn.pid === 'number' && + conn.pid > 0 + ) { + return conn as BrokerConnectionFile; + } + } catch { + // Invalid JSON is handled as no reusable broker. + } + return null; +} + +function readBrokerConnectionFile(connectionPath: string): BrokerConnectionFile | null { + try { + return parseBrokerConnectionFile(readFileSync(connectionPath, 'utf-8')); + } catch { + return null; + } +} + +function isPidRunning(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +function safeUnlinkSync(filePath: string): void { + try { + unlinkSync(filePath); + } catch { + // Best-effort cleanup. + } +} + +function sleepMs(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +/** Resolve the rollout selector without changing the legacy default. */ +export function resolveBrokerTransportMode( + explicitMode?: BrokerTransportMode, + env: NodeJS.ProcessEnv = process.env +): BrokerTransportMode { + const candidate = explicitMode ?? env.RELAYFLOWS_INTEGRATION_TRANSPORT ?? 'legacy'; + if (candidate === 'legacy' || candidate === 'shadow' || candidate === 'adapter') { + return candidate; + } + throw new Error( + `Invalid broker transport mode "${candidate}". Expected legacy, shadow, or adapter.` + ); +} + +class HarnessBrokerAgentHandle implements BrokerAgentHandle { + constructor(private readonly inner: SpawnedAgentHandle) {} + + get name(): string { + return this.inner.name; + } + + get runtime(): SpawnedAgentHandle['runtime'] { + return this.inner.runtime; + } + + get exitCode(): number | undefined { + return this.inner.exitCode; + } + + get exitSignal(): string | undefined { + return this.inner.exitSignal; + } + + waitForExit(timeoutMs?: number): ReturnType { + return this.inner.waitForExit(timeoutMs); + } + + waitForIdle(timeoutMs?: number): ReturnType { + return this.inner.waitForIdle(timeoutMs); + } + + release(reason?: string): Promise<{ name: string }> { + return this.inner.release(reason); + } +} + +/** Legacy-compatible transport backed by the published harness-driver and Relaycast SDK. */ +export class HarnessBrokerTransport implements BrokerTransportPort { + readonly mode: BrokerTransportMode; + private readonly cwd: string; + private readonly relayOptions: RuntimeSpawnOptions; + private readonly resolveRelayEnv: () => NodeJS.ProcessEnv | undefined; + private relay?: HarnessDriverClient; + private context?: BrokerRunContext; + private hooks?: BrokerTransportHooks; + private recoveryPromise?: Promise; + private relaycast?: RelayCast; + private relaycastAgent?: AgentClient; + private _apiKey?: string; + private _apiKeyAutoCreated = false; + /** @internal retained for focused lock/lease tests. */ + private sharedBrokerLease?: SharedBrokerLease; + private listenerDisposers: Array<() => void> = []; + + constructor(options: HarnessBrokerTransportOptions) { + this.mode = options.mode ?? 'legacy'; + this.cwd = options.cwd; + this.relayOptions = options.relay ?? {}; + this.resolveRelayEnv = options.resolveRelayEnv ?? (() => undefined); + } + + get apiKey(): string | undefined { + return this._apiKey; + } + + get apiKeyAutoCreated(): boolean { + return this._apiKeyAutoCreated; + } + + get connected(): boolean { + return this.relay !== undefined; + } + + async ensureApiKey(channel: string): Promise { + if (this._apiKey) return; + + const envKey = this.relayOptions.env?.RELAY_API_KEY ?? process.env.RELAY_API_KEY; + if (envKey) { + this._apiKey = envKey; + return; + } + + const workspaceName = `relay-${channel}-${randomBytes(4).toString('hex')}`; + const baseUrl = this.getRelaycastBaseUrl(); + const res = await fetch(`${baseUrl}/v1/workspaces`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ name: workspaceName }), + }); + if (!res.ok) { + throw new Error(`Failed to auto-create Relaycast workspace: ${res.status} ${await res.text()}`); + } + + const body = (await res.json()) as Record; + const data = (body.data ?? body) as Record; + const apiKey = data.api_key as string; + if (!apiKey) { + throw new Error('Relaycast workspace response missing api_key'); + } + + this._apiKey = apiKey; + this._apiKeyAutoCreated = true; + const dashboardPort = process.env.AGENT_RELAY_DASHBOARD_PORT || '3888'; + fetch(`http://127.0.0.1:${dashboardPort}/api/relay-config`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ apiKey }), + }) + .then((response) => { + if (!response.ok) { + console.warn(`[WorkflowRunner] dashboard key push failed: HTTP ${response.status}`); + } + }) + .catch(() => { + // Dashboard not running — silently ignore. + }); + } + + async start(context: BrokerRunContext, hooks: BrokerTransportHooks): Promise { + this.context = context; + this.hooks = hooks; + this.relaycast = undefined; + this.relaycastAgent = undefined; + await this.startOrReuseSharedBroker(context); + if (!this.relay) { + throw new Error('Broker client was not initialized'); + } + this.wireRelayClient(); + } + + async spawnPty(input: SpawnPtyInput, operation: string): Promise { + const handle = await this.withBrokerRecovery(operation, (relay) => relay.spawnPty(input)); + return new HarnessBrokerAgentHandle(handle); + } + + listAgents(operation: string): Promise { + return this.withBrokerRecovery(operation, (relay) => relay.listAgents()); + } + + release(name: string, reason: string | undefined, operation: string): Promise<{ name: string }> { + return this.withBrokerRecovery(operation, (relay) => relay.release(name, reason)); + } + + sendMessage( + input: SendMessageInput, + operation: string + ): Promise<{ event_id: string; targets: string[] }> { + return this.withBrokerRecovery(operation, (relay) => relay.sendMessage(input)); + } + + async sendInput(name: string, text: string, operation: string): Promise<'pty' | 'message'> { + const relay = this.relay; + if (!relay) { + throw new Error(`Broker unavailable while ${operation}`); + } + if (typeof relay.sendInput === 'function') { + await relay.sendInput(name, `${text}\r`); + return 'pty'; + } + await this.sendMessage({ from: 'workflow-runner', to: name, text }, operation); + return 'message'; + } + + async createAndJoinChannel(channel: string, topic?: string): Promise { + const agent = await this.ensureRelaycastRunnerAgent(); + try { + await agent.channels.create({ name: channel, ...(topic ? { topic } : {}) }); + } catch (error) { + if (!(error instanceof RelayError && error.code === 'name_conflict')) { + throw error; + } + } + await agent.channels.join(channel); + } + + async startExternalAgentHeartbeat( + name: string, + persona?: string + ): Promise<(() => void) | undefined> { + const agent = await this.registerRelaycastExternalAgent(name, persona); + if (!agent) return undefined; + const beat = () => { + agent.heartbeat().catch(() => {}); + }; + const timer = setInterval(beat, 30_000); + timer.unref(); + beat(); + return () => clearInterval(timer); + } + + async inviteAgent(channel: string, name: string): Promise { + const agent = await this.ensureRelaycastRunnerAgent(); + await agent.channels.invite(channel, name); + } + + async postToChannel(channel: string, text: string): Promise { + const agent = await this.ensureRelaycastRunnerAgent(); + await agent.send(channel, text); + } + + async shutdown(): Promise { + this.clearRelayListeners(); + const relay = this.relay; + const lease = this.sharedBrokerLease; + this.sharedBrokerLease = undefined; + this.relay = undefined; + this.recoveryPromise = undefined; + + if (!relay) { + if (lease) safeUnlinkSync(lease.leasePath); + this.resetRunState(); + return; + } + + if (!lease) { + await relay.shutdown(); + this.resetRunState(); + return; + } + + safeUnlinkSync(lease.leasePath); + const liveLeases = this.countLiveSharedBrokerLeases(lease.stateDir); + if (liveLeases === 0 && (lease.startedBroker || this.isWorkflowOwnedSharedBroker(lease))) { + await relay.shutdown(); + safeUnlinkSync(lease.connectionPath); + safeUnlinkSync(lease.ownerPath); + } else { + const disconnect = (relay as { disconnect?: () => void }).disconnect; + if (typeof disconnect === 'function') disconnect.call(relay); + } + this.resetRunState(); + } + + private resetRunState(): void { + this.context = undefined; + this.hooks = undefined; + this.relaycast = undefined; + this.relaycastAgent = undefined; + } + + private wireRelayClient(): void { + const relay = this.relay; + const hooks = this.hooks; + if (!relay || !hooks) return; + this.clearRelayListeners(); + this.listenerDisposers.push(relay.onEvent(hooks.onEvent)); + const unsubBrokerExit = relay.onBrokerExit?.((info) => { + if (this.relay?.brokerPid === info.pid) { + this.relay = undefined; + } + hooks.onLog( + `Broker exited (pid: ${info.pid ?? '?'}, code: ${info.code ?? '?'}, signal: ${info.signal ?? '?'})` + ); + }); + if (unsubBrokerExit) this.listenerDisposers.push(unsubBrokerExit); + relay.connectEvents(); + } + + private clearRelayListeners(): void { + for (const dispose of this.listenerDisposers) { + try { + dispose(); + } catch { + // Best-effort event cleanup. + } + } + this.listenerDisposers = []; + } + + private isRetryableProtocolError(error: unknown): boolean { + const candidate = error as { retryable?: unknown; status?: unknown; message?: unknown } | undefined; + if (candidate?.retryable === true) return true; + if (typeof candidate?.status === 'number' && candidate.status >= 500) return true; + const message = typeof candidate?.message === 'string' ? candidate.message : ''; + return /\b(fetch failed|econn|enotfound|eai_again|socket hang up|network|service unavailable|timed out)\b/i.test( + message + ); + } + + private async recoverBroker(reason: string): Promise { + const context = this.context; + const hooks = this.hooks; + if (!context || !hooks) { + throw new Error(`Broker unavailable and no recovery context exists (${reason})`); + } + if (this.recoveryPromise) { + await this.recoveryPromise; + return; + } + const activeAgents = hooks.getActiveAgentNames(); + if (activeAgents.length > 0) { + throw new Error( + `Broker recovery is unsafe while ${activeAgents.length} agent${activeAgents.length === 1 ? ' is' : 's are'} still active: ${activeAgents.slice(0, 3).join(', ')}` + ); + } + + this.recoveryPromise = (async () => { + hooks.onLog(`Broker unavailable (${reason}); restarting...`); + await this.shutdown().catch(() => undefined); + await this.start(context, hooks); + hooks.onLog('Broker restarted'); + })(); + try { + await this.recoveryPromise; + } finally { + this.recoveryPromise = undefined; + } + } + + private async withBrokerRecovery( + operation: string, + work: (relay: HarnessDriverClient) => Promise + ): Promise { + let lastError: unknown; + for (let attempt = 1; attempt <= BROKER_OPERATION_MAX_ATTEMPTS; attempt++) { + const relay = this.relay; + if (!relay) { + lastError = new Error(`Broker unavailable while ${operation}`); + } else { + try { + return await work(relay); + } catch (error) { + lastError = error; + if (!this.isRetryableProtocolError(error)) throw error; + } + } + if (attempt >= BROKER_OPERATION_MAX_ATTEMPTS) break; + await this.recoverBroker(`${operation} failed`); + await sleepMs(BROKER_OPERATION_RETRY_DELAY_MS * attempt); + } + const message = lastError instanceof Error ? lastError.message : String(lastError); + throw new Error(`Broker operation failed during ${operation}: ${message}`); + } + + private getBrokerCwd(): string { + return this.relayOptions.cwd ?? this.cwd; + } + + private getBrokerStateDir(brokerCwd: string): string { + const configured = + this.relayOptions.binaryArgs?.stateDir ?? + this.relayOptions.env?.AGENT_RELAY_STATE_DIR ?? + process.env.AGENT_RELAY_STATE_DIR; + return path.resolve(configured ?? path.join(brokerCwd, '.agentworkforce', 'relay')); + } + + private async tryConnectSharedBroker( + connectionPath: string, + brokerCwd: string + ): Promise { + const conn = readBrokerConnectionFile(connectionPath); + if (!conn) return null; + if (!isPidRunning(conn.pid)) { + safeUnlinkSync(connectionPath); + return null; + } + try { + const client = HarnessDriverClient.connect({ cwd: brokerCwd, connectionPath }); + await client.getStatus(); + return client; + } catch { + return null; + } + } + + private async acquireSharedBrokerStartLock( + stateDir: string, + startupTimeoutMs: number + ): Promise<() => void> { + mkdirSync(stateDir, { recursive: true }); + const lockDir = path.join(stateDir, SHARED_BROKER_LOCK_DIRNAME); + const deadline = Date.now() + Math.max(startupTimeoutMs + 5_000, 10_000); + const staleAfterMs = Math.max(startupTimeoutMs * 2, 30_000); + for (;;) { + try { + mkdirSync(lockDir); + writeFileSync( + path.join(lockDir, 'owner.json'), + JSON.stringify({ pid: process.pid, createdAt: new Date().toISOString() }), + 'utf-8' + ); + return () => rmSync(lockDir, { recursive: true, force: true }); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error; + } + try { + const stats = statSync(lockDir); + if (Date.now() - stats.mtimeMs > staleAfterMs) { + rmSync(lockDir, { recursive: true, force: true }); + continue; + } + } catch { + continue; + } + if (Date.now() > deadline) { + throw new Error(`Timed out waiting for shared broker startup lock at ${lockDir}`); + } + await sleepMs(SHARED_BROKER_LOCK_POLL_MS); + } + } + + private createSharedBrokerLease( + stateDir: string, + connectionPath: string, + runId: string, + startedBroker: boolean + ): SharedBrokerLease { + const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME); + const ownerPath = path.join(stateDir, SHARED_BROKER_OWNER_FILENAME); + mkdirSync(leaseDir, { recursive: true }); + const leasePath = path.join( + leaseDir, + `${process.pid}-${runId}-${randomBytes(4).toString('hex')}.json` + ); + writeFileSync( + leasePath, + JSON.stringify({ pid: process.pid, runId, startedBroker, createdAt: new Date().toISOString() }), + 'utf-8' + ); + return { stateDir, connectionPath, ownerPath, leasePath, startedBroker }; + } + + private writeSharedBrokerOwner(lease: SharedBrokerLease): void { + const conn = readBrokerConnectionFile(lease.connectionPath); + writeFileSync( + lease.ownerPath, + JSON.stringify({ pid: conn?.pid, createdByPid: process.pid, createdAt: new Date().toISOString() }), + 'utf-8' + ); + } + + private isWorkflowOwnedSharedBroker(lease: SharedBrokerLease): boolean { + const conn = readBrokerConnectionFile(lease.connectionPath); + if (!conn) return false; + try { + const owner = JSON.parse(readFileSync(lease.ownerPath, 'utf-8')) as { pid?: unknown }; + return owner.pid === conn.pid; + } catch { + return false; + } + } + + private countLiveSharedBrokerLeases(stateDir: string): number { + const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME); + let entries: Dirent[]; + try { + entries = readdirSync(leaseDir, { withFileTypes: true }); + } catch { + return 0; + } + let live = 0; + for (const entry of entries) { + if (!entry.isFile()) continue; + const leasePath = path.join(leaseDir, entry.name); + try { + const lease = JSON.parse(readFileSync(leasePath, 'utf-8')) as { pid?: unknown }; + if (typeof lease.pid === 'number' && lease.pid > 0 && isPidRunning(lease.pid)) { + live += 1; + } else { + safeUnlinkSync(leasePath); + } + } catch { + safeUnlinkSync(leasePath); + } + } + return live; + } + + private async startOrReuseSharedBroker(context: BrokerRunContext): Promise { + const brokerCwd = this.getBrokerCwd(); + const stateDir = this.getBrokerStateDir(brokerCwd); + const connectionPath = path.join(stateDir, BROKER_CONNECTION_FILENAME); + const startupTimeoutMs = + this.relayOptions.startupTimeoutMs ?? SHARED_BROKER_DEFAULT_STARTUP_TIMEOUT_MS; + const lease = this.createSharedBrokerLease(stateDir, connectionPath, context.runId, false); + this.sharedBrokerLease = lease; + + const existing = await this.tryConnectSharedBroker(connectionPath, brokerCwd); + if (existing) { + this.hooks?.onLog('Reusing shared broker...'); + this.relay = existing; + return; + } + + const releaseLock = await this.acquireSharedBrokerStartLock(stateDir, startupTimeoutMs); + try { + const lockedExisting = await this.tryConnectSharedBroker(connectionPath, brokerCwd); + if (lockedExisting) { + this.hooks?.onLog('Reusing shared broker...'); + this.relay = lockedExisting; + return; + } + + this.hooks?.onLog('Starting broker...'); + const relayEnv = { + ...(this.resolveRelayEnv() ?? {}), + AGENT_RELAY_STATE_DIR: stateDir, + }; + this.relay = await HarnessDriverClient.spawn({ + ...this.relayOptions, + cwd: brokerCwd, + brokerName: context.brokerName, + channels: context.relaycastDisabled ? [] : [context.channel], + binaryArgs: { ...(this.relayOptions.binaryArgs ?? {}), persist: true, stateDir }, + env: relayEnv, + requestTimeoutMs: this.relayOptions.requestTimeoutMs ?? 120_000, + onStderr: (line: string) => { + const trimmed = line.trim(); + if (!trimmed || (trimmed.startsWith('{') && trimmed.endsWith('}'))) return; + console.log(`${chalk.dim.yellow('[broker]')} ${line}`); + }, + }); + lease.startedBroker = true; + this.writeSharedBrokerOwner(lease); + } finally { + releaseLock(); + } + } + + private getRelaycastBaseUrl(): string { + return ( + this.relayOptions.env?.RELAYCAST_BASE_URL ?? + process.env.RELAYCAST_BASE_URL ?? + 'https://api.relaycast.dev' + ); + } + + private getRelaycastClient(): RelayCast { + if (!this._apiKey) throw new Error('No Relaycast API key available'); + if (!this.relaycast) { + this.relaycast = new RelayCast({ apiKey: this._apiKey, baseUrl: this.getRelaycastBaseUrl() }); + } + return this.relaycast; + } + + private async ensureRelaycastRunnerAgent(): Promise { + if (this.relaycastAgent) return this.relaycastAgent; + const rc = this.getRelaycastClient(); + let registration; + try { + registration = await rc.agents.register({ name: 'WorkflowRunner', type: 'agent' }); + } catch (error) { + if (error instanceof RelayError && error.code === 'name_conflict') { + registration = await rc.agents.register({ + name: `WorkflowRunner-${randomBytes(4).toString('hex')}`, + type: 'agent', + }); + } else { + throw error; + } + } + this.relaycastAgent = rc.as(registration.token); + return this.relaycastAgent; + } + + private async registerRelaycastExternalAgent( + name: string, + persona?: string + ): Promise { + const rc = this.getRelaycastClient(); + try { + const registration = await rc.agents.register({ + name, + type: 'agent', + ...(persona ? { persona } : {}), + }); + return rc.as(registration.token); + } catch (error) { + if (error instanceof RelayError && error.code === 'name_conflict') return null; + throw error; + } + } +} diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 2269a0b..956abd4 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -1,5 +1,6 @@ export * from './types.js'; export * from './runner.js'; +export * from './broker-transport.js'; export * from './custom-steps.js'; export * from './cli-session-collector.js'; export * from './channel-messenger.js'; diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 7344a0d..8c503f3 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -14,9 +14,7 @@ import { readFileSync, readdirSync, renameSync, - rmSync, statSync, - unlinkSync, watch, writeFileSync, } from 'node:fs'; @@ -144,12 +142,16 @@ import { // ── Broker client / messaging imports ─────────────────────────────────────── -// Broker / PTY / lifecycle is driven by the harness-driver client; messaging -// uses @relaycast/sdk (below). -import { HarnessDriverClient } from '@agent-relay/harness-driver'; -import type { RuntimeSpawnOptions, SpawnPtyInput } from '@agent-relay/harness-driver'; +// Broker / PTY / Relaycast lifecycle is isolated behind BrokerTransportPort. +import type { RuntimeSpawnOptions } from '@agent-relay/harness-driver'; import { WorkflowAgentHandle } from './agent-handle.js'; -import { RelayCast, RelayError, type AgentClient } from '@relaycast/sdk'; +import { + HarnessBrokerTransport, + resolveBrokerTransportMode, + type BrokerRunContext, + type BrokerTransportMode, + type BrokerTransportPort, +} from './broker-transport.js'; import { SlackClient } from '@relayflows/slack-primitive'; import { RelayfileSetup, RelayFileClient, type ChangeEvent, type FilesystemEvent, type Subscription } from '@relayfile/sdk'; @@ -230,72 +232,6 @@ function filteredEnv(extra?: Record): Record 0 - ) { - return conn as BrokerConnectionFile; - } - } catch { - // Invalid JSON is handled as no reusable broker. - } - return null; -} - -function readBrokerConnectionFile(connectionPath: string): BrokerConnectionFile | null { - try { - return parseBrokerConnectionFile(readFileSync(connectionPath, 'utf-8')); - } catch { - return null; - } -} - -function isPidRunning(pid: number): boolean { - try { - process.kill(pid, 0); - return true; - } catch { - return false; - } -} - -function safeUnlinkSync(filePath: string): void { - try { - unlinkSync(filePath); - } catch { - // Best-effort cleanup. - } -} - function sleepMs(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } @@ -409,6 +345,10 @@ export interface WorkflowRunnerOptions { db?: WorkflowDb; workspaceId?: string; relay?: RuntimeSpawnOptions; + /** Explicit broker transport. Takes precedence over the rollout selector. */ + brokerTransport?: BrokerTransportPort; + /** Per-run rollout selector; falls back to RELAYFLOWS_INTEGRATION_TRANSPORT. */ + brokerTransportMode?: BrokerTransportMode; cwd?: string; summaryDir?: string; executor?: RunnerStepExecutor; @@ -596,8 +536,6 @@ const DEFAULT_ANSWER_FILE_POLL_MS = 5_000; const DEFAULT_WORKFLOW_MAX_RETRIES = 2; const DEFAULT_WORKFLOW_REPAIR_RETRIES = 2; const DEFAULT_WORKFLOW_RETRY_DELAY_MS = 1000; -const BROKER_OPERATION_MAX_ATTEMPTS = 3; -const BROKER_OPERATION_RETRY_DELAY_MS = 1_000; const AGENT_TRANSIENT_NETWORK_MAX_ATTEMPTS = 3; const AGENT_TRANSIENT_NETWORK_RETRY_DELAY_MS = 1_000; @@ -834,13 +772,6 @@ interface ChannelEvidenceOptions { origin?: CompletionEvidenceChannelOrigin; } -interface BrokerRunContext { - runId: string; - brokerName: string; - channel: string; - relaycastDisabled: boolean; -} - // ── CLI resolution ─────────────────────────────────────────────────────────── /** @@ -884,15 +815,7 @@ export class WorkflowRunner { private readonly envSecrets?: Record; private readonly templateResolver: TemplateResolver; private readonly channelMessenger: ChannelMessenger; - - /** @internal exposed for CLI signal-handler shutdown only */ - relay?: HarnessDriverClient; - private currentBrokerContext?: BrokerRunContext; - private brokerRecoveryPromise?: Promise; - private relaycast?: RelayCast; - private relaycastAgent?: AgentClient; - private relayApiKey?: string; - private relayApiKeyAutoCreated = false; + private readonly brokerTransport: BrokerTransportPort; private channel?: string; private trajectory?: WorkflowTrajectory; private abortController?: AbortController; @@ -965,11 +888,6 @@ export class WorkflowRunner { private workersFileLock: Promise = Promise.resolve(); /** Timestamp when the current workflow run started, for elapsed-time logging. */ private runStartTime?: number; - /** Unsubscribe handle for broker stderr listener wired during a run. */ - private unsubBrokerStderr?: () => void; - private unsubRelayListeners: Array<() => void> = []; - /** Local lease metadata for the shared workflow broker, when broker init was needed. */ - private sharedBrokerLease?: SharedBrokerLease; /** Tracks last idle log time per agent to debounce idle warnings (30s multiples). */ private readonly lastIdleLog = new Map(); /** Tracks last logged activity type per agent to avoid duplicate status lines. */ @@ -1005,6 +923,17 @@ export class WorkflowRunner { this.executor = options.executor; this.processBackend = options.processBackend; this.envSecrets = options.envSecrets; + this.brokerTransport = + options.brokerTransport ?? + new HarnessBrokerTransport({ + cwd: this.cwd, + relay: this.relayOptions, + mode: resolveBrokerTransportMode(options.brokerTransportMode, { + ...process.env, + ...(this.relayOptions.env ?? {}), + }), + resolveRelayEnv: () => this.getRelayEnv() ?? filteredEnv(), + }); if (!this.executor && !this.processBackend) { // Only reached when the caller injected neither an executor nor a // backend. The config's provider defaults to `none`, which yields @@ -2255,54 +2184,6 @@ export class WorkflowRunner { }; } - private clearRelayListeners(): void { - for (const off of this.unsubRelayListeners) { - try { - off(); - } catch { - /* ignore */ - } - } - this.unsubRelayListeners = []; - } - - private wireRelayClient(runId: string): void { - if (!this.relay) return; - - this.clearRelayListeners(); - this.unsubRelayListeners.push(this.relay.onEvent(this.createBrokerEventHandler(runId))); - const unsubBrokerExit = this.relay.onBrokerExit?.((info) => { - if (this.relay?.brokerPid === info.pid) { - this.relay = undefined; - } - this.log( - `Broker exited (pid: ${info.pid ?? '?'}, code: ${info.code ?? '?'}, signal: ${info.signal ?? '?'})` - ); - }); - if (unsubBrokerExit) { - this.unsubRelayListeners.push(unsubBrokerExit); - } - this.relay.connectEvents(); - } - - private async startBroker(context: BrokerRunContext): Promise { - await this.startOrReuseSharedBroker(context.runId, context.channel, context.relaycastDisabled); - if (!this.relay) { - throw new Error('Broker client was not initialized'); - } - this.wireRelayClient(context.runId); - } - - private isRetryableProtocolError(error: unknown): boolean { - const candidate = error as { retryable?: unknown; status?: unknown; message?: unknown } | undefined; - if (candidate?.retryable === true) return true; - if (typeof candidate?.status === 'number' && candidate.status >= 500) return true; - const message = typeof candidate?.message === 'string' ? candidate.message : ''; - return /\b(fetch failed|econn|enotfound|eai_again|socket hang up|network|service unavailable|timed out)\b/i.test( - message - ); - } - private isTransientAgentNetworkError(error: unknown): boolean { const candidate = error as { retryable?: unknown; status?: unknown; message?: unknown } | undefined; if (candidate?.retryable === true) return true; @@ -2311,129 +2192,6 @@ export class WorkflowRunner { return /\b(fetch failed|econn|enotfound|eai_again|socket hang up|network error|connection reset|connection refused|service unavailable)\b/i.test(message); } - private async recoverBroker(reason: string): Promise { - if (!this.currentBrokerContext) { - throw new Error(`Broker unavailable and no recovery context exists (${reason})`); - } - if (this.brokerRecoveryPromise) { - await this.brokerRecoveryPromise; - return; - } - if (this.activeAgentHandles.size > 0) { - const activeAgents = [...this.activeAgentHandles.keys()]; - throw new Error( - `Broker recovery is unsafe while ${activeAgents.length} agent${activeAgents.length === 1 ? ' is' : 's are'} still active: ${activeAgents.slice(0, 3).join(', ')}` - ); - } - - this.brokerRecoveryPromise = (async () => { - this.log(`Broker unavailable (${reason}); restarting...`); - this.clearRelayListeners(); - await this.shutdownRelay().catch(() => undefined); - await this.startBroker(this.currentBrokerContext!); - this.log('Broker restarted'); - })(); - - try { - await this.brokerRecoveryPromise; - } finally { - this.brokerRecoveryPromise = undefined; - } - } - - private async withBrokerRecovery(operation: string, work: (relay: HarnessDriverClient) => Promise): Promise { - let lastError: unknown; - for (let attempt = 1; attempt <= BROKER_OPERATION_MAX_ATTEMPTS; attempt++) { - const relay = this.relay; - if (!relay) { - lastError = new Error(`Broker unavailable while ${operation}`); - } else { - try { - return await work(relay); - } catch (error) { - lastError = error; - if (!this.isRetryableProtocolError(error)) { - throw error; - } - } - } - - if (attempt >= BROKER_OPERATION_MAX_ATTEMPTS) { - break; - } - await this.recoverBroker(`${operation} failed`); - await this.delay(BROKER_OPERATION_RETRY_DELAY_MS * attempt); - } - - const message = lastError instanceof Error ? lastError.message : String(lastError); - throw new Error(`Broker operation failed during ${operation}: ${message}`); - } - - // ── Relaycast auto-provisioning ──────────────────────────────────────── - - /** - * Ensure a Relaycast workspace API key is available for the broker. - * Resolution order: - * 1. RELAY_API_KEY environment variable (explicit override) - * 2. Auto-create a fresh workspace via the Relaycast API - * - * Each workflow run gets its own isolated workspace — no caching, no sharing. - */ - private async ensureRelaycastApiKey(channel: string): Promise { - if (this.relayApiKey) return; - - // Explicit override from relayOptions or environment takes priority. - const envKey = this.relayOptions.env?.RELAY_API_KEY ?? process.env.RELAY_API_KEY; - if (envKey) { - this.relayApiKey = envKey; - return; - } - - // Always create a fresh workspace — each run gets full isolation. - const workspaceName = `relay-${channel}-${randomBytes(4).toString('hex')}`; - const baseUrl = - this.relayOptions.env?.RELAYCAST_BASE_URL ?? - process.env.RELAYCAST_BASE_URL ?? - 'https://api.relaycast.dev'; - const res = await fetch(`${baseUrl}/v1/workspaces`, { - method: 'POST', - headers: { 'content-type': 'application/json' }, - body: JSON.stringify({ name: workspaceName }), - }); - - if (!res.ok) { - throw new Error(`Failed to auto-create Relaycast workspace: ${res.status} ${await res.text()}`); - } - - const body = (await res.json()) as Record; - const data = (body.data ?? body) as Record; - const apiKey = data.api_key as string; - - if (!apiKey) { - throw new Error('Relaycast workspace response missing api_key'); - } - - this.relayApiKey = apiKey; - this.relayApiKeyAutoCreated = true; - - // Best-effort: push the key to a co-running dashboard (agent-relay up) so it - // can make Relaycast API calls without any file or manual env var setup. - const dashboardPort = process.env.AGENT_RELAY_DASHBOARD_PORT || '3888'; - fetch(`http://127.0.0.1:${dashboardPort}/api/relay-config`, { - method: 'POST', - headers: { 'content-type': 'application/json' }, - body: JSON.stringify({ apiKey }), - }) - .then((res) => { - if (!res.ok) { - console.warn(`[WorkflowRunner] dashboard key push failed: HTTP ${res.status}`); - } - }) - .catch(() => { - // Dashboard not running — silently ignore. - }); - } - private async loadCredentialProxyModule(): Promise { try { const dynamicImport = new Function('specifier', 'return import(specifier)') as ( @@ -2613,7 +2371,7 @@ export class WorkflowRunner { return { ...process.env, ...(this.relayOptions.env ?? {}), - ...(this.relayApiKey ? { RELAY_API_KEY: this.relayApiKey } : {}), + ...(this.brokerTransport.apiKey ? { RELAY_API_KEY: this.brokerTransport.apiKey } : {}), }; } @@ -2623,7 +2381,7 @@ export class WorkflowRunner { const inheritedProxyToken = resolveProxyTokenFromEnv(env); if ( - !this.relayApiKey && + !this.brokerTransport.apiKey && !this.relayOptions.env && !proxyMode && !(inheritedProxyUrl && inheritedProxyToken) @@ -2648,267 +2406,8 @@ export class WorkflowRunner { return env; } - private getBrokerCwd(): string { - return this.relayOptions.cwd ?? this.cwd; - } - - private getBrokerStateDir(brokerCwd: string): string { - const configured = - this.relayOptions.binaryArgs?.stateDir ?? - this.relayOptions.env?.AGENT_RELAY_STATE_DIR ?? - process.env.AGENT_RELAY_STATE_DIR; - return path.resolve(configured ?? path.join(brokerCwd, '.agentworkforce', 'relay')); - } - - private async tryConnectSharedBroker( - connectionPath: string, - brokerCwd: string - ): Promise { - const conn = readBrokerConnectionFile(connectionPath); - if (!conn) { - return null; - } - - if (!isPidRunning(conn.pid)) { - safeUnlinkSync(connectionPath); - return null; - } - - try { - const client = HarnessDriverClient.connect({ cwd: brokerCwd, connectionPath }); - await client.getStatus(); - return client; - } catch { - return null; - } - } - - private async acquireSharedBrokerStartLock( - stateDir: string, - startupTimeoutMs: number - ): Promise<() => void> { - mkdirSync(stateDir, { recursive: true }); - const lockDir = path.join(stateDir, SHARED_BROKER_LOCK_DIRNAME); - const deadline = Date.now() + Math.max(startupTimeoutMs + 5_000, 10_000); - const staleAfterMs = Math.max(startupTimeoutMs * 2, 30_000); - - for (;;) { - try { - mkdirSync(lockDir); - writeFileSync( - path.join(lockDir, 'owner.json'), - JSON.stringify({ pid: process.pid, createdAt: new Date().toISOString() }), - 'utf-8' - ); - return () => { - rmSync(lockDir, { recursive: true, force: true }); - }; - } catch (err) { - if ((err as NodeJS.ErrnoException).code !== 'EEXIST') { - throw err; - } - } - - try { - const stat = statSync(lockDir); - if (Date.now() - stat.mtimeMs > staleAfterMs) { - rmSync(lockDir, { recursive: true, force: true }); - continue; - } - } catch { - continue; - } - - if (Date.now() > deadline) { - throw new Error(`Timed out waiting for shared broker startup lock at ${lockDir}`); - } - await sleepMs(SHARED_BROKER_LOCK_POLL_MS); - } - } - - private createSharedBrokerLease( - stateDir: string, - connectionPath: string, - runId: string, - startedBroker: boolean - ): SharedBrokerLease { - const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME); - const ownerPath = path.join(stateDir, SHARED_BROKER_OWNER_FILENAME); - mkdirSync(leaseDir, { recursive: true }); - const leasePath = path.join( - leaseDir, - `${process.pid}-${runId}-${randomBytes(4).toString('hex')}.json` - ); - writeFileSync( - leasePath, - JSON.stringify({ - pid: process.pid, - runId, - startedBroker, - createdAt: new Date().toISOString(), - }), - 'utf-8' - ); - return { stateDir, connectionPath, ownerPath, leasePath, startedBroker }; - } - - private writeSharedBrokerOwner(lease: SharedBrokerLease): void { - const conn = readBrokerConnectionFile(lease.connectionPath); - writeFileSync( - lease.ownerPath, - JSON.stringify({ - pid: conn?.pid, - createdByPid: process.pid, - createdAt: new Date().toISOString(), - }), - 'utf-8' - ); - } - - private isWorkflowOwnedSharedBroker(lease: SharedBrokerLease): boolean { - const conn = readBrokerConnectionFile(lease.connectionPath); - if (!conn) { - return false; - } - try { - const owner = JSON.parse(readFileSync(lease.ownerPath, 'utf-8')) as { pid?: unknown }; - return owner.pid === conn.pid; - } catch { - return false; - } - } - - private disconnectRelayClient(relay: HarnessDriverClient): void { - const disconnect = (relay as { disconnect?: () => void }).disconnect; - if (typeof disconnect === 'function') { - disconnect.call(relay); - } - } - - private countLiveSharedBrokerLeases(stateDir: string): number { - const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME); - let entries: Dirent[]; - try { - entries = readdirSync(leaseDir, { withFileTypes: true }); - } catch { - return 0; - } - - let live = 0; - for (const entry of entries) { - if (!entry.isFile()) continue; - const leasePath = path.join(leaseDir, entry.name); - try { - const lease = JSON.parse(readFileSync(leasePath, 'utf-8')) as { pid?: unknown }; - if (typeof lease.pid === 'number' && lease.pid > 0 && isPidRunning(lease.pid)) { - live += 1; - } else { - safeUnlinkSync(leasePath); - } - } catch { - safeUnlinkSync(leasePath); - } - } - return live; - } - - private async startOrReuseSharedBroker( - runId: string, - channel: string, - relaycastDisabled: boolean - ): Promise { - const brokerCwd = this.getBrokerCwd(); - const stateDir = this.getBrokerStateDir(brokerCwd); - const connectionPath = path.join(stateDir, BROKER_CONNECTION_FILENAME); - const startupTimeoutMs = - this.relayOptions.startupTimeoutMs ?? SHARED_BROKER_DEFAULT_STARTUP_TIMEOUT_MS; - const lease = this.createSharedBrokerLease(stateDir, connectionPath, runId, false); - this.sharedBrokerLease = lease; - - const existing = await this.tryConnectSharedBroker(connectionPath, brokerCwd); - if (existing) { - this.log('Reusing shared broker...'); - this.relay = existing; - return; - } - - const releaseLock = await this.acquireSharedBrokerStartLock(stateDir, startupTimeoutMs); - try { - const lockedExisting = await this.tryConnectSharedBroker(connectionPath, brokerCwd); - if (lockedExisting) { - this.log('Reusing shared broker...'); - this.relay = lockedExisting; - return; - } - - this.log('Starting broker...'); - // Include a short run ID suffix in the broker name so a newly-created - // broker keeps the same Relaycast identity behavior as previous runs. - const brokerBaseName = path.basename(this.cwd) || 'workflow'; - const brokerName = `${brokerBaseName}-${runId.slice(0, 8)}`; - const relayEnv = { - ...(this.getRelayEnv() ?? filteredEnv()), - AGENT_RELAY_STATE_DIR: stateDir, - }; - this.relay = await HarnessDriverClient.spawn({ - ...this.relayOptions, - cwd: brokerCwd, - brokerName, - channels: relaycastDisabled ? [] : [channel], - binaryArgs: { - ...(this.relayOptions.binaryArgs ?? {}), - persist: true, - stateDir, - }, - env: relayEnv, - // Workflows spawn agents across multiple waves; each spawn requires a PTY + - // Relaycast registration. 60s is too tight when the broker is saturated with - // long-running PTY processes from earlier steps. 120s gives room to breathe. - requestTimeoutMs: this.relayOptions.requestTimeoutMs ?? 120_000, - // Wire broker stderr to console for observability — skip empty and - // JSON event lines (already surfaced via the broker:event emitter). - onStderr: (line: string) => { - const trimmed = line.trim(); - if (!trimmed) return; - if (trimmed.startsWith('{') && trimmed.endsWith('}')) return; - console.log(`${chalk.dim.yellow('[broker]')} ${line}`); - }, - }); - lease.startedBroker = true; - this.writeSharedBrokerOwner(lease); - } finally { - releaseLock(); - } - } - async shutdownRelay(): Promise { - const relay = this.relay; - const lease = this.sharedBrokerLease; - this.sharedBrokerLease = undefined; - - if (!relay) { - if (lease) { - safeUnlinkSync(lease.leasePath); - } - return; - } - - this.relay = undefined; - - if (!lease) { - await relay.shutdown(); - return; - } - - safeUnlinkSync(lease.leasePath); - const liveLeases = this.countLiveSharedBrokerLeases(lease.stateDir); - if (liveLeases === 0 && (lease.startedBroker || this.isWorkflowOwnedSharedBroker(lease))) { - await relay.shutdown(); - safeUnlinkSync(lease.connectionPath); - safeUnlinkSync(lease.ownerPath); - } else { - this.disconnectRelayClient(relay); - } + await this.brokerTransport.shutdown(); } private async provisionAgents(config: RelayYamlConfig): Promise { @@ -2964,88 +2463,6 @@ export class WorkflowRunner { ); } - private getRelaycastBaseUrl(): string { - return ( - this.relayOptions.env?.RELAYCAST_BASE_URL ?? - process.env.RELAYCAST_BASE_URL ?? - 'https://api.relaycast.dev' - ); - } - - private getRelaycastClient(): RelayCast { - if (!this.relayApiKey) { - throw new Error('No Relaycast API key available'); - } - if (!this.relaycast) { - this.relaycast = new RelayCast({ - apiKey: this.relayApiKey, - baseUrl: this.getRelaycastBaseUrl(), - }); - } - return this.relaycast; - } - - private async ensureRelaycastRunnerAgent(): Promise { - if (this.relaycastAgent) return this.relaycastAgent; - - const rc = this.getRelaycastClient(); - let registration; - try { - registration = await rc.agents.register({ name: 'WorkflowRunner', type: 'agent' }); - } catch (err) { - if (err instanceof RelayError && err.code === 'name_conflict') { - registration = await rc.agents.register({ - name: `WorkflowRunner-${randomBytes(4).toString('hex')}`, - type: 'agent', - }); - } else { - throw err; - } - } - - this.relaycastAgent = rc.as(registration.token); - return this.relaycastAgent; - } - - private async createAndJoinRelaycastChannel(channel: string, topic?: string): Promise { - const agent = await this.ensureRelaycastRunnerAgent(); - try { - await agent.channels.create({ name: channel, ...(topic ? { topic } : {}) }); - } catch (err) { - if (!(err instanceof RelayError && err.code === 'name_conflict')) { - throw err; - } - } - await agent.channels.join(channel); - } - - private async registerRelaycastExternalAgent(name: string, persona?: string): Promise { - const rc = this.getRelaycastClient(); - try { - const registration = await rc.agents.register({ - name, - type: 'agent', - ...(persona ? { persona } : {}), - }); - return rc.as(registration.token); - } catch (err) { - if (err instanceof RelayError && err.code === 'name_conflict') { - return null; - } - throw err; - } - } - - private startRelaycastHeartbeat(agent: AgentClient, intervalMs = 30_000): () => void { - const beat = () => { - agent.heartbeat().catch(() => {}); - }; - const timer = setInterval(beat, intervalMs); - timer.unref(); - beat(); - return () => clearInterval(timer); - } - // ── Event subscription ────────────────────────────────────────────────── on(listener: WorkflowEventListener): () => void { @@ -4258,32 +3675,33 @@ export class WorkflowRunner { if (requiresBroker) { if (!relaycastDisabled) { this.log('Resolving Relaycast API key...'); - await this.ensureRelaycastApiKey(channel); + await this.brokerTransport.ensureApiKey(channel); this.log('API key resolved'); - if (this.relayApiKeyAutoCreated) { + if (this.brokerTransport.apiKeyAutoCreated) { for (const line of formatObserverGuidance(channel)) { this.log(line); } } } - this.currentBrokerContext = { + const brokerContext: BrokerRunContext = { runId, brokerName: this.buildBrokerName(runId), channel, relaycastDisabled, }; - await this.startBroker(this.currentBrokerContext); - - this.relaycast = undefined; - this.relaycastAgent = undefined; + await this.brokerTransport.start(brokerContext, { + onEvent: this.createBrokerEventHandler(runId), + onLog: (message) => this.log(message), + getActiveAgentNames: () => [...this.activeAgentHandles.keys()], + }); if (!relaycastDisabled) { this.log(`Creating channel: ${channel}...`); if (isResume) { - await this.createAndJoinRelaycastChannel(channel); + await this.brokerTransport.createAndJoinChannel(channel); } else { - await this.createAndJoinRelaycastChannel(channel, workflow.description); + await this.brokerTransport.createAndJoinChannel(channel, workflow.description); } this.log('Channel ready'); @@ -4453,10 +3871,6 @@ export class WorkflowRunner { this.ptyOutputBuffers.clear(); this.ptyListeners.clear(); - this.unsubBrokerStderr?.(); - this.unsubBrokerStderr = undefined; - - this.clearRelayListeners(); this.lastIdleLog.clear(); this.lastActivity.clear(); this.clearPendingHumanQuestionDrafts(); @@ -4471,11 +3885,7 @@ export class WorkflowRunner { this.log('Shutting down broker...'); await this.shutdownRelay(); - this.currentBrokerContext = undefined; - this.brokerRecoveryPromise = undefined; this.runStartTime = undefined; - this.relaycast = undefined; - this.relaycastAgent = undefined; this.channel = undefined; this.trajectory = undefined; this.abortController = undefined; @@ -6378,14 +5788,12 @@ export class WorkflowRunner { } private async releaseStaleRetryAgents(baseRequestedName: string, stepName: string): Promise { - if (!this.relay) { + if (!this.brokerTransport.connected) { return; } const staleAgents = ( - await this.withBrokerRecovery(`listing stale retry agents for step "${stepName}"`, (relay) => - relay.listAgents() - ) + await this.brokerTransport.listAgents(`listing stale retry agents for step "${stepName}"`) ).filter((agent) => agent.name === baseRequestedName || agent.name.startsWith(`${baseRequestedName}-r`)); if (staleAgents.length === 0) { return; @@ -6395,17 +5803,17 @@ export class WorkflowRunner { this.log(`[${stepName}] Releasing stale retry agent(s): ${staleNames.join(', ')}`); for (const name of staleNames) { - await this.withBrokerRecovery(`releasing stale retry agent "${name}"`, (relay) => - relay.release(name, `workflow retry cleanup for step "${stepName}"`) + await this.brokerTransport.release( + name, + `workflow retry cleanup for step "${stepName}"`, + `releasing stale retry agent "${name}"` ); } const deadline = Date.now() + 5_000; while (Date.now() < deadline) { const remaining = ( - await this.withBrokerRecovery(`confirming retry cleanup for step "${stepName}"`, (relay) => - relay.listAgents() - ) + await this.brokerTransport.listAgents(`confirming retry cleanup for step "${stepName}"`) ) .map((agent) => agent.name) .filter((name) => staleNames.includes(name)); @@ -6646,6 +6054,14 @@ export class WorkflowRunner { }; } catch (error) { const message = error instanceof Error ? error.message : String(error); + // A worker can settle on the same turn that the owner fails. Give its + // completion handler one turn to mark the handle released before issuing + // cleanup, otherwise the transport boundary can observe a duplicate + // release under scheduler contention. + await Promise.race([ + workerSettled, + new Promise((resolve) => setTimeout(resolve, 0)), + ]); if (!workerReleased && workerHandle) { await workerHandle.release().catch(() => undefined); } @@ -7280,6 +6696,13 @@ export class WorkflowRunner { logicalName: reviewerDef.name, onSpawned: ({ agent }) => { reviewerHandle = agent; + // A fast broker can stream the decision before the spawn response + // reaches this callback. Complete the deferred release once the + // handle becomes available so the review cannot wait forever. + if (completedReview && !reviewerReleased) { + reviewerReleased = true; + void agent.release().catch(() => undefined); + } }, onChunk: ({ chunk }) => { const nextOutput = reviewOutput + WorkflowRunner.stripAnsi(chunk); @@ -7661,17 +7084,14 @@ export class WorkflowRunner { // Register agent in Relaycast for observability let stopHeartbeat: (() => void) | undefined; - if (this.relayApiKey) { - const agentClient = await this.registerRelaycastExternalAgent( + if (this.brokerTransport.apiKey) { + stopHeartbeat = await this.brokerTransport.startExternalAgentHeartbeat( agentName, `Non-interactive workflow agent for step "${step.name}" (${agentCli})` ).catch((err) => { console.warn(`[WorkflowRunner] Failed to register ${agentName} in Relaycast:`, err?.message ?? err); - return null; + return undefined; }); - if (agentClient) { - stopHeartbeat = this.startRelaycastHeartbeat(agentClient); - } } // Post assignment notification (no task content — task arrives via direct broker injection) @@ -7964,7 +7384,7 @@ export class WorkflowRunner { const interactiveSpawnPolicy = resolveSpawnPolicy({ AGENT_NAME: agentName, AGENT_CLI: agentCli, - RELAY_API_KEY: this.relayApiKey ?? 'workflow-runner', + RELAY_API_KEY: this.brokerTransport.apiKey ?? 'workflow-runner', AGENT_CHANNELS: (agentChannels ?? []).join(','), }); const proxyMode = await this.resolveAgentProxyMode(agentDef, this.currentConfig); @@ -8000,11 +7420,12 @@ export class WorkflowRunner { `[${step.name}] Spawning ${personaResolution ? `persona ${personaResolution.resolved.spec.id}` : agentCli} (pty)` ); agent = new WorkflowAgentHandle( - await this.withBrokerRecovery(`spawning agent for step "${step.name}"`, (relay) => - relay.spawnPty({ + await this.brokerTransport.spawnPty( + { ...(spawnOptions as Record), cli: agentCli, - } as SpawnPtyInput) + } as import('@agent-relay/harness-driver').SpawnPtyInput, + `spawning agent for step "${step.name}"` ) ); @@ -8020,9 +7441,8 @@ export class WorkflowRunner { `Persona "${personaResolution.resolved.spec.id}" failed harness readiness (${ready})` ); } - const registered = await this.withBrokerRecovery( - `verifying broker registration for persona step "${step.name}"`, - (relay) => relay.listAgents() + const registered = await this.brokerTransport.listAgents( + `verifying broker registration for persona step "${step.name}"` ); if (!registered.some((candidate) => candidate.name === agent?.name)) { throw new Error( @@ -8066,8 +7486,8 @@ export class WorkflowRunner { // Register in workers.json so `agents:kill` can find this agent let workerPid: number | undefined; try { - const rawAgents = await this.withBrokerRecovery(`listing spawned agents for step "${step.name}"`, (relay) => - relay.listAgents() + const rawAgents = await this.brokerTransport.listAgents( + `listing spawned agents for step "${step.name}"` ); workerPid = rawAgents.find((a) => a.name === agentName)?.pid ?? undefined; } catch { @@ -8076,8 +7496,8 @@ export class WorkflowRunner { this.registerWorker(agentName, agentCli, step.task ?? '', workerPid); // Register the spawned agent in Relaycast for observability + start heartbeat - if (this.relayApiKey) { - const agentClient = await this.registerRelaycastExternalAgent( + if (this.brokerTransport.apiKey) { + stopHeartbeat = await this.brokerTransport.startExternalAgentHeartbeat( liveAgent.name, `Workflow agent for step "${step.name}" (${agentCli})` ).catch((err) => { @@ -8085,19 +7505,13 @@ export class WorkflowRunner { `[WorkflowRunner] Failed to register ${liveAgent.name} in Relaycast:`, err?.message ?? err ); - return null; + return undefined; }); - - // Keep the agent online in the dashboard while it's working - if (agentClient) { - stopHeartbeat = this.startRelaycastHeartbeat(agentClient); - } } // Invite the spawned agent to the workflow channel - if (this.channel && this.relayApiKey) { - const channelAgent = await this.ensureRelaycastRunnerAgent().catch(() => null); - await channelAgent?.channels.invite(this.channel, agent.name).catch(() => {}); + if (this.channel && this.brokerTransport.apiKey) { + await this.brokerTransport.inviteAgent(this.channel, agent.name).catch(() => {}); } // Keep operational assignment chatter out of the agent coordination channel. @@ -8720,12 +8134,13 @@ export class WorkflowRunner { .join('\n'); for (const agentName of targets) { - await this.withBrokerRecovery(`injecting Relayfile event into "${agentName}"`, (relay) => - relay.sendMessage({ + await this.brokerTransport.sendMessage( + { from: 'workflow-runner', to: agentName, text, - }) + }, + `injecting Relayfile event into "${agentName}"` ); } } @@ -10500,34 +9915,25 @@ export class WorkflowRunner { `Cannot inject ${input.source} answer into "${input.agentName}" because that agent is not active` ); } - if (!this.relay) { + if (!this.brokerTransport.connected) { throw new Error('Cannot inject human answer because the workflow broker is not connected'); } this.log(`[${input.stepName}] Injecting ${input.source} answer into ${input.agentName}`); - if (typeof this.relay.sendInput === 'function') { - void this.relay.sendInput(input.agentName, `${input.text}\r`).catch((err: unknown) => { + void this.brokerTransport + .sendInput( + input.agentName, + input.text, + `injecting ${input.source} answer into "${input.agentName}"` + ) + .catch((err: unknown) => { const message = err instanceof Error ? err.message : String(err); if (/aborted due to timeout|operation was aborted|timeout/i.test(message)) return; this.log(`[${input.stepName}] PTY input dispatch reported an error for ${input.agentName}: ${message}`); }); - // The current broker writes PTY input before replying, but does not ack the - // successful write. Give the dispatch a moment to leave this process. - await this.delay(500); - } else { - await Promise.race([ - this.withBrokerRecovery(`injecting ${input.source} answer into "${input.agentName}"`, (relay) => - relay.sendMessage({ - from: 'workflow-runner', - to: input.agentName, - text: input.text, - }) - ), - this.delay(20_000).then(() => { - throw new Error(`Timed out injecting ${input.source} answer into "${input.agentName}"`); - }), - ]); - } + // The current broker writes PTY input before replying, but does not ack the + // successful write. Give the dispatch a moment to leave this process. + await this.delay(500); this.log(`[${input.stepName}] Injected ${input.source} answer into ${input.agentName}`); this.postToChannel(`**[${input.stepName}]** Injected ${input.source} answer into \`${input.agentName}\``); } @@ -10822,18 +10228,19 @@ export class WorkflowRunner { agentDef: AgentDefinition, step: WorkflowStep ): Promise { - if (!this.relay) return; + if (!this.brokerTransport.connected) return; const hubAgent = this.resolveHubForNudge(agentDef); if (hubAgent) { // Hub-mediated: tell the hub to check on the idle agent (sent as the hub). try { - await this.withBrokerRecovery(`nudging idle agent "${agent.name}" via hub`, (relay) => - relay.sendMessage({ + await this.brokerTransport.sendMessage( + { from: hubAgent.name, to: agent.name, text: `Agent ${agent.name} appears idle on step "${step.name}". Check on them and remind them to /exit when done.`, - }) + }, + `nudging idle agent "${agent.name}" via hub` ); return; // Hub nudge succeeded } catch { @@ -10842,13 +10249,15 @@ export class WorkflowRunner { } // Direct system injection from the workflow runner. - await this.withBrokerRecovery(`nudging idle agent "${agent.name}" directly`, (relay) => - relay.sendMessage({ - from: 'workflow-runner', - to: agent.name, - text: "You appear idle. If you've completed your task, output /exit. If still working, continue.", - }) - ) + await this.brokerTransport + .sendMessage( + { + from: 'workflow-runner', + to: agent.name, + text: "You appear idle. If you've completed your task, output /exit. If still working, continue.", + }, + `nudging idle agent "${agent.name}" directly` + ) .catch(() => { // Non-critical — don't break workflow }); @@ -11263,7 +10672,7 @@ export class WorkflowRunner { /** Post a message to the workflow channel. Fire-and-forget — never throws or blocks. */ private postToChannel(text: string, options: ChannelEvidenceOptions = {}): void { - if (!this.relayApiKey || !this.channel) return; + if (!this.brokerTransport.apiKey || !this.channel) return; this.recordChannelEvidence(text, options); const stepName = options.stepName ?? this.inferStepNameFromChannelText(text); @@ -11280,8 +10689,8 @@ export class WorkflowRunner { }); } - this.ensureRelaycastRunnerAgent() - .then((agent) => agent.send(this.channel!, text)) + this.brokerTransport + .postToChannel(this.channel, text) .catch(() => { // Non-critical — don't break workflow execution }); From 33d5120a170ec292c111b7c8f09c6bfc6f6c59d6 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Wed, 26 Aug 2026 14:52:05 +0200 Subject: [PATCH 2/3] fix: address broker-transport review findings - injectAnswerToAgent awaits the transport with a bounded ack window: a rejected message-fallback delivery now throws instead of being logged as injected, and PTY dispatches log confirmed-vs-best-effort honestly - HarnessBrokerTransport merges caller-configured relay.env as the base of the spawned broker environment so embedders without resolveRelayEnv keep their credentials and runtime settings - the defaults-to-legacy test pins RELAYFLOWS_INTEGRATION_TRANSPORT so it cannot depend on the invoking shell Co-Authored-By: Claude Opus 4.7 --- .../src/__tests__/broker-transport.test.ts | 16 +++++-- packages/core/src/broker-transport.ts | 5 ++ packages/core/src/runner.ts | 46 +++++++++++++------ 3 files changed, 48 insertions(+), 19 deletions(-) diff --git a/packages/core/src/__tests__/broker-transport.test.ts b/packages/core/src/__tests__/broker-transport.test.ts index d0ec5b4..018ada8 100644 --- a/packages/core/src/__tests__/broker-transport.test.ts +++ b/packages/core/src/__tests__/broker-transport.test.ts @@ -8,10 +8,18 @@ import { WorkflowRunner } from '../runner.js'; describe('broker transport selection', () => { it('defaults to legacy mode', () => { - expect(resolveBrokerTransportMode(undefined, {})).toBe('legacy'); - const runner = new WorkflowRunner(); - expect((runner as any).brokerTransport).toBeInstanceOf(HarnessBrokerTransport); - expect((runner as any).brokerTransport.mode).toBe('legacy'); + // WorkflowRunner resolves the mode from process.env when no selector is + // passed; pin the variable so the assertion cannot depend on the shell. + const saved = process.env.RELAYFLOWS_INTEGRATION_TRANSPORT; + delete process.env.RELAYFLOWS_INTEGRATION_TRANSPORT; + try { + expect(resolveBrokerTransportMode(undefined, {})).toBe('legacy'); + const runner = new WorkflowRunner(); + expect((runner as any).brokerTransport).toBeInstanceOf(HarnessBrokerTransport); + expect((runner as any).brokerTransport.mode).toBe('legacy'); + } finally { + if (saved !== undefined) process.env.RELAYFLOWS_INTEGRATION_TRANSPORT = saved; + } }); it('prefers the per-run selector over the environment selector', () => { diff --git a/packages/core/src/broker-transport.ts b/packages/core/src/broker-transport.ts index 32c188a..ede6dde 100644 --- a/packages/core/src/broker-transport.ts +++ b/packages/core/src/broker-transport.ts @@ -655,7 +655,12 @@ export class HarnessBrokerTransport implements BrokerTransportPort { } this.hooks?.onLog('Starting broker...'); + // relayOptions.env is the caller-configured base (credentials, runtime + // settings); resolveRelayEnv is the runner's richer per-run resolution + // and wins where both define a key. Without the base merge, an embedder + // passing relay.env but no resolveRelayEnv loses its environment. const relayEnv = { + ...(this.relayOptions.env ?? {}), ...(this.resolveRelayEnv() ?? {}), AGENT_RELAY_STATE_DIR: stateDir, }; diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 8c503f3..9717251 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -9920,21 +9920,37 @@ export class WorkflowRunner { } this.log(`[${input.stepName}] Injecting ${input.source} answer into ${input.agentName}`); - void this.brokerTransport - .sendInput( - input.agentName, - input.text, - `injecting ${input.source} answer into "${input.agentName}"` - ) - .catch((err: unknown) => { - const message = err instanceof Error ? err.message : String(err); - if (/aborted due to timeout|operation was aborted|timeout/i.test(message)) return; - this.log(`[${input.stepName}] PTY input dispatch reported an error for ${input.agentName}: ${message}`); - }); - // The current broker writes PTY input before replying, but does not ack the - // successful write. Give the dispatch a moment to leave this process. - await this.delay(500); - this.log(`[${input.stepName}] Injected ${input.source} answer into ${input.agentName}`); + // Await the transport so a failed message-fallback delivery is never + // reported as injected. The PTY path writes before replying but a broker + // may still reply late, so the wait is bounded rather than indefinite. + let outcome: 'pty' | 'message' | 'unconfirmed'; + try { + outcome = await Promise.race([ + this.brokerTransport.sendInput( + input.agentName, + input.text, + `injecting ${input.source} answer into "${input.agentName}"` + ), + this.delay(10_000).then(() => 'unconfirmed' as const), + ]); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + this.log( + `[${input.stepName}] Failed to inject ${input.source} answer into ${input.agentName}: ${message}` + ); + throw new Error( + `Failed to inject ${input.source} answer into "${input.agentName}": ${message}` + ); + } + if (outcome === 'unconfirmed') { + this.log( + `[${input.stepName}] ${input.source} answer dispatch to ${input.agentName} not acknowledged within 10s; continuing best-effort` + ); + } else { + this.log( + `[${input.stepName}] Injected ${input.source} answer into ${input.agentName} via ${outcome}` + ); + } this.postToChannel(`**[${input.stepName}]** Injected ${input.source} answer into \`${input.agentName}\``); } From 564e7f0524636325b285061547ef7a01c7090f18 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Wed, 26 Aug 2026 15:06:26 +0200 Subject: [PATCH 3/3] fix: clear the injection ack timer so fast dispatches do not delay process exit Co-Authored-By: Claude Opus 4.7 --- packages/core/src/runner.ts | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 9717251..a9cfc19 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -9924,6 +9924,7 @@ export class WorkflowRunner { // reported as injected. The PTY path writes before replying but a broker // may still reply late, so the wait is bounded rather than indefinite. let outcome: 'pty' | 'message' | 'unconfirmed'; + let ackTimer: NodeJS.Timeout | undefined; try { outcome = await Promise.race([ this.brokerTransport.sendInput( @@ -9931,7 +9932,9 @@ export class WorkflowRunner { input.text, `injecting ${input.source} answer into "${input.agentName}"` ), - this.delay(10_000).then(() => 'unconfirmed' as const), + new Promise<'unconfirmed'>((resolve) => { + ackTimer = setTimeout(() => resolve('unconfirmed'), 10_000); + }), ]); } catch (err) { const message = err instanceof Error ? err.message : String(err); @@ -9941,6 +9944,10 @@ export class WorkflowRunner { throw new Error( `Failed to inject ${input.source} answer into "${input.agentName}": ${message}` ); + } finally { + // A pending race-loser timer would keep the event loop alive and delay + // CLI process exit after every fast dispatch. + if (ackTimer) clearTimeout(ackTimer); } if (outcome === 'unconfirmed') { this.log(