From 945a40e6fb1ab5e04f24dc4dc9e059c1cd4568b5 Mon Sep 17 00:00:00 2001 From: Yash Singh Date: Thu, 24 Sep 2026 20:59:11 -0500 Subject: [PATCH 1/2] fix(server): allow forks from provider-finished runs --- .../src/orchestration-v2/Orchestrator.ts | 10 +- .../ThreadForkService.test.ts | 108 +++++++++++++----- .../src/orchestration-v2/ThreadForkService.ts | 27 ++++- 3 files changed, 111 insertions(+), 34 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 05565a40e2eb..59b24f0a26c3 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -94,7 +94,11 @@ import { delegatedTaskProgress, subagentThreadTitle, } from "./SubagentProjection.ts"; -import { ThreadForkServiceV2 } from "./ThreadForkService.ts"; +import { + forkableSourceRunStatusError, + isForkableSourceRunStatus, + ThreadForkServiceV2, +} from "./ThreadForkService.ts"; import { planThreadDeletion } from "./ThreadDeletion.ts"; export class OrchestratorDispatchError extends Schema.TaggedError()( @@ -3110,11 +3114,11 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio cause: `No stable source run was found for fork source ${command.sourcePoint.type}.`, }); } - if (sourceRun.status !== "completed") { + if (!isForkableSourceRunStatus(sourceRun.status)) { return yield* new OrchestratorDispatchError({ commandId: command.commandId, commandType: command.type, - cause: `Fork source run ${sourceRun.id} is ${sourceRun.status}; only completed runs are supported.`, + cause: forkableSourceRunStatusError(sourceRun), }); } const sourceProviderThread = providerThreadForRun(sourceProjection, sourceRun); diff --git a/apps/server/src/orchestration-v2/ThreadForkService.test.ts b/apps/server/src/orchestration-v2/ThreadForkService.test.ts index 9315783c8541..f734dfdb68db 100644 --- a/apps/server/src/orchestration-v2/ThreadForkService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadForkService.test.ts @@ -15,7 +15,12 @@ import { import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; -import { layer, ThreadForkServiceV2 } from "./ThreadForkService.ts"; +import { + forkableSourceRunStatusError, + isForkableSourceRunStatus, + layer, + ThreadForkServiceV2, +} from "./ThreadForkService.ts"; const sourceThreadId = ThreadId.make("thread:fork-snoozed-source"); const targetThreadId = ThreadId.make("thread:fork-awake-target"); @@ -62,7 +67,7 @@ function makeSourceThread(): OrchestrationV2AppThread { }; } -function makeCompletedSourceRun(): OrchestrationV2Run { +function makeSourceRun(status: OrchestrationV2Run["status"]): OrchestrationV2Run { return { id: sourceRunId, threadId: sourceThreadId, @@ -73,7 +78,7 @@ function makeCompletedSourceRun(): OrchestrationV2Run { userMessageId: MessageId.make("message:fork-snoozed-source"), rootNodeId: null, activeAttemptId: null, - status: "completed", + status, queuePosition: null, requestedAt: sourceCreatedAt, startedAt: sourceCreatedAt, @@ -83,33 +88,34 @@ function makeCompletedSourceRun(): OrchestrationV2Run { }; } -it.effect("keeps a fork awake when its source thread is snoozed", () => +function makeSourceProjection(sourceRun: OrchestrationV2Run): OrchestrationV2ThreadProjection { + return { + thread: makeSourceThread(), + runs: [sourceRun], + attempts: [], + nodes: [], + subagents: [], + providerSessions: [], + providerThreads: [], + providerTurns: [], + runtimeRequests: [], + messages: [], + plans: [], + turnItems: [], + checkpointScopes: [], + checkpoints: [], + contextHandoffs: [], + contextTransfers: [], + visibleTurnItems: [], + updatedAt: snoozedAt, + }; +} + +const planFork = (sourceRun: OrchestrationV2Run) => Effect.gen(function* () { - const sourceThread = makeSourceThread(); - const sourceRun = makeCompletedSourceRun(); - const sourceProjection: OrchestrationV2ThreadProjection = { - thread: sourceThread, - runs: [sourceRun], - attempts: [], - nodes: [], - subagents: [], - providerSessions: [], - providerThreads: [], - providerTurns: [], - runtimeRequests: [], - messages: [], - plans: [], - turnItems: [], - checkpointScopes: [], - checkpoints: [], - contextHandoffs: [], - contextTransfers: [], - visibleTurnItems: [], - updatedAt: snoozedAt, - }; const service = yield* ThreadForkServiceV2; - const result = yield* service.plan({ - sourceProjection, + return yield* service.plan({ + sourceProjection: makeSourceProjection(sourceRun), sourceRun, sourceProviderThread: undefined, canonicalSourcePoint: { @@ -123,6 +129,26 @@ it.effect("keeps a fork awake when its source thread is snoozed", () => creationSource: "mobile", createdAt: forkCreatedAt, }); + }).pipe(Effect.provide(layer)); + +it("treats usage-limited and other provider-finished runs as forkable", () => { + assert.isTrue(isForkableSourceRunStatus("completed")); + assert.isTrue(isForkableSourceRunStatus("waiting")); + assert.isTrue(isForkableSourceRunStatus("failed")); + assert.isTrue(isForkableSourceRunStatus("interrupted")); + assert.isTrue(isForkableSourceRunStatus("cancelled")); + assert.isFalse(isForkableSourceRunStatus("running")); + assert.isFalse(isForkableSourceRunStatus("starting")); + assert.isFalse(isForkableSourceRunStatus("queued")); + assert.isFalse(isForkableSourceRunStatus("preparing")); + assert.isFalse(isForkableSourceRunStatus("rolled_back")); +}); + +it.effect("keeps a fork awake when its source thread is snoozed", () => + Effect.gen(function* () { + const sourceThread = makeSourceThread(); + const sourceRun = makeSourceRun("completed"); + const result = yield* planFork(sourceRun); assert.isNull(result.targetThread.snoozedUntil); assert.isNull(result.targetThread.snoozedAt); @@ -144,5 +170,29 @@ it.effect("keeps a fork awake when its source thread is snoozed", () => threadId: sourceThreadId, runId: sourceRunId, }); - }).pipe(Effect.provide(layer)), + }), +); + +it.effect("forks from a usage-limited failed run", () => + Effect.gen(function* () { + const result = yield* planFork(makeSourceRun("failed")); + assert.deepEqual(result.targetThread.forkedFrom, { + type: "run", + threadId: sourceThreadId, + runId: sourceRunId, + }); + }), +); + +it.effect("rejects in-progress and rolled-back fork sources", () => + Effect.gen(function* () { + for (const status of ["running", "rolled_back"] as const) { + const sourceRun = makeSourceRun(status); + const error = yield* planFork(sourceRun).pipe(Effect.flip); + assert.equal(error._tag, "ThreadForkPlanError"); + assert.equal(error.sourceThreadId, sourceThreadId); + assert.equal(error.targetThreadId, targetThreadId); + assert.equal(error.cause, forkableSourceRunStatusError(sourceRun)); + } + }), ); diff --git a/apps/server/src/orchestration-v2/ThreadForkService.ts b/apps/server/src/orchestration-v2/ThreadForkService.ts index 59d94d236b65..56c7aa81693d 100644 --- a/apps/server/src/orchestration-v2/ThreadForkService.ts +++ b/apps/server/src/orchestration-v2/ThreadForkService.ts @@ -30,6 +30,29 @@ export class ThreadForkPlanError extends Schema.TaggedError }, ) {} +/** + * Fork copies a provider-finished conversation. Usage-limited and other + * failed, interrupted, or cancelled turns still have a native thread (or a + * portable transcript) even though the run did not complete successfully. + * `waiting` is provider-finished with checkpoint capture still pending. + * In-progress and rolled-back runs are not forkable. + */ +export function isForkableSourceRunStatus(status: OrchestrationV2Run["status"]): boolean { + return ( + status === "completed" || + status === "waiting" || + status === "failed" || + status === "interrupted" || + status === "cancelled" + ); +} + +export function forkableSourceRunStatusError( + run: Pick, +): string { + return `Fork source run ${run.id} is ${run.status}; in-progress and rolled-back runs cannot be forked.`; +} + export interface ThreadForkServiceV2Shape { readonly plan: (input: { readonly sourceProjection: Pick; @@ -55,11 +78,11 @@ export const layer: Layer.Layer = Layer.succeed( ThreadForkServiceV2.of({ plan: (input) => Effect.gen(function* () { - if (input.sourceRun.status !== "completed") { + if (!isForkableSourceRunStatus(input.sourceRun.status)) { return yield* new ThreadForkPlanError({ sourceThreadId: input.sourceProjection.thread.id, targetThreadId: input.targetThreadId, - cause: `Fork source run ${input.sourceRun.id} is ${input.sourceRun.status}.`, + cause: forkableSourceRunStatusError(input.sourceRun), }); } const targetThread: OrchestrationV2AppThread = { From 778fae5ff516febaca721c418410ff6d141404b3 Mon Sep 17 00:00:00 2001 From: Yash Singh Date: Sat, 26 Sep 2026 03:12:40 -0500 Subject: [PATCH 2/2] fix(server): forks of unfinished runs no longer pull in later turns - Only fork natively when the source run completed or is waiting; failed, interrupted, and cancelled runs fall back to a portable transcript that stops at the selected run - Pass the source run status into decideForkExecution - Add execution tests for Codex and Claude forks of failed, interrupted, and cancelled runs --- .../orchestration-v2/CommandPolicy.test.ts | 4 + .../src/orchestration-v2/CommandPolicy.ts | 5 + .../src/orchestration-v2/Orchestrator.ts | 3 +- .../ThreadFork.execution.test.ts | 243 ++++++++++++++++++ 4 files changed, 254 insertions(+), 1 deletion(-) create mode 100644 apps/server/src/orchestration-v2/ThreadFork.execution.test.ts diff --git a/apps/server/src/orchestration-v2/CommandPolicy.test.ts b/apps/server/src/orchestration-v2/CommandPolicy.test.ts index 2f9f88db12f4..d663184dfd46 100644 --- a/apps/server/src/orchestration-v2/CommandPolicy.test.ts +++ b/apps/server/src/orchestration-v2/CommandPolicy.test.ts @@ -473,6 +473,7 @@ layer("CommandPolicyV2", (it) => { capabilities: CodexProviderCapabilitiesV2, sameProvider: true, hasStrongNativeSource: true, + sourceRunStatus: "completed", fromSpecificTurn: true, }); @@ -491,6 +492,7 @@ layer("CommandPolicyV2", (it) => { capabilities: CursorProviderCapabilitiesV2, sameProvider: true, hasStrongNativeSource: true, + sourceRunStatus: "completed", fromSpecificTurn: true, }); @@ -509,6 +511,7 @@ layer("CommandPolicyV2", (it) => { capabilities: GrokProviderCapabilitiesV2, sameProvider: true, hasStrongNativeSource: true, + sourceRunStatus: "completed", fromSpecificTurn: true, }); @@ -538,6 +541,7 @@ layer("CommandPolicyV2", (it) => { })), sameProvider: true, hasStrongNativeSource: true, + sourceRunStatus: "completed", fromSpecificTurn: true, }) .pipe(Effect.flip); diff --git a/apps/server/src/orchestration-v2/CommandPolicy.ts b/apps/server/src/orchestration-v2/CommandPolicy.ts index cc7b88e12767..5f7d130c17fd 100644 --- a/apps/server/src/orchestration-v2/CommandPolicy.ts +++ b/apps/server/src/orchestration-v2/CommandPolicy.ts @@ -3,6 +3,7 @@ import { ModelSelection, type OrchestrationV2Command, OrchestrationV2ProviderCapabilities, + type OrchestrationV2Run, OrchestrationV2ThreadProjection, ProviderInstanceId, ProviderTurnId, @@ -202,6 +203,7 @@ export interface CommandPolicyV2Shape { input: CapabilityCheckInput & { readonly sameProvider: boolean; readonly hasStrongNativeSource: boolean; + readonly sourceRunStatus: OrchestrationV2Run["status"]; readonly fromSpecificTurn: boolean; }, ) => Effect.Effect; @@ -361,7 +363,10 @@ const ensureContextHandoff: CommandPolicyV2Shape["ensureContextHandoff"] = (inpu }; const decideForkExecution: CommandPolicyV2Shape["decideForkExecution"] = (input) => { + // Unsuccessful runs may have no native turn or assistant cursor. Forking + // those at native head can include later turns, so use the bounded transcript. const canForkNatively = + (input.sourceRunStatus === "completed" || input.sourceRunStatus === "waiting") && input.sameProvider && input.hasStrongNativeSource && input.capabilities.threads.canForkThread && diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 59b24f0a26c3..caf458716ce6 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -5138,7 +5138,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ); const forkExecution = - pendingForkTransfer === undefined + pendingForkTransfer === undefined || sourceRun === null ? null : yield* enforceCommandPolicy(command)( commandPolicy.decideForkExecution({ @@ -5149,6 +5149,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio sameProvider: pendingForkTransfer.sourceProviderInstanceId === modelSelection.instanceId, hasStrongNativeSource: sourceProviderThread?.nativeThreadRef?.strength === "strong", + sourceRunStatus: sourceRun.status, fromSpecificTurn: sourceRun !== null, }), ); diff --git a/apps/server/src/orchestration-v2/ThreadFork.execution.test.ts b/apps/server/src/orchestration-v2/ThreadFork.execution.test.ts new file mode 100644 index 000000000000..2f86f5946329 --- /dev/null +++ b/apps/server/src/orchestration-v2/ThreadFork.execution.test.ts @@ -0,0 +1,243 @@ +import { assert, it } from "@effect/vitest"; +import { + CommandId, + EventId, + MessageId, + NodeId, + ProjectId, + ProviderDriverKind, + ProviderInstanceId, + ProviderThreadId, + ProviderTurnId, + RunAttemptId, + RunId, + ThreadId, + TurnItemId, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; + +import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; +import { ClaudeProviderCapabilitiesV2 } from "./Adapters/ClaudeAdapterV2.ts"; +import { EventSinkV2 } from "./EventSink.ts"; +import { OrchestratorV2 } from "./Orchestrator.ts"; +import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; +import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; +import { makeOrchestratorV2ReplayLayerWithRegistry } from "./testkit/ProviderReplayHarness.ts"; + +for (const driverName of ["codex", "claudeAgent"] as const) { + const driver = ProviderDriverKind.make(driverName); + const instanceId = ProviderInstanceId.make(driver); + const modelSelection = { instanceId, model: "test-model" }; + const adapter: ProviderAdapterV2Shape = { + instanceId, + driver, + getCapabilities: () => + Effect.succeed( + driver === "codex" ? CodexProviderCapabilitiesV2 : ClaudeProviderCapabilitiesV2, + ), + planSelectionTransition: () => Effect.succeed({ type: "apply_on_next_turn" }), + openSession: () => Effect.die("Execution is paused after dispatch for handoff inspection"), + }; + const layer = makeOrchestratorV2ReplayLayerWithRegistry( + { name: `fork-boundary-${driver}` }, + ProviderAdapterRegistry.makeLayer([adapter]), + { runEffectWorker: false }, + ); + + for (const status of ["failed", "interrupted", "cancelled"] as const) { + it.effect(`bounds ${driver} context when continuing a fork of a ${status} run`, () => + Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const eventSink = yield* EventSinkV2; + const now = yield* DateTime.now; + const sourceThreadId = ThreadId.make("fork-boundary-source"); + const targetThreadId = ThreadId.make("fork-boundary-target"); + const providerThreadId = ProviderThreadId.make("fork-boundary-native-thread"); + const sourceRunId = RunId.make("fork-boundary-source-run"); + const attemptId = RunAttemptId.make("interrupted-source-attempt"); + const providerTurnId = ProviderTurnId.make("interrupted-source-turn"); + const rootNodeId = NodeId.make("interrupted-source-root"); + + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create-source"), + threadId: sourceThreadId, + projectId: ProjectId.make("fork-boundary-project"), + title: "Fork boundary source", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + yield* eventSink.write({ + events: [ + { + id: EventId.make("source-provider-thread"), + type: "provider-thread.updated", + threadId: sourceThreadId, + occurredAt: now, + payload: { + id: providerThreadId, + driver, + providerInstanceId: instanceId, + providerSessionId: null, + appThreadId: sourceThreadId, + ownerNodeId: null, + nativeThreadRef: { driver, nativeId: "native-source", strength: "strong" }, + nativeConversationHeadRef: null, + status: "idle", + firstRunOrdinal: 1, + lastRunOrdinal: 2, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + ], + }); + // A cancelled queue entry has no provider turn; an early interruption + // can have a turn but no native assistant cursor. + if (status === "interrupted") { + yield* eventSink.write({ + events: [ + { + id: EventId.make("source-attempt"), + type: "run-attempt.created", + threadId: sourceThreadId, + runId: sourceRunId, + occurredAt: now, + payload: { + id: attemptId, + runId: sourceRunId, + attemptOrdinal: 1, + rootNodeId, + providerInstanceId: instanceId, + providerThreadId, + providerTurnId, + reason: "initial", + status, + startedAt: now, + completedAt: now, + }, + }, + { + id: EventId.make("source-provider-turn"), + type: "provider-turn.updated", + threadId: sourceThreadId, + occurredAt: now, + payload: { + id: providerTurnId, + providerThreadId, + nodeId: rootNodeId, + runAttemptId: attemptId, + nativeTurnRef: { driver, nativeId: "turn:synthetic", strength: "weak" }, + ordinal: 1, + status, + startedAt: now, + completedAt: now, + }, + }, + ], + }); + } + for (const ordinal of [1, 2]) { + const runId = ordinal === 1 ? sourceRunId : RunId.make("later-run"); + const messageId = MessageId.make(`source-message-${ordinal}`); + yield* eventSink.write({ + events: [ + { + id: EventId.make(`run-${ordinal}`), + type: "run.created", + threadId: sourceThreadId, + runId, + occurredAt: now, + payload: { + id: runId, + threadId: sourceThreadId, + ordinal, + providerInstanceId: instanceId, + modelSelection, + providerThreadId, + userMessageId: messageId, + rootNodeId: null, + activeAttemptId: ordinal === 1 && status === "interrupted" ? attemptId : null, + status: ordinal === 1 ? status : "completed", + queuePosition: null, + requestedAt: now, + startedAt: now, + completedAt: now, + checkpointId: null, + contextHandoffId: null, + }, + }, + { + id: EventId.make(`item-${ordinal}`), + type: "turn-item.updated", + threadId: sourceThreadId, + runId, + occurredAt: now, + payload: { + id: TurnItemId.make(`item-${ordinal}`), + threadId: sourceThreadId, + runId, + nodeId: null, + providerThreadId, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal, + status: "completed", + title: null, + startedAt: now, + completedAt: now, + updatedAt: now, + type: "user_message", + createdBy: "user", + creationSource: "web", + inputIntent: "turn_start", + messageId, + text: ordinal === 1 ? "INCLUDED_SOURCE_MARKER" : "EXCLUDED_LATER_MARKER", + attachments: [], + }, + }, + ], + }); + } + yield* orchestrator.dispatch({ + type: "thread.fork", + commandId: CommandId.make("fork-source"), + sourceThreadId, + targetThreadId, + sourcePoint: { type: "run", runId: sourceRunId }, + createdBy: "user", + creationSource: "web", + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("continue-fork"), + threadId: targetThreadId, + messageId: MessageId.make("continue-fork"), + text: "Continue from the selected source run", + attachments: [], + modelSelection, + dispatchMode: { type: "start_immediately" }, + createdBy: "user", + creationSource: "web", + }); + const target = yield* orchestrator.getThreadProjection(targetThreadId); + assert.equal(target.contextTransfers[0]?.resolution?.strategy, "portable_context"); + assert.lengthOf(target.contextHandoffs, 1); + const handoff = target.contextHandoffs[0]!; + const history = handoff.history?.messages.map((message) => message.text).join("\n") ?? ""; + assert.include(`${handoff.summaryText}\n${history}`, "INCLUDED_SOURCE_MARKER"); + assert.notInclude(`${handoff.summaryText}\n${history}`, "EXCLUDED_LATER_MARKER"); + assert.isNull(target.providerThreads[0]?.forkedFrom); + }).pipe(Effect.provide(layer)), + ); + } +}