Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
161 changes: 159 additions & 2 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ import * as AcpSessionRuntime from "../../provider/acp/AcpSessionRuntime.ts";
import {
extractXAiAcpSubagentEndNotice,
extractXAiAcpSubagentUpdate,
makeXAiPromptCompletionRuntime,
normalizeXAiAcpToolCallState,
registerXAiBackgroundTaskTracking,
} from "../../provider/acp/XAiAcpExtension.ts";
Expand Down Expand Up @@ -2975,7 +2976,11 @@ describe("AcpAdapterV2", () => {
idAllocator,
serverConfig,
selfInvocation: yield* resolveSelfInvocation(),
makeRuntime: makeMockRuntime({ childProcessSpawner, mockAgentPath, protocolEvents }),
// Production Grok runtimes are wrapped by the x.ai prompt runtime.
makeRuntime: (input) =>
makeMockRuntime({ childProcessSpawner, mockAgentPath, protocolEvents })(input).pipe(
Effect.flatMap(makeXAiPromptCompletionRuntime),
),
});
yield* adapter.openSession({
threadId: ThreadId.make(`grok-model-${model}`),
Expand Down Expand Up @@ -3027,7 +3032,11 @@ describe("AcpAdapterV2", () => {
idAllocator,
serverConfig,
selfInvocation: yield* resolveSelfInvocation(),
makeRuntime: makeMockRuntime({ childProcessSpawner, mockAgentPath, protocolEvents }),
// Production Grok runtimes are wrapped by the x.ai prompt runtime.
makeRuntime: (input) =>
makeMockRuntime({ childProcessSpawner, mockAgentPath, protocolEvents })(input).pipe(
Effect.flatMap(makeXAiPromptCompletionRuntime),
),
});
const threadId = ThreadId.make("grok-model-switch-back");
const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({
Expand Down Expand Up @@ -6078,6 +6087,154 @@ describe("AcpAdapterV2", () => {
}).pipe(Effect.provide(testLayer), Effect.scoped),
);

it.effect("finishes a settled root's carryover subagent from its structured end", () =>
Effect.gen(function* () {
const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const fileSystem = yield* FileSystem.FileSystem;
const idAllocator = yield* IdAllocatorV2;
const path = yield* Path.Path;
const serverConfig = yield* ServerConfig;
const selfInvocation = yield* resolveSelfInvocation();
const mockAgentPath = yield* path.fromFileUrl(
new URL("../../../scripts/acp-mock-agent.ts", import.meta.url),
);
const protocolEvents = yield* Queue.bounded<EffectAcpProtocol.AcpProtocolLogEvent>(256);
const instanceId = ProviderInstanceId.make("acp-test");
const childSessionId = "019f44a6-4820-7402-925d-bc862ee711dd";
let finishSubagent: AcpAdapterV2ExtensionContext["finishSubagent"] | undefined;
const adapter = makeAcpAdapterV2({
crypto: yield* Crypto.Crypto,
instanceId,
flavor: {
driver: ACP_TEST_DRIVER,
capabilities: AcpProviderCapabilitiesV2,
enablePostSettleContinuation: true,
// The spawn tool reports a background subagent that keeps running.
extractSubagentUpdate: (toolCall) =>
toolCall.toolCallId !== "tool-call-generic-1"
? undefined
: {
nativeTaskId: "task-generic-1",
prompt: "background subagent",
title: "background subagent",
model: null,
status: "running",
childSessionId,
result: null,
},
registerExtensions: (context) =>
Effect.sync(() => {
finishSubagent = context.finishSubagent;
}),
makeRuntime: makeMockRuntime({
childProcessSpawner,
mockAgentPath,
environment: { T3_ACP_EMIT_GENERIC_TOOL_PLACEHOLDERS: "1" },
protocolEvents,
}),
},
fileSystem,
idAllocator,
serverConfig,
selfInvocation,
continuationRequests: { offer: () => Effect.void },
});
const threadId = ThreadId.make("thread-acp-carryover-subagent-finished");
const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({
runtimeMode: "full-access",
interactionMode: "default",
cwd: process.cwd(),
});
const modelSelection = { instanceId, model: "default" } as const;
const runtime = yield* adapter.openSession({
threadId,
providerSessionId: ProviderSessionId.make("provider-session-acp-carryover-finished"),
modelSelection,
runtimePolicy,
});
if (runtime.hasPendingBackgroundWork === undefined) {
return yield* Effect.die("post-settle continuation must expose hasPendingBackgroundWork");
}
const hasPendingBackgroundWork = runtime.hasPendingBackgroundWork;
const events = yield* Queue.unbounded<ProviderAdapterV2Event>();
yield* runtime.events.pipe(
Stream.runForEach((event) => Queue.offer(events, event)),
Effect.forkScoped,
);
const providerThread = yield* runtime.ensureThread({
threadId,
modelSelection,
runtimePolicy,
});
yield* runtime.startTurn(
makeTurnInput({
threadId,
providerThread,
instanceId,
runtimePolicy,
now: yield* DateTime.now,
}),
);
const providerTurnId = idAllocator.derive.providerTurn({
driver: ACP_TEST_DRIVER,
nativeTurnId: acpScopedNativeId(instanceId, "mock-session-1:turn:1"),
});
let rootStatus: string | null = null;
while (rootStatus === null) {
const event = yield* Queue.take(events);
if (event.type === "turn.terminal" && event.providerTurnId === providerTurnId) {
rootStatus = event.status;
}
}
// The root completed with the subagent still running: it is carryover.
assert.equal(rootStatus, "completed");
assert.isTrue(yield* hasPendingBackgroundWork);
assert.isDefined(finishSubagent);
const subagentStatuses = Effect.gen(function* () {
// Adapter events reach this queue through the events stream fiber.
yield* Effect.yieldNow;
yield* Effect.yieldNow;
const statuses: Array<[string, string | null]> = [];
let polled = yield* Queue.poll(events);
while (Option.isSome(polled)) {
const event = polled.value;
if (event.type === "turn_item.updated" && event.turnItem.type === "subagent") {
statuses.push([event.turnItem.status, event.turnItem.result]);
}
polled = yield* Queue.poll(events);
}
return statuses;
});
yield* subagentStatuses;

// A nested subagent reports to its own parent session, not the root.
yield* finishSubagent!({
sessionId: "some-other-session",
childSessionId,
status: "completed",
result: "WRONG_PARENT",
});
assert.deepEqual(yield* subagentStatuses, [], "non-root notices must be dropped");
assert.isTrue(yield* hasPendingBackgroundWork);

yield* finishSubagent!({
sessionId: "mock-session-1",
childSessionId,
status: "failed",
result: "tool crashed",
});
assert.deepEqual(
yield* subagentStatuses,
[["failed", "tool crashed"]],
"the completed root still owns the run, so the carryover end projects at once",
);
assert.isFalse(
yield* hasPendingBackgroundWork,
"a finished carryover subagent stops pinning",
);
}).pipe(Effect.provide(testLayer), Effect.scoped),
);

it.effect("projects completed-root carryover eagerly and drain cannot resurrect it", () =>
Effect.gen(function* () {
const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner;
Expand Down
61 changes: 61 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,17 @@ export interface AcpAdapterV2ExtensionContext {
readonly sessionId: string;
readonly failure: OrchestrationV2ProviderFailure;
}) => Effect.Effect<void>;
/**
* A subagent's structured end on the root session (Grok `subagent_finished`),
* keyed by its child session id. Finishes the subagent row, in the turn that
* holds it or in the carryover of a settled one.
*/
readonly finishSubagent: (notice: {
readonly sessionId: string;
readonly childSessionId: string;
readonly status: "completed" | "failed" | "cancelled";
readonly result: string | null;
}) => Effect.Effect<void>;
/**
* Session-scoped background-task lifecycle reported via extension
* notifications (e.g. Grok `x.ai/task_backgrounded`; older builds use the
Expand Down Expand Up @@ -5005,6 +5016,45 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
return true;
});

const finishSubagentFromNotice = Effect.fnUntraced(function* (notice: {
readonly childSessionId: string;
readonly status: "completed" | "failed" | "cancelled";
readonly result: string | null;
}) {
const context = yield* Ref.get(activeTurn);
const subagent =
context === null ? undefined : context.subagentsBySessionId.get(notice.childSessionId);
if (context !== null && subagent !== undefined && !context.finalized) {
if (!acpSubagentStatusBlocksTurnSettlement(subagent.task.status)) return;
yield* emitSubagent(context, {
nativeTaskId: subagent.task.nativeTaskRef?.nativeId ?? notice.childSessionId,
prompt: subagent.task.prompt,
title: subagent.task.title,
model: subagent.task.model,
status: notice.status,
childSessionId: notice.childSessionId,
result: notice.result,
suppressNormalTool: true,
});
yield* rearmDeferredFinalize(context);
return;
}
// The root turn already settled: the subagent is carryover. Project
// its end while the completed root still owns the run.
const carryover = yield* Ref.get(carryoverSubagents);
yield* updateCarryoverSubagentStatus(
notice.childSessionId,
notice.status,
notice.result,
{
project:
carryover !== null &&
carryover.sessionId === (yield* Ref.get(activeSessionId)) &&
carryover.rootTerminalStatus === "completed",
},
);
});

applyFinalizedActiveTurnSubagentTerminal = Effect.fnUntraced(function* (
context: ActiveAcpTurn,
notification: EffectAcpSchema.SessionNotification,
Expand Down Expand Up @@ -5631,6 +5681,17 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
requestUserInput,
captureProposedPlan,
lastProposedPlanMarkdown,
finishSubagent: (notice) =>
runRuntimeCallbackAtGeneration(
handlerGeneration,
Effect.gen(function* () {
if (yield* Ref.get(stoppedRunQuarantine)) return;
// Root-session notices only; nested subagents report to
// their own parent session.
if ((yield* Ref.get(activeSessionId)) !== notice.sessionId) return;
yield* finishSubagentFromNotice(notice);
}),
).pipe(Effect.asVoid),
applyBackgroundTaskMutation: (mutation) =>
runRuntimeCallbackAtGeneration(
handlerGeneration,
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/GrokAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ import {
extractXAiMonitorTaskId,
isXAiPersistentMonitor,
extractXAiExitPlanMarkdown,
registerXAiSubagentFinished,
makeXAiAskUserQuestionCancelledResponse,
makeXAiAskUserQuestionResponse,
makeXAiExitPlanModeCapturedResponse,
Expand Down Expand Up @@ -127,10 +128,12 @@ const registerGrokAcpExtensions: NonNullable<AcpAdapterV2Flavor["registerExtensi
runtime,
requestUserInput,
applyBackgroundTaskMutation,
finishSubagent,
captureProposedPlan,
lastProposedPlanMarkdown,
}) =>
registerXAiBackgroundTaskTracking(runtime, applyBackgroundTaskMutation).pipe(
Effect.andThen(registerXAiSubagentFinished(runtime, finishSubagent)),
Effect.andThen(registerGrokAskUserQuestionExtensions({ runtime, requestUserInput })),
Effect.andThen(
registerGrokExitPlanModeExtensions({
Expand Down
Loading
Loading