diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts index 28c241ba82f8..4981140ad992 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts @@ -5834,12 +5834,17 @@ describe("AcpAdapterV2", () => { ); const protocolEvents = yield* Queue.bounded(256); const continuationRequests: Array = []; + const promptSettled = yield* Deferred.make(); const instanceId = ProviderInstanceId.make("acp-test"); const childSessionId = "mock-child-session-post-settle"; let subagentPhase: "spawn" | "complete" = "spawn"; type RuntimeService = AcpSessionRuntime.AcpSessionRuntime["Service"]; let sessionUpdateHandler: Parameters[0] | undefined; const adapter = makeAcpAdapterV2({ + testHooks: { + afterPromptSettledWithBackgroundWork: () => + Deferred.succeed(promptSettled, undefined).pipe(Effect.asVoid), + }, crypto: yield* Crypto.Crypto, instanceId, flavor: { @@ -5929,28 +5934,13 @@ describe("AcpAdapterV2", () => { yield* runtime.startTurn( makeTurnInput({ threadId, providerThread, instanceId, runtimePolicy, now }), ); - yield* Stream.fromQueue(protocolEvents).pipe( - Stream.filter( - (event) => - event.direction === "incoming" && - event.stage === "raw" && - typeof event.payload === "string" && - event.payload.includes('"stopReason"'), - ), - Stream.runHead, - ); - yield* Effect.yieldNow; - yield* Effect.yieldNow; + yield* Deferred.await(promptSettled); const firstProviderTurnId = idAllocator.derive.providerTurn({ driver: ACP_TEST_DRIVER, nativeTurnId: acpScopedNativeId(instanceId, "mock-session-1:turn:1"), }); - const interruptFiber = yield* runtime - .interruptTurn({ providerThread, providerTurnId: firstProviderTurnId }) - .pipe(Effect.forkScoped); - yield* TestClock.adjust("10 seconds"); - yield* Fiber.join(interruptFiber); + yield* runtime.interruptTurn({ providerThread, providerTurnId: firstProviderTurnId }); let firstTerminalStatus: string | null = null; while (firstTerminalStatus === null) { @@ -7006,12 +6996,17 @@ describe("AcpAdapterV2", () => { ); const protocolEvents = yield* Queue.bounded(256); const continuationRequests: Array = []; + const promptSettled = yield* Deferred.make(); const instanceId = ProviderInstanceId.make("acp-test"); const childSessionId = "019f5470-bf92-7a90-afb3-5a6cea5b34a3"; let subagentPhase: "spawn" | "complete" = "spawn"; type RuntimeService = AcpSessionRuntime.AcpSessionRuntime["Service"]; let sessionUpdateHandler: Parameters[0] | undefined; const adapter = makeAcpAdapterV2({ + testHooks: { + afterPromptSettledWithBackgroundWork: () => + Deferred.succeed(promptSettled, undefined).pipe(Effect.asVoid), + }, crypto: yield* Crypto.Crypto, instanceId, flavor: { @@ -7110,28 +7105,13 @@ describe("AcpAdapterV2", () => { yield* runtime.startTurn( makeTurnInput({ threadId, providerThread, instanceId, runtimePolicy, now }), ); - yield* Stream.fromQueue(protocolEvents).pipe( - Stream.filter( - (event) => - event.direction === "incoming" && - event.stage === "raw" && - typeof event.payload === "string" && - event.payload.includes('"stopReason"'), - ), - Stream.runHead, - ); - yield* Effect.yieldNow; - yield* Effect.yieldNow; + yield* Deferred.await(promptSettled); const firstProviderTurnId = idAllocator.derive.providerTurn({ driver: ACP_TEST_DRIVER, nativeTurnId: acpScopedNativeId(instanceId, "mock-session-1:turn:1"), }); - const interruptFiber = yield* runtime - .interruptTurn({ providerThread, providerTurnId: firstProviderTurnId }) - .pipe(Effect.forkScoped); - yield* TestClock.adjust("10 seconds"); - yield* Fiber.join(interruptFiber); + yield* runtime.interruptTurn({ providerThread, providerTurnId: firstProviderTurnId }); let firstTerminalStatus: string | null = null; while (firstTerminalStatus === null) { @@ -8539,7 +8519,12 @@ describe("AcpAdapterV2", () => { const authoritativeSubagentText = "authoritative subagent carryover text"; type RuntimeService = AcpSessionRuntime.AcpSessionRuntime["Service"]; let sessionUpdateHandler: Parameters[0] | undefined; + const promptSettled = yield* Deferred.make(); const adapter = makeAcpAdapterV2({ + testHooks: { + afterPromptSettledWithBackgroundWork: () => + Deferred.succeed(promptSettled, undefined).pipe(Effect.asVoid), + }, crypto: yield* Crypto.Crypto, instanceId, flavor: { @@ -8629,18 +8614,7 @@ describe("AcpAdapterV2", () => { yield* runtime.startTurn( makeTurnInput({ threadId, providerThread, instanceId, runtimePolicy, now }), ); - yield* Stream.fromQueue(protocolEvents).pipe( - Stream.filter( - (event) => - event.direction === "incoming" && - event.stage === "raw" && - typeof event.payload === "string" && - event.payload.includes('"stopReason"'), - ), - Stream.runHead, - ); - yield* Effect.yieldNow; - yield* Effect.yieldNow; + yield* Deferred.await(promptSettled); // Project an authoritative v2 assistant-message upsert onto the carryover // subagent while the deferred turn is still active. diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts index d53a5f33039f..805b693b6e78 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts @@ -516,6 +516,7 @@ export interface AcpAdapterV2Options { * by exactly that on this receipt. */ readonly onDeferredFinalizeScheduled?: (debounce: Duration.Input) => Effect.Effect; + readonly afterPromptSettledWithBackgroundWork?: () => Effect.Effect; readonly afterNativeResponseTransportClosed?: () => Effect.Effect; readonly afterHardTeardownTransportDrained?: () => Effect.Effect; readonly beforeNativeResponseAdmissionCheck?: ( @@ -6935,6 +6936,9 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV // The agent finished this prompt's reply. Background work // holds the run open, not the text it already sent. yield* closeTextStreams(context); + yield* ( + options.testHooks?.afterPromptSettledWithBackgroundWork?.() ?? Effect.void + ); return; } yield* finalizeTurn(context, status);