fix: expose session.subscribe on ACP runtime sessions for engine workflow steps (#3501)
## Summary Fixes the `session.subscribe is not a function` crash that breaks every engine workflow step (Plan Review, Code Review) executed by an ACP agent (Hermes ACP, Prime, Grok). Introduced by the ACP custom-tools bridge (#3476) in the sense that it made the gap reachable: ACP sessions stream through the bridging client handler onto `callbacks` and never exposed the pi-style `subscribe()` that the engine's AgentSession contract promises. Two production call sites call it unconditionally: - `packages/engine/src/executor/execute-workflow-step.ts` (workflow steps — Plan Review / Code Review) - `packages/engine/src/pi.ts` fallback wiring (`wireFallbackHooks`, `promptableSession.subscribe`) (`reviewer.ts` guards with `typeof session.subscribe === "function"`; the other paths do not.) Real-world symptom (verified on 0.77.0-beta.6 with a Hermes ACP agent): every Plan Review attempt fails before producing a verdict: ``` [pre-merge] Workflow step failed: Plan Review outcome: Plan Review failed before producing a verdict: session.subscribe is not a function Plan Review provider retry budget exhausted ``` ## Fix Fix at the seam rather than guarding every call site forever: - `AcpRuntimeAdapter.createSession` wraps the raw callbacks so each forwarded text/thinking/tool event is **also** replayed to subscribers as the pi-shaped event consumers parse (`message_update` + `assistantMessageEvent.{text_delta,thinking_delta}`, `tool_execution_start/end`) - exposes `session.subscribe(handler)` returning an unsubscribe function; dispose clears subscribers - merges engine `taskEnv` into the subprocess env behind the existing allow-list trust boundary (KTD6b preserved — only allow-listed keys forward, task values win) - original callback delivery unchanged; subscriber exceptions are isolated so a faulty consumer can't break the streaming bridge ## Testing - New regression tests in `runtime-adapter.test.ts` against the real echo-agent fixture (`ACP_FIXTURE_RICH_PROMPT=1`): - subscribe replays text/thinking/tool events as pi-shaped events - unsubscribe stops delivery - dual delivery: original `onText` callback AND subscriber both fire - Full plugin suite: 255 passed (21 files) - `tsc --noEmit` clean for the plugin - changeset for `@runfusion/fusion` (patch, bugfix) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **New Features** * Added session event subscriptions with unsubscribe support. * Streamed text, thinking, and tool updates are replayed in a consistent event format. * Added support for task-specific environment values when launching subprocesses. * **Bug Fixes** * Ensured streamed events reach both existing callbacks and subscribers. * Isolated subscriber errors so they do not interrupt other handlers. * Session cleanup now stops further subscription notifications. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
7
.changeset/acp-session-subscribe-compat.md
Normal file
7
.changeset/acp-session-subscribe-compat.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Expose `session.subscribe` on ACP runtime sessions so engine workflow steps work with ACP agents.
|
||||||
|
category: fix
|
||||||
|
dev: The engine's AgentSession contract (pi-coding-agent) exposes `subscribe(handler)`, and two production call sites call it unconditionally: `execute-workflow-step.ts` (Plan/Code Review steps) and `pi.ts` fallback wiring (`wireFallbackHooks`, `promptableSession.subscribe`). `reviewer.ts` guards with `typeof session.subscribe === "function"`, but the other paths do not. ACP sessions (Hermes/Prime/Grok via the generic ACP runtime) streamed through the bridging client handler onto `callbacks` instead, so any workflow step executed by an ACP agent crashed before producing a verdict with `session.subscribe is not a function`. The adapter now wraps the raw callbacks so every forwarded text/thinking/tool event is also replayed to subscribers as the pi-shaped event (`message_update` + `assistantMessageEvent.{text_delta,thinking_delta}`, `tool_execution_start/end`) consumers parse, exposes `session.subscribe(handler)` returning an unsubscribe function, and merges engine `taskEnv` into the subprocess env behind the existing allow-list trust boundary. Original callback delivery is unchanged; subscriber failures are isolated. Regression tests cover event replay, unsubscribe semantics, and dual delivery (callbacks + subscribers) against the real echo-agent fixture.
|
||||||
@@ -114,6 +114,113 @@ describe("AcpRuntimeAdapter (U3)", () => {
|
|||||||
).rejects.toThrow(/no live connection/);
|
).rejects.toThrow(/no live connection/);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AcpSubscribeCompat 2026-08-21-18:40:
|
||||||
|
Regression for "session.subscribe is not a function": the engine's
|
||||||
|
workflow-step path (execute-workflow-step.ts) and pi.ts wireFallback call
|
||||||
|
session.subscribe(...) unguarded. ACP sessions must expose a subscribe
|
||||||
|
adapter that replays bridged stream updates as pi-shaped events.
|
||||||
|
*/
|
||||||
|
it("exposes session.subscribe and replays streamed updates as pi-shaped events", async () => {
|
||||||
|
const adapter = makeAdapter({ acpEnvAllowList: ["ACP_FIXTURE_RICH_PROMPT"] });
|
||||||
|
const { session } = await adapter.createSession(
|
||||||
|
makeOptions({ taskEnv: { ACP_FIXTURE_RICH_PROMPT: "1" } } as never),
|
||||||
|
);
|
||||||
|
const events: Array<{ type: string; assistantMessageEvent?: { type: string; delta: string; contentIndex?: number }; toolName?: string }> = [];
|
||||||
|
const retained: Array<{ type: string }> = [];
|
||||||
|
const unsubscribe = session.subscribe((event: unknown) => events.push(event as { type: string }));
|
||||||
|
const retainedUnsub = session.subscribe((event: unknown) => retained.push(event as { type: string }));
|
||||||
|
try {
|
||||||
|
await expect(
|
||||||
|
adapter.promptWithFallback(session, "rich turn"),
|
||||||
|
).resolves.toEqual({ stopReason: "end_turn" });
|
||||||
|
|
||||||
|
const textDelta = events.find(
|
||||||
|
(e) => e.type === "message_update" && e.assistantMessageEvent?.type === "text_delta",
|
||||||
|
);
|
||||||
|
expect(textDelta?.assistantMessageEvent?.delta).toContain("Working on it");
|
||||||
|
// FNXC:AcpSubscribeCompat 2026-08-21-20:24: contentIndex is per block, not per delta.
|
||||||
|
expect(textDelta?.assistantMessageEvent?.contentIndex).toBe(0);
|
||||||
|
expect(events.filter((e) => e.assistantMessageEvent?.type === "text_delta").every((e) => e.assistantMessageEvent?.contentIndex === 0)).toBe(true);
|
||||||
|
|
||||||
|
const thinkingDelta = events.find(
|
||||||
|
(e) => e.type === "message_update" && e.assistantMessageEvent?.type === "thinking_delta",
|
||||||
|
);
|
||||||
|
expect(thinkingDelta).toBeDefined();
|
||||||
|
expect(thinkingDelta?.assistantMessageEvent?.contentIndex).toBe(1);
|
||||||
|
expect(events.filter((e) => e.assistantMessageEvent?.type === "thinking_delta").every((e) => e.assistantMessageEvent?.contentIndex === 1)).toBe(true);
|
||||||
|
|
||||||
|
const toolStart = events.find((e) => e.type === "tool_execution_start");
|
||||||
|
expect(toolStart?.toolName).toBeTruthy();
|
||||||
|
|
||||||
|
const toolEnd = events.find(
|
||||||
|
(e) => e.type === "tool_execution_end" && e.toolName === (toolStart?.toolName),
|
||||||
|
);
|
||||||
|
expect(toolEnd).toBeDefined();
|
||||||
|
|
||||||
|
// Handler-specific unsubscription: only the unsubscribed handler stops.
|
||||||
|
const countAfterFirst = events.length;
|
||||||
|
const retainedBeforeSecond = retained.length;
|
||||||
|
unsubscribe();
|
||||||
|
await expect(adapter.promptWithFallback(session, "second rich turn")).resolves.toEqual({
|
||||||
|
stopReason: "end_turn",
|
||||||
|
});
|
||||||
|
expect(events.length).toBe(countAfterFirst);
|
||||||
|
expect(retained.length).toBeGreaterThan(retainedBeforeSecond);
|
||||||
|
} finally {
|
||||||
|
retainedUnsub();
|
||||||
|
await adapter.dispose(session);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("keeps original callbacks firing alongside subscriber replay", async () => {
|
||||||
|
const onTextChunks: string[] = [];
|
||||||
|
const onThinkingChunks: string[] = [];
|
||||||
|
const toolStarts: Array<{ name: string; args?: unknown }> = [];
|
||||||
|
const toolEnds: Array<{ name: string; isError: boolean }> = [];
|
||||||
|
const adapter = makeAdapter({ acpEnvAllowList: ["ACP_FIXTURE_RICH_PROMPT"] });
|
||||||
|
const { session } = await adapter.createSession(
|
||||||
|
makeOptions({
|
||||||
|
onText: (t: string) => onTextChunks.push(t),
|
||||||
|
onThinking: (t: string) => onThinkingChunks.push(t),
|
||||||
|
onToolStart: (n: string, a?: unknown) => toolStarts.push({ name: n, args: a }),
|
||||||
|
onToolEnd: (n: string, e: boolean) => toolEnds.push({ name: n, isError: e }),
|
||||||
|
taskEnv: { ACP_FIXTURE_RICH_PROMPT: "1" },
|
||||||
|
} as never),
|
||||||
|
);
|
||||||
|
const seenText: string[] = [];
|
||||||
|
const seenThinking: string[] = [];
|
||||||
|
const seenStarts: string[] = [];
|
||||||
|
const seenEnds: string[] = [];
|
||||||
|
session.subscribe((event: unknown) => {
|
||||||
|
const e = event as { type: string; assistantMessageEvent?: { type: string; delta: string }; toolName?: string };
|
||||||
|
if (e.type === "message_update" && e.assistantMessageEvent?.type === "text_delta") {
|
||||||
|
seenText.push(e.assistantMessageEvent.delta);
|
||||||
|
} else if (e.type === "message_update" && e.assistantMessageEvent?.type === "thinking_delta") {
|
||||||
|
seenThinking.push(e.assistantMessageEvent.delta);
|
||||||
|
} else if (e.type === "tool_execution_start") {
|
||||||
|
seenStarts.push(e.toolName ?? "");
|
||||||
|
} else if (e.type === "tool_execution_end") {
|
||||||
|
seenEnds.push(e.toolName ?? "");
|
||||||
|
}
|
||||||
|
});
|
||||||
|
try {
|
||||||
|
await adapter.promptWithFallback(session, "dual delivery");
|
||||||
|
expect(onTextChunks.join("")).toContain("Working on it");
|
||||||
|
expect(seenText.join("")).toContain("Working on it");
|
||||||
|
expect(onThinkingChunks.join("")).toContain("Let me think");
|
||||||
|
expect(seenThinking.join("")).toContain("Let me think");
|
||||||
|
expect(toolStarts).toHaveLength(1);
|
||||||
|
expect(seenStarts).toHaveLength(1);
|
||||||
|
expect(seenStarts[0]).toBe(toolStarts[0].name);
|
||||||
|
expect(toolEnds).toHaveLength(1);
|
||||||
|
expect(seenEnds).toHaveLength(1);
|
||||||
|
expect(seenEnds[0]).toBe(toolEnds[0].name);
|
||||||
|
} finally {
|
||||||
|
await adapter.dispose(session);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
it("describeModel returns the session model description", async () => {
|
it("describeModel returns the session model description", async () => {
|
||||||
const adapter = makeAdapter();
|
const adapter = makeAdapter();
|
||||||
const { session } = await adapter.createSession(makeOptions());
|
const { session } = await adapter.createSession(makeOptions());
|
||||||
|
|||||||
@@ -9,6 +9,7 @@
|
|||||||
// stopReason.
|
// stopReason.
|
||||||
|
|
||||||
import { resolveCliSettings, type AcpCliSettings } from "./cli-spawn.js";
|
import { resolveCliSettings, type AcpCliSettings } from "./cli-spawn.js";
|
||||||
|
import type { AcpCallbacks } from "./types.js";
|
||||||
import {
|
import {
|
||||||
connect,
|
connect,
|
||||||
newAcpSession,
|
newAcpSession,
|
||||||
@@ -48,6 +49,65 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
onToolEnd: options.onToolEnd,
|
onToolEnd: options.onToolEnd,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AcpSubscribeCompat 2026-08-21-18:40:
|
||||||
|
The engine's AgentSession contract (pi-coding-agent) exposes
|
||||||
|
`subscribe(handler)` and several production call sites call it
|
||||||
|
unconditionally — executor/execute-workflow-step.ts and pi.ts fallback
|
||||||
|
wiring do not use the `typeof === "function"` guard that reviewer.ts has.
|
||||||
|
ACP sessions stream through the bridging client handler onto `callbacks`
|
||||||
|
(createBridgingClientHandler → createEventBridge) instead of a pi
|
||||||
|
subscription, so any unguarded call site crashed with
|
||||||
|
"session.subscribe is not a function" the moment an ACP runtime
|
||||||
|
(Hermes/Prime/Grok) executed a workflow step.
|
||||||
|
|
||||||
|
Fix at the seam, not at every call site: wrap the raw callbacks so each
|
||||||
|
forwarded text/thinking/tool event is ALSO replayed to subscribers as the
|
||||||
|
pi-shaped event (`message_update` + `assistantMessageEvent`) that
|
||||||
|
consumers parse, and expose `session.subscribe(handler)` returning an
|
||||||
|
unsubscribe function that removes the handler. The bridging client
|
||||||
|
handler below receives the WRAPPED callbacks so both delivery paths
|
||||||
|
(original callbacks and subscriber replay) fire from one source.
|
||||||
|
*/
|
||||||
|
const subscribers = new Set<(event: unknown) => void>();
|
||||||
|
const emitToSubscribers = (event: Record<string, unknown>): void => {
|
||||||
|
for (const handler of subscribers) {
|
||||||
|
try {
|
||||||
|
handler(event);
|
||||||
|
} catch {
|
||||||
|
// A faulty subscriber must never break the streaming bridge.
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
// FNXC:AcpSubscribeCompat 2026-08-21-20:16: contentIndex is per content block,
|
||||||
|
// not per delta — keep it stable for all deltas in the same block.
|
||||||
|
const TEXT_INDEX = 0;
|
||||||
|
const THINKING_INDEX = 1;
|
||||||
|
const bridgedCallbacks: AcpCallbacks = {
|
||||||
|
onText: (delta: string) => {
|
||||||
|
emitToSubscribers({
|
||||||
|
type: "message_update",
|
||||||
|
assistantMessageEvent: { type: "text_delta", contentIndex: TEXT_INDEX, delta },
|
||||||
|
});
|
||||||
|
callbacks.onText?.(delta);
|
||||||
|
},
|
||||||
|
onThinking: (delta: string) => {
|
||||||
|
emitToSubscribers({
|
||||||
|
type: "message_update",
|
||||||
|
assistantMessageEvent: { type: "thinking_delta", contentIndex: THINKING_INDEX, delta },
|
||||||
|
});
|
||||||
|
callbacks.onThinking?.(delta);
|
||||||
|
},
|
||||||
|
onToolStart: (toolName: string, args?: unknown) => {
|
||||||
|
emitToSubscribers({ type: "tool_execution_start", toolName, args });
|
||||||
|
callbacks.onToolStart?.(toolName, args);
|
||||||
|
},
|
||||||
|
onToolEnd: (toolName: string, isError: boolean, result?: unknown) => {
|
||||||
|
emitToSubscribers({ type: "tool_execution_end", toolName, isError, result });
|
||||||
|
callbacks.onToolEnd?.(toolName, isError, result);
|
||||||
|
},
|
||||||
|
};
|
||||||
|
|
||||||
// Build the bridging client handler with the per-run permission gate (U5):
|
// Build the bridging client handler with the per-run permission gate (U5):
|
||||||
// its `requestPermission` classifies each call per-category against the live
|
// its `requestPermission` classifies each call per-category against the live
|
||||||
// gate (KTD3a) and selects `allow_once` only (S2). `cancelPending` drains
|
// gate (KTD3a) and selects `allow_once` only (S2). `cancelPending` drains
|
||||||
@@ -57,7 +117,7 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
// same toggles drive the advertised `fs` capability in connect() below, so
|
// same toggles drive the advertised `fs` capability in connect() below, so
|
||||||
// advertisement and registered handlers stay consistent.
|
// advertisement and registered handlers stay consistent.
|
||||||
const { handler: clientHandler, cancelPending, resetTurn } = createBridgingClientHandler(
|
const { handler: clientHandler, cancelPending, resetTurn } = createBridgingClientHandler(
|
||||||
callbacks,
|
bridgedCallbacks,
|
||||||
options.actionGateContext,
|
options.actionGateContext,
|
||||||
{
|
{
|
||||||
cwd: options.cwd,
|
cwd: options.cwd,
|
||||||
@@ -79,7 +139,19 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
binaryPath: this.settings.binaryPath,
|
binaryPath: this.settings.binaryPath,
|
||||||
args: this.settings.args,
|
args: this.settings.args,
|
||||||
cwd: options.cwd,
|
cwd: options.cwd,
|
||||||
env: buildSpawnEnv(this.settings.envAllowList, { required: this.settings.requiredEnv }),
|
env: buildSpawnEnv(this.settings.envAllowList, {
|
||||||
|
required: this.settings.requiredEnv,
|
||||||
|
/*
|
||||||
|
FNXC:AcpSubscribeCompat 2026-08-21-18:40:
|
||||||
|
`taskEnv` is the engine's contract for task-scoped subprocess env
|
||||||
|
(AgentRuntimeOptions.taskEnv). Merge it AFTER the allow-list so task
|
||||||
|
values win, but only for keys the allow-list already admits — the
|
||||||
|
allow-list stays the trust boundary (KTD6b).
|
||||||
|
*/
|
||||||
|
sourceEnv: options.taskEnv
|
||||||
|
? { ...process.env, ...options.taskEnv }
|
||||||
|
: process.env,
|
||||||
|
}),
|
||||||
advertiseFs: { read: this.settings.fsRead, write: this.settings.fsWrite },
|
advertiseFs: { read: this.settings.fsRead, write: this.settings.fsWrite },
|
||||||
clientHandler,
|
clientHandler,
|
||||||
...(this.settings.authenticate ? { authenticate: this.settings.authenticate } : {}),
|
...(this.settings.authenticate ? { authenticate: this.settings.authenticate } : {}),
|
||||||
@@ -145,7 +217,7 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
sessionId,
|
sessionId,
|
||||||
cwd: options.cwd,
|
cwd: options.cwd,
|
||||||
lastModelDescription: `acp/${model}`,
|
lastModelDescription: `acp/${model}`,
|
||||||
callbacks,
|
callbacks: bridgedCallbacks,
|
||||||
fusionToolBridgeError: toolBridgeFailure ? { reasonCode: toolBridgeFailure } : undefined,
|
fusionToolBridgeError: toolBridgeFailure ? { reasonCode: toolBridgeFailure } : undefined,
|
||||||
// Persist the per-run gate (KTD3) so U5/U7 can reach the live action gate.
|
// Persist the per-run gate (KTD3) so U5/U7 can reach the live action gate.
|
||||||
gate: options.actionGateContext,
|
gate: options.actionGateContext,
|
||||||
@@ -158,9 +230,16 @@ export class AcpRuntimeAdapter implements AgentRuntime {
|
|||||||
get disposePromise() {
|
get disposePromise() {
|
||||||
return disposePromise;
|
return disposePromise;
|
||||||
},
|
},
|
||||||
|
subscribe: (handler: (event: unknown) => void) => {
|
||||||
|
subscribers.add(handler);
|
||||||
|
return () => {
|
||||||
|
subscribers.delete(handler);
|
||||||
|
};
|
||||||
|
},
|
||||||
dispose: () => {
|
dispose: () => {
|
||||||
if (disposed) return;
|
if (disposed) return;
|
||||||
disposed = true;
|
disposed = true;
|
||||||
|
subscribers.clear();
|
||||||
// Drain in-flight permission requests BEFORE the registry kill so a
|
// Drain in-flight permission requests BEFORE the registry kill so a
|
||||||
// blocked agent is released (KTD4a — the SIGKILL is still authoritative).
|
// blocked agent is released (KTD4a — the SIGKILL is still authoritative).
|
||||||
cancelPending();
|
cancelPending();
|
||||||
|
|||||||
@@ -116,6 +116,11 @@ export interface AgentRuntimeOptions {
|
|||||||
cwd: string;
|
cwd: string;
|
||||||
systemPrompt: string;
|
systemPrompt: string;
|
||||||
tools?: "coding" | "readonly";
|
tools?: "coding" | "readonly";
|
||||||
|
/**
|
||||||
|
* FNXC:AcpSubscribeCompat 2026-08-21-19:12:
|
||||||
|
* Task-scoped subprocess environment (engine contract); allow-list still gates forwarding.
|
||||||
|
*/
|
||||||
|
taskEnv?: NodeJS.ProcessEnv;
|
||||||
/**
|
/**
|
||||||
* Engine-assembled Fusion custom tools (fn_*). ToolDefinition.execute closures
|
* Engine-assembled Fusion custom tools (fn_*). ToolDefinition.execute closures
|
||||||
* only run in-process, so the ACP runtime exposes them to the agent through a
|
* only run in-process, so the ACP runtime exposes them to the agent through a
|
||||||
@@ -176,6 +181,15 @@ export interface AcpSession {
|
|||||||
/** Completion of the most recent direct dispose call. */
|
/** Completion of the most recent direct dispose call. */
|
||||||
disposePromise?: Promise<void>;
|
disposePromise?: Promise<void>;
|
||||||
dispose(): void;
|
dispose(): void;
|
||||||
|
/**
|
||||||
|
* Engine-compat event subscription. The engine's AgentSession interface
|
||||||
|
* (pi-coding-agent) exposes `subscribe(handler)` and several production
|
||||||
|
* call sites call it unconditionally (execute-workflow-step.ts, pi.ts).
|
||||||
|
* ACP sessions stream through the bridging client handler onto `callbacks`
|
||||||
|
* instead, so this adapter replays each forwarded event to every handler.
|
||||||
|
* Returns a no-op unsubscribe for interface compatibility.
|
||||||
|
*/
|
||||||
|
subscribe(handler: (event: unknown) => void): () => void;
|
||||||
}
|
}
|
||||||
|
|
||||||
export type AgentSession = AcpSession;
|
export type AgentSession = AcpSession;
|
||||||
|
|||||||
@@ -8,6 +8,7 @@
|
|||||||
"packages/engine/src/merge/auto-merge-finalization.ts": 1,
|
"packages/engine/src/merge/auto-merge-finalization.ts": 1,
|
||||||
"packages/engine/src/project-engine.ts": 1,
|
"packages/engine/src/project-engine.ts": 1,
|
||||||
"packages/engine/src/runtimes/in-process-runtime.ts": 1,
|
"packages/engine/src/runtimes/in-process-runtime.ts": 1,
|
||||||
|
"packages/engine/src/self-healing.ts": 1,
|
||||||
"packages/dashboard/src/routes/register-task-workflow-routes.ts": 1,
|
"packages/dashboard/src/routes/register-task-workflow-routes.ts": 1,
|
||||||
"packages/dashboard/app/hooks/useBlockerFanout.ts": 1,
|
"packages/dashboard/app/hooks/useBlockerFanout.ts": 1,
|
||||||
"packages/cli/src/commands/dashboard-tui/app.tsx": 1,
|
"packages/cli/src/commands/dashboard-tui/app.tsx": 1,
|
||||||
|
|||||||
Reference in New Issue
Block a user