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
4 changes: 4 additions & 0 deletions apps/server/src/orchestration-v2/CommandPolicy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -473,6 +473,7 @@ layer("CommandPolicyV2", (it) => {
capabilities: CodexProviderCapabilitiesV2,
sameProvider: true,
hasStrongNativeSource: true,
sourceRunStatus: "completed",
fromSpecificTurn: true,
});

Expand All @@ -491,6 +492,7 @@ layer("CommandPolicyV2", (it) => {
capabilities: CursorProviderCapabilitiesV2,
sameProvider: true,
hasStrongNativeSource: true,
sourceRunStatus: "completed",
fromSpecificTurn: true,
});

Expand All @@ -509,6 +511,7 @@ layer("CommandPolicyV2", (it) => {
capabilities: GrokProviderCapabilitiesV2,
sameProvider: true,
hasStrongNativeSource: true,
sourceRunStatus: "completed",
fromSpecificTurn: true,
});

Expand Down Expand Up @@ -538,6 +541,7 @@ layer("CommandPolicyV2", (it) => {
})),
sameProvider: true,
hasStrongNativeSource: true,
sourceRunStatus: "completed",
fromSpecificTurn: true,
})
.pipe(Effect.flip);
Expand Down
5 changes: 5 additions & 0 deletions apps/server/src/orchestration-v2/CommandPolicy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
ModelSelection,
type OrchestrationV2Command,
OrchestrationV2ProviderCapabilities,
type OrchestrationV2Run,
OrchestrationV2ThreadProjection,
ProviderInstanceId,
ProviderTurnId,
Expand Down Expand Up @@ -202,6 +203,7 @@ export interface CommandPolicyV2Shape {
input: CapabilityCheckInput & {
readonly sameProvider: boolean;
readonly hasStrongNativeSource: boolean;
readonly sourceRunStatus: OrchestrationV2Run["status"];
readonly fromSpecificTurn: boolean;
},
) => Effect.Effect<ForkExecutionPolicyV2, CommandPolicyV2Error>;
Expand Down Expand Up @@ -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 &&
Expand Down
13 changes: 9 additions & 4 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,11 @@ import {
delegatedTaskProgress,
subagentThreadTitle,
} from "./SubagentProjection.ts";
import { ThreadForkServiceV2 } from "./ThreadForkService.ts";
import {
forkableSourceRunStatusError,
isForkableSourceRunStatus,
ThreadForkServiceV2,
} from "./ThreadForkService.ts";
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
import { planThreadDeletion } from "./ThreadDeletion.ts";

export class OrchestratorDispatchError extends Schema.TaggedError<OrchestratorDispatchError>()(
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -5134,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({
Expand All @@ -5145,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,
}),
);
Expand Down
243 changes: 243 additions & 0 deletions apps/server/src/orchestration-v2/ThreadFork.execution.test.ts
Original file line number Diff line number Diff line change
@@ -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)),
);
}
}
Loading
Loading