diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index a64dc21cec91..02f518140800 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -1,3 +1,4 @@ +import * as EffectOutbox from "./EffectOutbox.ts"; import * as SourceControlProviderRegistry from "../sourceControl/SourceControlProviderRegistry.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; @@ -6,6 +7,7 @@ import { EventId, MessageId, type ModelSelection, + type OrchestrationV2ThreadProjection, NodeId, ProjectId, ProviderDriverKind, @@ -34,7 +36,9 @@ import * as WorkspacePaths from "../workspace/WorkspacePaths.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as EventSink from "./EventSink.ts"; import * as Orchestrator from "./Orchestrator.ts"; +import * as ProviderRuntimeRecoveryService from "./ProviderRuntimeRecoveryService.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; +import { makeSubagentChildThread } from "./SubagentProjection.ts"; import { OrchestrationV2EventSinkLayerLive, OrchestrationV2LayerLive, @@ -384,6 +388,315 @@ it.layer(TestLayer)("delegated completion delivery repairs", (it) => { }), ); + it.effect( + "recovers a child that runtime recovery cancelled, deferring pending continuations", + () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const sink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:delegated-recovery"); + const runId = RunId.make("run:delegated-recovery"); + const taskId = NodeId.make("node:delegated-recovery-task"); + const childThreadId = ThreadId.make("thread:delegated-recovery-child"); + yield* seedParentWithTerminalTask({ + threadId, + runId, + projectId: ProjectId.make("project:delegated-recovery"), + rootNodeId: NodeId.make("node:delegated-recovery-root"), + taskId, + deliveryState: "claimed", + completionWake: "always", + now, + }); + const parent = yield* orchestrator.getThreadProjection(threadId); + const parentRun = parent.runs[0]!; + const childRunId = RunId.make("run:delegated-recovery-child"); + yield* sink.write({ + commandId: CommandId.make("command:delegated-recovery:seed"), + events: [ + { + id: EventId.make("event:delegated-recovery:child"), + type: "thread.created", + threadId: childThreadId, + driver, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: makeSubagentChildThread({ + parentThread: parent.thread, + childThreadId, + parentNodeId: taskId, + activeProviderThreadId: null, + providerInstanceId: modelSelection.instanceId, + modelSelection, + title: "Recovered child", + now, + createdBy: "agent", + creationSource: "server", + }), + }, + { + id: EventId.make("event:delegated-recovery:task"), + type: "subagent.updated", + threadId, + runId, + nodeId: taskId, + driver, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { + ...parent.subagents[0]!, + childThreadId, + status: "running", + result: null, + completedAt: null, + completionDelivery: undefined, + }, + }, + ], + }); + yield* sink.write({ + commandId: CommandId.make("command:delegated-recovery:running-child"), + events: [ + { + id: EventId.make("event:delegated-recovery:child-run"), + type: "run.updated", + threadId: childThreadId, + runId: childRunId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { + ...parentRun, + id: childRunId, + threadId: childThreadId, + userMessageId: MessageId.make("message:delegated-recovery-child"), + rootNodeId: null, + activeAttemptId: null, + status: "running", + completedAt: null, + delegatedCompletion: undefined, + }, + }, + ], + }); + const resultTransfers = (projection: { + readonly contextTransfers: ReadonlyArray<{ + readonly type: string; + readonly sourceThreadId: ThreadId; + }>; + }) => + projection.contextTransfers.filter( + (row) => row.type === "subagent_result" && row.sourceThreadId === childThreadId, + ); + + // The report pass must run after real runtime recovery cancels the run. + const runtimeRecovery = + yield* ProviderRuntimeRecoveryService.ProviderRuntimeRecoveryService; + yield* orchestrator.recoverDelegatedTaskReports(); + assert.lengthOf(resultTransfers(yield* orchestrator.getThreadProjection(threadId)), 0); + yield* runtimeRecovery.recover; + const child = yield* orchestrator.getThreadProjection(childThreadId); + assert.equal(child.runs.at(-1)?.status, "cancelled"); + // The terminal reactor ignores runtime-reconcile events. + assert.lengthOf(resultTransfers(yield* orchestrator.getThreadProjection(threadId)), 0); + // An unreadable continuation must not publish a false cancellation. + yield* orchestrator.recoverDelegatedTaskReports(() => + Effect.fail( + new EffectOutbox.EffectOutboxError({ + operation: "get", + cause: new Error("outbox unavailable"), + }), + ), + ); + assert.lengthOf(resultTransfers(yield* orchestrator.getThreadProjection(threadId)), 0); + // A child whose restart continuation is still pending is left alone. + yield* orchestrator.recoverDelegatedTaskReports(() => Effect.succeed(true)); + assert.lengthOf(resultTransfers(yield* orchestrator.getThreadProjection(threadId)), 0); + + yield* orchestrator.recoverDelegatedTaskReports((run) => + runtimeRecovery.isRestartContinuationPending(run.id), + ); + const recovered = yield* orchestrator.getThreadProjection(threadId); + assert.lengthOf(resultTransfers(recovered), 1); + const task = recovered.subagents.find((row) => row.id === taskId); + assert.equal(task?.status, "cancelled"); + assert.isDefined(task?.completionDelivery); + + // Running the pass again publishes nothing new. + yield* orchestrator.recoverDelegatedTaskReports(); + assert.lengthOf(resultTransfers(yield* orchestrator.getThreadProjection(threadId)), 1); + }), + ); + + it.effect("re-offers a parent wake that startup recovery cancelled with a sibling", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const sink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:delegated-recovery-wake"); + const runId = RunId.make("run:delegated-recovery-wake"); + const taskId = NodeId.make("node:delegated-recovery-wake-first"); + const siblingId = NodeId.make("node:delegated-recovery-wake-sibling"); + const siblingThreadId = ThreadId.make("thread:delegated-recovery-wake-sibling"); + const wakeMessageId = MessageId.make(`message:delegated-delivery:${threadId}`); + const wakeRunId = RunId.make("run:delegated-recovery-wake-delivery"); + const siblingRunId = RunId.make("run:delegated-recovery-wake-sibling"); + yield* seedParentWithTerminalTask({ + threadId, + runId, + projectId: ProjectId.make("project:delegated-recovery-wake"), + rootNodeId: NodeId.make("node:delegated-recovery-wake-root"), + taskId, + deliveryState: "claimed", + completionWake: "always", + deliveryTaskIds: [taskId], + now, + }); + const parent = yield* orchestrator.getThreadProjection(threadId); + const parentRun = parent.runs[0]!; + const runPayload = ( + id: RunId, + thread: ThreadId, + messageId: MessageId, + status: "running" | "cancelled", + ) => ({ + ...parentRun, + id, + threadId: thread, + ordinal: 2, + userMessageId: messageId, + rootNodeId: null, + activeAttemptId: null, + status, + completedAt: status === "cancelled" ? now : null, + delegatedCompletion: undefined, + }); + // Task A's wake is running in the parent while sibling B still works. + yield* sink.write({ + commandId: CommandId.make("command:delegated-recovery-wake:seed"), + events: [ + { + id: EventId.make("event:delegated-recovery-wake:message"), + type: "message.updated", + threadId, + runId: wakeRunId, + occurredAt: now, + payload: { + id: wakeMessageId, + threadId, + runId: wakeRunId, + nodeId: null, + role: "user", + text: "Background task finished", + attachments: [], + streaming: false, + createdBy: "agent", + creationSource: "server", + createdAt: now, + updatedAt: now, + delegatedCompletion: { parentRunId: runId, generation: 1, taskIds: [taskId] }, + }, + }, + { + id: EventId.make("event:delegated-recovery-wake:wake-run"), + type: "run.updated", + threadId, + runId: wakeRunId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: runPayload(wakeRunId, threadId, wakeMessageId, "running"), + }, + { + id: EventId.make("event:delegated-recovery-wake:sibling-thread"), + type: "thread.created", + threadId: siblingThreadId, + driver, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: makeSubagentChildThread({ + parentThread: parent.thread, + childThreadId: siblingThreadId, + parentNodeId: siblingId, + activeProviderThreadId: null, + providerInstanceId: modelSelection.instanceId, + modelSelection, + title: "Sibling", + now, + createdBy: "agent", + creationSource: "server", + }), + }, + { + id: EventId.make("event:delegated-recovery-wake:sibling-task"), + type: "subagent.updated", + threadId, + runId, + nodeId: siblingId, + driver, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { + ...parent.subagents[0]!, + id: siblingId, + childThreadId: siblingThreadId, + status: "running", + result: null, + completedAt: null, + completionDelivery: undefined, + }, + }, + ], + }); + // Startup runtime recovery cancels both the wake and the sibling. + yield* sink.write({ + commandId: CommandId.make("command:runtime-reconcile:delegated-recovery-wake"), + events: [ + { + id: EventId.make("event:delegated-recovery-wake:wake-cancelled"), + type: "run.updated", + threadId, + runId: wakeRunId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: runPayload(wakeRunId, threadId, wakeMessageId, "cancelled"), + }, + { + id: EventId.make("event:delegated-recovery-wake:sibling-cancelled"), + type: "run.updated", + threadId: siblingThreadId, + runId: siblingRunId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: runPayload( + siblingRunId, + siblingThreadId, + MessageId.make("message:delegated-recovery-wake-sibling"), + "cancelled", + ), + }, + ], + }); + const cohortOf = (projection: OrchestrationV2ThreadProjection) => + projection.runs.find((run) => run.id === runId)?.delegatedCompletion; + // Until the startup pass runs, nothing re-offers the cancelled wake. + assert.equal( + cohortOf(yield* orchestrator.getThreadProjection(threadId))?.delivery?.generation, + 1, + ); + + yield* orchestrator.recoverDelegatedTaskReports(() => Effect.succeed(false)); + const recovered = yield* orchestrator.getThreadProjection(threadId); + assert.isTrue( + recovered.contextTransfers.some( + (row) => row.type === "subagent_result" && row.sourceThreadId === siblingThreadId, + ), + ); + const delivery = cohortOf(recovered)?.delivery; + assert.equal(delivery?.generation, 2); + assert.deepEqual([...(delivery?.taskIds ?? [])].toSorted(), [siblingId, taskId].toSorted()); + }), + ); + it.effect("builds completion text and metadata from the same live cohort", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; diff --git a/apps/server/src/orchestration-v2/EffectWorker.ts b/apps/server/src/orchestration-v2/EffectWorker.ts index 48dcfa78e96b..f519e35e519e 100644 --- a/apps/server/src/orchestration-v2/EffectWorker.ts +++ b/apps/server/src/orchestration-v2/EffectWorker.ts @@ -109,10 +109,18 @@ export const executorLayer: Layer.Layer< const willRetry = options?.willRetry ?? false; switch (effect.request.type) { case "provider-runtime.continue": + // A delegated child that will not be continued stays cancelled, so + // its parent must get that result now rather than at next startup. return continueRestartedRun({ threadId: effect.threadId, sourceRunId: effect.request.sourceRunId, }).pipe( + Effect.tapCause(() => + willRetry ? Effect.void : threads.reconcileAppOwnedSubagentResult(effect.threadId), + ), + Effect.tap((started) => + started ? Effect.void : threads.reconcileAppOwnedSubagentResult(effect.threadId), + ), Effect.provideService(ThreadManagementService.ThreadManagementService, threads), Effect.provideService(ServerSettings.ServerSettingsService, settings), Effect.mapError( @@ -123,6 +131,7 @@ export const executorLayer: Layer.Layer< cause, }), ), + Effect.asVoid, ); case "provider-session.detach": return providerSessions diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index df9805d504c9..2f2ea1f646d9 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -74,7 +74,11 @@ import { notificationTurnItem } from "./Notification.ts"; import { isRestartNoteSource } from "./RestartBackgroundNote.ts"; import { isUndeliveredMailboxSteer } from "./NotificationMailbox.ts"; import { EventSinkV2 } from "./EventSink.ts"; -import type { OrchestrationEffectRequestV2, PendingOrchestrationEffectV2 } from "./EffectOutbox.ts"; +import type { + EffectOutboxError, + OrchestrationEffectRequestV2, + PendingOrchestrationEffectV2, +} from "./EffectOutbox.ts"; import { IdAllocatorV2 } from "./IdAllocator.ts"; import { ThreadCommandExecutor, @@ -239,6 +243,17 @@ export interface OrchestratorV2DispatchResult { export interface OrchestratorV2Shape { readonly resumeQueuedRuns: Effect.Effect; + /** + * Recovers delegated-task reports after startup runtime recovery: publishes + * results of app-owned children whose last run is terminal but unreported, + * then settles parent completion wakes that recovery terminalized. + * `deferChild` skips a child whose work may still resume. + */ + readonly recoverDelegatedTaskReports: ( + deferChild?: (run: OrchestrationV2Run) => Effect.Effect, + ) => Effect.Effect; + /** Publishes one app-owned child's terminal result to its parent, if still owed. */ + readonly reconcileAppOwnedSubagentResult: (childThreadId: ThreadId) => Effect.Effect; readonly dispatch: ( command: OrchestrationV2ServerCommand, ) => Effect.Effect; @@ -9518,97 +9533,130 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio Effect.forkDetach, ); - // Recover child results from projections. Queue recovery instead holds - // unstarted runs until an explicit queue.resume command arrives. - yield* projectionStore.getRecoveryThreadIds("subagent-results").pipe( - Effect.flatMap((threadIds) => - Effect.forEach( - threadIds, - (threadId) => - Effect.gen(function* () { - const thread = yield* projectionStore.getThreadShell(threadId); - const parentThreadId = thread?.lineage.parentThreadId; - if (parentThreadId === undefined || parentThreadId === null) return; - yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(threadId)); - }).pipe( - Effect.catchCause((cause) => - Effect.logWarning("Failed to recover terminal app-owned subagent", { - childThreadId: threadId, - cause, - }), - ), - ), - { concurrency: 8, discard: true }, + const reconcileAppOwnedSubagentResult = (childThreadId: ThreadId) => + Effect.gen(function* () { + const parentThreadId = yield* appOwnedSubagentParentThreadId(childThreadId); + if (parentThreadId === undefined) return; + yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(childThreadId)); + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("Failed to recover terminal app-owned subagent", { + childThreadId, + cause, + }), ), - ), - Effect.catchCause((cause) => - Effect.logWarning("Failed to inspect app-owned subagents during recovery", { - cause, - }), - ), - ); - yield* projectionStore.getRecoveryThreadIds("delegated-completions").pipe( - Effect.flatMap((threadIds) => - Effect.forEach( - threadIds, - (threadId) => - threadDispatch - .withLock( - threadId, - Effect.gen(function* () { - const projection = yield* projectionStore.getThreadRecords(threadId, [ - "runs", - "messages", - ]); - const terminalDeliveryRunIds = projection.runs - .filter((run) => delegatedTaskTerminalStatus(run.status) !== null) - .filter((run) => - projection.messages.some( - (message) => - message.id === run.userMessageId && - message.delegatedCompletion !== undefined, - ), - ) - .map((run) => run.id); - for (const runId of terminalDeliveryRunIds) { - yield* finalizeDelegatedCompletionDelivery(threadId, runId); - } - const refreshed = - terminalDeliveryRunIds.length === 0 - ? projection - : yield* projectionStore.getThreadRecords(threadId, ["runs", "messages"], { - messageRoles: ["user"], - }); - for (const run of refreshed.runs) { - if ( - run.delegatedCompletion?.delivery !== null && - run.delegatedCompletion !== undefined - ) { - yield* offerDelegatedCompletionDelivery(threadId, run.id); - } - } - }), - ) - .pipe( + ); + + const recoverAppOwnedSubagentResults = ( + deferChild?: (run: OrchestrationV2Run) => Effect.Effect, + ) => + projectionStore.getRecoveryThreadIds("subagent-results").pipe( + Effect.flatMap((threadIds) => + Effect.forEach( + threadIds, + (threadId) => + Effect.gen(function* () { + if (deferChild !== undefined) { + const { runs } = yield* projectionStore.getThreadRecords(threadId, ["runs"]); + const latest = runs.at(-1); + if (latest !== undefined && (yield* deferChild(latest))) return; + } + yield* reconcileAppOwnedSubagentResult(threadId); + }).pipe( Effect.catchCause((cause) => - Effect.logWarning("Failed to recover delegated completion delivery", { - threadId, + Effect.logWarning("Failed to recover terminal app-owned subagent", { + childThreadId: threadId, cause, }), ), ), - { concurrency: 8, discard: true }, + { concurrency: 8, discard: true }, + ), ), - ), - Effect.catchCause((cause) => - Effect.logWarning("Failed to inspect delegated completion delivery during recovery", { - cause, - }), - ), - ); + Effect.catchCause((cause) => + Effect.logWarning("Failed to inspect app-owned subagents during recovery", { + cause, + }), + ), + ); + // Runs once startup runtime recovery has terminalized interrupted work, so + // nothing is reported before a pending restart continuation is known. + // Queue recovery instead holds unstarted runs until queue.resume arrives. + const recoverDelegatedTaskReports = ( + deferChild?: (run: OrchestrationV2Run) => Effect.Effect, + ) => + recoverAppOwnedSubagentResults(deferChild).pipe( + Effect.andThen( + projectionStore.getRecoveryThreadIds("delegated-completions").pipe( + Effect.flatMap((threadIds) => + Effect.forEach( + threadIds, + (threadId) => + threadDispatch + .withLock( + threadId, + Effect.gen(function* () { + const projection = yield* projectionStore.getThreadRecords(threadId, [ + "runs", + "messages", + ]); + const terminalDeliveryRunIds = projection.runs + .filter((run) => delegatedTaskTerminalStatus(run.status) !== null) + .filter((run) => + projection.messages.some( + (message) => + message.id === run.userMessageId && + message.delegatedCompletion !== undefined, + ), + ) + .map((run) => run.id); + for (const runId of terminalDeliveryRunIds) { + yield* finalizeDelegatedCompletionDelivery(threadId, runId); + } + const refreshed = + terminalDeliveryRunIds.length === 0 + ? projection + : yield* projectionStore.getThreadRecords( + threadId, + ["runs", "messages"], + { + messageRoles: ["user"], + }, + ); + for (const run of refreshed.runs) { + if ( + run.delegatedCompletion?.delivery !== null && + run.delegatedCompletion !== undefined + ) { + yield* offerDelegatedCompletionDelivery(threadId, run.id); + } + } + }), + ) + .pipe( + Effect.catchCause((cause) => + Effect.logWarning("Failed to recover delegated completion delivery", { + threadId, + cause, + }), + ), + ), + { concurrency: 8, discard: true }, + ), + ), + Effect.catchCause((cause) => + Effect.logWarning("Failed to inspect delegated completion delivery during recovery", { + cause, + }), + ), + ), + ), + ); return OrchestratorV2.of({ resumeQueuedRuns, + recoverDelegatedTaskReports, + reconcileAppOwnedSubagentResult, dispatch: dispatchWithReceipt, getTimelinePage: (threadId, options) => projectionStore @@ -9718,6 +9766,8 @@ export const layer: Layer.Layer< const layerUnavailable: Layer.Layer = Layer.succeed( OrchestratorV2, OrchestratorV2.of({ + recoverDelegatedTaskReports: () => Effect.void, + reconcileAppOwnedSubagentResult: () => Effect.void, resumeQueuedRuns: Effect.fail( new OrchestratorDispatchError({ commandId: CommandId.make("command:system:resume-queued-runs"), diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts index 8df85818a761..1a990c320170 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts @@ -1272,3 +1272,28 @@ it.effect( }).pipe(Effect.provide(layer)); }, ); + +it.effect("propagates an unreadable restart continuation instead of reporting it absent", () => { + const failure = new EffectOutbox.EffectOutboxError({ + operation: "get", + cause: new Error("offline"), + }); + const layer = ProviderRuntimeRecovery.layer.pipe( + Layer.provide(ServerSettings.layerTest()), + Layer.provide( + Layer.mergeAll( + Layer.mock(ProjectionStore.ProjectionStoreV2)({}), + Layer.mock(EventSink.EventSinkV2)({}), + IdAllocator.layer, + Layer.mock(EffectOutbox.EffectOutboxV2)({ get: () => Effect.fail(failure) }), + ), + ), + ); + return Effect.gen(function* () { + const recovery = yield* ProviderRuntimeRecovery.ProviderRuntimeRecoveryService; + const error = yield* Effect.flip( + recovery.isRestartContinuationPending(RunId.make("run:unreadable")), + ); + assert.strictEqual(error, failure); + }).pipe(Effect.provide(layer)); +}); diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 67ae22f1883b..8414381b6504 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -2,6 +2,7 @@ import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { CommandId, type OrchestrationV2DomainEvent, + type RunId, type ProviderThreadId, type OrchestrationV2RestartCancelledBackgroundWork, type OrchestrationV2ThreadProjection, @@ -11,6 +12,7 @@ import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; import * as EffectOutbox from "./EffectOutbox.ts"; @@ -62,6 +64,10 @@ export class ProviderRuntimeRecoveryService extends Context.Service< ) => Effect.Effect; readonly prepareForShutdown: Effect.Effect; readonly recover: Effect.Effect; + /** Whether a restart continuation for this run is still waiting to settle. */ + readonly isRestartContinuationPending: ( + runId: RunId, + ) => Effect.Effect; } >()("t3/orchestration-v2/ProviderRuntimeRecoveryService") {} @@ -798,7 +804,23 @@ export const make = Effect.gen(function* () { return (yield* reconcile("startup")) satisfies ProviderRuntimeRecoverySummary; }); - return ProviderRuntimeRecoveryService.of({ reconcile, prepareForShutdown, recover }); + const isRestartContinuationPending = (runId: RunId) => + outbox + .get(`effect:restart-continuation:${runId}`) + .pipe( + Effect.map( + (effect) => + Option.isSome(effect) && + (effect.value.status === "pending" || effect.value.status === "running"), + ), + ); + + return ProviderRuntimeRecoveryService.of({ + reconcile, + prepareForShutdown, + recover, + isRestartContinuationPending, + }); }); export const layer = Layer.effect(ProviderRuntimeRecoveryService, make); diff --git a/apps/server/src/orchestration-v2/RestartContinuation.test.ts b/apps/server/src/orchestration-v2/RestartContinuation.test.ts index f43f4cd01666..64f9cc8f7351 100644 --- a/apps/server/src/orchestration-v2/RestartContinuation.test.ts +++ b/apps/server/src/orchestration-v2/RestartContinuation.test.ts @@ -22,6 +22,13 @@ import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "./IdAllocator.ts"; import * as EffectWorker from "./EffectWorker.ts"; import * as EffectOutbox from "./EffectOutbox.ts"; +import * as CheckpointRollbackService from "./CheckpointRollbackService.ts"; +import * as ProviderSessionManager from "./ProviderSessionManager.ts"; +import * as ProviderTurnControlService from "./ProviderTurnControlService.ts"; +import * as ProviderTurnStartService from "./ProviderTurnStartService.ts"; +import * as RunFinalizationService from "./RunFinalizationService.ts"; +import * as RuntimeRequestService from "./RuntimeRequestService.ts"; +import * as ThreadTitleRegenerationService from "./ThreadTitleRegenerationService.ts"; const threadId = ThreadId.make("thread:restart"); const runId = RunId.make("run:restart"); @@ -513,3 +520,63 @@ it.effect("does not cancel or resume a run that completes while shutdown intent assert.isFalse(dispatched); }), ); + +it.effect("reconciles a delegated child when its restart continuation will not run", () => + Effect.gen(function* () { + const projection = makeProjection(); + const cancelled = { ...projection, runs: [{ ...projection.runs[0]!, status: "cancelled" }] }; + const reconciled: ThreadId[] = []; + let dispatchFails = false; + const execute = (settingEnabled: boolean, willRetry = false) => + Effect.gen(function* () { + const executor = yield* EffectWorker.OrchestrationEffectExecutorV2; + return yield* executor + .execute( + { + id: `effect:restart-continuation:${runId}`, + threadId, + request: { type: "provider-runtime.continue", sourceRunId: runId }, + } as EffectOutbox.OrchestrationEffectV2, + { willRetry }, + ) + .pipe(Effect.exit); + }).pipe( + Effect.provide( + EffectWorker.executorLayer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadRecords: () => Effect.succeed(cancelled as never), + dispatch: () => + dispatchFails ? Effect.die("dispatch failed") : Effect.succeed({} as never), + reconcileAppOwnedSubagentResult: (childThreadId) => + Effect.sync(() => void reconciled.push(childThreadId)), + }), + ServerSettings.layerTest({ continueThreadsAfterServerUpdate: settingEnabled }), + Layer.mock(RunFinalizationService.RunFinalizationService)({}), + Layer.mock(CheckpointRollbackService.CheckpointRollbackServiceV2)({}), + Layer.mock(ProviderSessionManager.ProviderSessionManagerV2)({}), + Layer.mock(ProviderTurnControlService.ProviderTurnControlServiceV2)({}), + Layer.mock(ProviderTurnStartService.ProviderTurnStartServiceV2)({}), + Layer.mock(RuntimeRequestService.RuntimeRequestServiceV2)({}), + Layer.mock(ThreadTitleRegenerationService.ThreadTitleRegenerationService)({}), + ), + ), + ), + ), + ); + + // A continuation that starts leaves the result to the continued run. + yield* execute(true); + assert.deepEqual(reconciled, []); + // A skipped continuation publishes the cancelled result now. + yield* execute(false); + assert.deepEqual(reconciled, [threadId]); + // A failure that will be retried waits; the final failure publishes. + dispatchFails = true; + yield* execute(true, true); + assert.deepEqual(reconciled, [threadId]); + yield* execute(true, false); + assert.deepEqual(reconciled, [threadId, threadId]); + }), +); diff --git a/apps/server/src/orchestration-v2/RestartContinuation.ts b/apps/server/src/orchestration-v2/RestartContinuation.ts index 216b830de41e..76ce9149be4c 100644 --- a/apps/server/src/orchestration-v2/RestartContinuation.ts +++ b/apps/server/src/orchestration-v2/RestartContinuation.ts @@ -89,11 +89,12 @@ export function restartContinuationRun( return run; } +/** Returns whether a continuation run was dispatched; false means it was skipped. */ export const continueRestartedRun = Effect.fn("RestartContinuation.continueRestartedRun")( function* (input: { readonly threadId: ThreadId; readonly sourceRunId: RunId }) { const settings = yield* ServerSettings.ServerSettingsService; const enabled = yield* settings.getSettings.pipe(Effect.orElseSucceed(() => null)); - if (!enabled) return; + if (!enabled) return false; const threads = yield* ThreadManagementService.ThreadManagementService; const messageId = MessageId.make(`message:restart-continuation:${input.sourceRunId}`); const projection = yield* threads.getThreadRecords( @@ -105,18 +106,19 @@ export const continueRestartedRun = Effect.fn("RestartContinuation.continueResta !resolveProjectSettings(enabled, projection.thread.projectId).settings .continueThreadsAfterServerUpdate ) - return; - if (projection.thread.archivedAt !== null || projection.thread.deletedAt !== null) return; + return false; + if (projection.thread.archivedAt !== null || projection.thread.deletedAt !== null) return false; - if (projection.messages.some((message) => message.id === messageId)) return; + // Already dispatched by an earlier attempt of this effect. + if (projection.messages.some((message) => message.id === messageId)) return true; const source = projection.runs.find((run) => run.id === input.sourceRunId); // A settled source prompts with the note of the background work it lost. const noteSource = source !== undefined && isRestartNoteSource(source, projection.providerTurns); - if (!source || (source.status !== "cancelled" && !noteSource)) return; + if (!source || (source.status !== "cancelled" && !noteSource)) return false; // A user submission after reconciliation takes precedence over an automatic prompt. - if (projection.runs.some((run) => run.ordinal > source.ordinal)) return; - if (projection.thread.providerInstanceId !== source.providerInstanceId) return; + if (projection.runs.some((run) => run.ordinal > source.ordinal)) return false; + if (projection.thread.providerInstanceId !== source.providerInstanceId) return false; yield* threads.dispatch({ type: "message.dispatch", commandId: CommandId.make(`command:restart-continuation:${input.sourceRunId}`), @@ -132,5 +134,6 @@ export const continueRestartedRun = Effect.fn("RestartContinuation.continueResta creationSource: "server", restartContinuationOfRunId: input.sourceRunId, }); + return true; }, ); diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.ts b/apps/server/src/orchestration-v2/ThreadManagementService.ts index b489f6890967..87851e1babb5 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.ts @@ -277,6 +277,8 @@ export interface ThreadManagementServiceShape { readonly getTimelinePage: Orchestrator.OrchestratorV2["Service"]["getTimelinePage"]; readonly getMessageCount: Orchestrator.OrchestratorV2["Service"]["getMessageCount"]; readonly getThreadRecords: Orchestrator.OrchestratorV2["Service"]["getThreadRecords"]; + readonly recoverDelegatedTaskReports: Orchestrator.OrchestratorV2["Service"]["recoverDelegatedTaskReports"]; + readonly reconcileAppOwnedSubagentResult: Orchestrator.OrchestratorV2["Service"]["reconcileAppOwnedSubagentResult"]; readonly getThreadProjection: ( threadId: ThreadId, ) => Effect.Effect; @@ -723,6 +725,8 @@ const make = Effect.gen(function* () { ensureProjectionTranscript(threadId).pipe( Effect.andThen(orchestrator.getThreadRecords(threadId, fields, filter)), ), + recoverDelegatedTaskReports: orchestrator.recoverDelegatedTaskReports, + reconcileAppOwnedSubagentResult: orchestrator.reconcileAppOwnedSubagentResult, getThreadProjection, getCheckpointContext, getThreadSnapshot, diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index 30ada288bd76..0e07f0feb57a 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -194,6 +194,8 @@ const makeTestRelay = Effect.fnUntraced(function* ( dispatch: unused, getTimelinePage: () => Effect.die("Unused timeline read"), getMessageCount: () => Effect.die("unused message count"), + recoverDelegatedTaskReports: unused, + reconcileAppOwnedSubagentResult: unused, getThreadRecords: () => Effect.die("unused record read"), getThreadProjection: unused, getCheckpointContext: unused, diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index e3f92058e70c..750acb4d7109 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -508,7 +508,23 @@ const make = (options?: StartupOptions) => ), ), ), - recover: runStartupPhase("orchestration-v2.recovery", providerRuntimeRecovery.recover), + recover: runStartupPhase( + "orchestration-v2.recovery", + providerRuntimeRecovery.recover.pipe( + // Report delegated tasks only once recovery has cancelled + // interrupted work, leaving a child with a pending restart + // continuation to that continuation. + Effect.tap(() => + ThreadManagement.ThreadManagementService.pipe( + Effect.flatMap((threads) => + threads.recoverDelegatedTaskReports((run) => + providerRuntimeRecovery.isRestartContinuationPending(run.id), + ), + ), + ), + ), + ), + ), startEffectWorker: runStartupPhase( "orchestration-v2.effect-worker.start", startEffectWorkerWithRelay({