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: 2 additions & 2 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4382,8 +4382,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
const source = projection.runs.find((run) => run.id === command.restartContinuationOfRunId);
if (
!source ||
(source.status !== "cancelled" &&
!isRestartNoteSource(source, projection.providerTurns)) ||
source.status !== "cancelled" ||
isRestartNoteSource(source, projection.providerTurns) ||
projection.thread.archivedAt !== null ||
projection.thread.deletedAt !== null ||
projection.thread.providerInstanceId !== source.providerInstanceId ||
Expand Down
36 changes: 2 additions & 34 deletions apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -149,33 +149,6 @@ function resolveStaleBackgroundItemProviderInstanceId(
return projection.providerThreads[0]?.providerInstanceId ?? projection.thread.providerInstanceId;
}

/**
* Provider threads with background work reconciliation would cancel and
* record for their next turn.
*/
function providerThreadsWithOpenBackgroundWork(
projection: ProjectionStore.ProjectionRuntimeRecoveryState,
): ReadonlySet<ProviderThreadId> {
const ids = new Set<ProviderThreadId>();
for (const item of projection.turnItems ?? []) {
if (
!isBackgroundCapableTurnItemType(item.type) ||
!isNonterminalTurnItemStatus(item.status) ||
isAppOwnedDelegationItem(item)
)
continue;
const providerThreadId =
item.providerThreadId ??
projection.runs.find((run) => run.id === item.runId)?.providerThreadId;
if (providerThreadId != null) ids.add(providerThreadId);
}
for (const thread of projection.providerThreads ?? []) {
if (thread.ownerNodeId === null && providerThreadHasPendingBackgroundTasks(thread))
ids.add(thread.id);
}
return ids;
}

/**
* A provider thread's latest started run: the last turn that provider saw.
* Restart recovery records the thread's cancelled background work on it, and
Expand Down Expand Up @@ -668,7 +641,7 @@ export const make = Effect.gen(function* () {
}
const continuationRun =
continueAfterRestart && trigger === "startup"
? restartContinuationRun(projection, new Set(cancelledBackgroundWork.keys()))
? restartContinuationRun(projection)
: undefined;
const effects: Array<EffectOutbox.PendingOrchestrationEffectV2> = continuationRun
? [
Expand Down Expand Up @@ -804,12 +777,7 @@ export const make = Effect.gen(function* () {
.continueThreadsAfterServerUpdate
)
return;
// Shutdown reconciliation cancels the background work below, so a
// settled thread's continuation must be captured while it is still open.
const run = restartContinuationRun(
projection,
providerThreadsWithOpenBackgroundWork(projection),
);
const run = restartContinuationRun(projection);
if (!run) return;
const commandId = CommandId.make(`command:restart-prepare:${run.id}`);
yield* eventSink.writeWithEffects({
Expand Down
209 changes: 144 additions & 65 deletions apps/server/src/orchestration-v2/RestartContinuation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ const driver = ProviderDriverKind.make("codex");
const providerThreadId = ProviderThreadId.make("provider-thread:restart");
const sessionId = ProviderSessionId.make("session:restart");
const attemptId = RunAttemptId.make("attempt:restart");
// "No project" threads belong to the environment's Scratch project.
const scratchProjectId = ProjectId.make("project:scratch");

function makeProjection() {
return {
Expand Down Expand Up @@ -155,76 +157,63 @@ it("recovers an admitted continuation after another crash before provider start"
assert.equal(restartContinuationRun(starting)?.id, runId);
});

it("continues a settled root run only when the restart cancelled its background work", () => {
it("does not continue settled root runs with restart-cancelled background work", () => {
const projection = makeProjection();
const settled = {
...projection,
runs: [{ ...projection.runs[0]!, status: "completed" as const }],
runs: [
{
...projection.runs[0]!,
status: "completed" as const,
restartCancelledBackgroundWork: [{ kind: "shell" as const, label: "sleep 25" }],
},
],
providerThreads: [{ ...projection.providerThreads[0]!, status: "idle" as const }],
providerSessions: [{ ...projection.providerSessions[0]!, status: "stopped" as const }],
providerTurns: [],
};
const lostWork = new Set([providerThreadId]);
assert.isUndefined(restartContinuationRun(settled));
assert.equal(restartContinuationRun(settled, lostWork)?.id, runId);
// Work an older provider thread launched (before a provider switch) is not
// this run's: its provider was never told about it and cannot continue it.
assert.isUndefined(
restartContinuationRun(settled, new Set([ProviderThreadId.make("provider-thread:claude")])),
);
for (const invalid of [
{ ...settled, thread: { ...settled.thread, archivedAt: {} } },
{ ...settled, thread: { ...settled.thread, deletedAt: {} } },
{ ...settled, runs: [{ ...settled.runs[0]!, status: "failed" as const }] },
{
...settled,
providerThreads: [{ ...settled.providerThreads[0]!, nativeThreadRef: null }],
},
])
assert.isUndefined(
restartContinuationRun(invalid as OrchestrationV2ThreadProjection, lostWork),
);
for (const projectId of [scratchProjectId, projection.thread.projectId])
for (const status of ["completed", "waiting"] as const)
assert.isUndefined(
restartContinuationRun({
...settled,
thread: { ...settled.thread, projectId },
runs: [{ ...settled.runs[0]!, status }],
}),
);
});

it.effect("prompts a settled thread's continuation with the note of its lost work", () =>
Effect.gen(function* () {
const base = makeProjection();
const work = [{ kind: "shell" as const, label: "sleep 25 && echo DONE" }];
const projection = {
...base,
runs: [{ ...base.runs[0]!, status: "completed", restartCancelledBackgroundWork: work }],
providerTurns: [{ ...base.providerTurns[0]!, status: "completed" }],
} as unknown as OrchestrationV2ThreadProjection;
const commands: Parameters<
ThreadManagementService.ThreadManagementService["Service"]["dispatch"]
>[0][] = [];
yield* continueRestartedRun({ threadId, sourceRunId: runId }).pipe(
Effect.provide(
Layer.merge(
Layer.mock(ThreadManagementService.ThreadManagementService)({
getThreadRecords: () => Effect.succeed(projection),
recoverDelegatedTask: () => Effect.void,
dispatch: (command) => {
commands.push(command);
return Effect.succeed({} as never);
},
}),
ServerSettings.layerTest({ continueThreadsAfterServerUpdate: true }),
it.effect.each(["completed", "waiting", "cancelled"] as const)(
"ignores a pending restart continuation for a %s run whose provider turn settled",
(status) =>
Effect.gen(function* () {
const base = makeProjection();
const work = [{ kind: "shell" as const, label: "sleep 25 && echo DONE" }];
const projection = {
...base,
runs: [{ ...base.runs[0]!, status, restartCancelledBackgroundWork: work }],
providerTurns: [{ ...base.providerTurns[0]!, status: "completed" }],
} as unknown as OrchestrationV2ThreadProjection;
const commands: Parameters<
ThreadManagementService.ThreadManagementService["Service"]["dispatch"]
>[0][] = [];
yield* continueRestartedRun({ threadId, sourceRunId: runId }).pipe(
Effect.provide(
Layer.merge(
Layer.mock(ThreadManagementService.ThreadManagementService)({
getThreadRecords: () => Effect.succeed(projection),
recoverDelegatedTask: () => Effect.void,
dispatch: (command) => {
commands.push(command);
return Effect.succeed({} as never);
},
}),
ServerSettings.layerTest({ continueThreadsAfterServerUpdate: true }),
),
),
),
);
assert.lengthOf(commands, 1);
const command = commands[0]!;
assert.equal(
command.type === "message.dispatch" ? command.restartContinuationOfRunId : null,
runId,
);
assert.include(
command.type === "message.dispatch" ? command.text : "",
"sleep 25 && echo DONE",
);
assert.notInclude(command.type === "message.dispatch" ? command.text : "", "Continue where");
}),
);
assert.lengthOf(commands, 0);
}),
);

it.effect("does not continue a failed run that lost background work", () =>
Expand Down Expand Up @@ -419,7 +408,7 @@ it.effect("prepares no continuation for background work another provider thread
);
yield* recovery.prepareForShutdown;
assert.lengthOf(writes, 0);
// The same work on the Codex run's own provider thread is continued.
// Work on the settled run's own provider thread must not wake it either.
const ownWork = {
...projection,
turnItems: [{ ...projection.turnItems[0]!, providerThreadId }],
Expand All @@ -446,13 +435,103 @@ it.effect("prepares no continuation for background work another provider thread
),
);
yield* ownRecovery.prepareForShutdown;
assert.deepEqual(
writes.map((write) => write.effects[0]?.request),
[{ type: "provider-runtime.continue", sourceRunId: runId }],
);
assert.lengthOf(writes, 0);
}),
);

it.effect.each([
["completed", scratchProjectId],
["completed", ProjectId.make("restart-project")],
["waiting", scratchProjectId],
["waiting", ProjectId.make("restart-project")],
] as const)(
"cleans up background work without waking a %s run in project %s",
([status, projectId]) =>
Effect.gen(function* () {
const base = makeProjection();
const projection = {
...base,
thread: { ...base.thread, projectId },
runs: [{ ...base.runs[0]!, status }],
providerTurns: [{ ...base.providerTurns[0]!, status: "completed" }],
turnItems: [
{
id: "turn-item:background-subagent",
runId,
nodeId: null,
providerThreadId,
providerInstanceId: instanceId,
type: "subagent",
subagentId: "subagent:background",
title: "Background reviewer",
status: "running",
},
],
subagents: [
{
id: "subagent:background",
runId,
driver,
providerInstanceId: instanceId,
status: "running",
},
],
} as unknown as OrchestrationV2ThreadProjection;
for (const trigger of ["startup", "shutdown"] as const) {
const commits: Parameters<EventSink.EventSinkV2["Service"]["commitCommand"]>[0][] = [];
const writes: Parameters<EventSink.EventSinkV2["Service"]["writeWithEffects"]>[0][] = [];
const recovery = yield* ProviderRuntimeRecovery.make.pipe(
Effect.provide(
Layer.mergeAll(
ServerSettings.layerTest({ continueThreadsAfterServerUpdate: true }),
Layer.mock(ProjectionStore.ProjectionStoreV2)({
getRecoveryThreadIds: () => Effect.succeed([threadId]),
getRuntimeRecoveryProjection: () => Effect.succeed(projection),
}),
Layer.mock(EventSink.EventSinkV2)({
writeWithEffects: (input) =>
Effect.sync(() => {
writes.push(input);
return [];
}),
commitCommand: (input) =>
Effect.sync(() => {
commits.push(input);
return { committed: true, cancelledEffectCount: 0 } as never;
}),
}),
IdAllocator.layer,
Layer.mock(EffectWorker.OrchestrationEffectWorkerV2)({
runRecoveryOnce: Effect.succeed(false),
}),
Layer.mock(EffectOutbox.EffectOutboxV2)({
listByCommandId: () => Effect.succeed([]),
reconcileAfterProcessLoss: Effect.succeed({ requeued: 0, cancelled: 0 }),
}),
),
),
);
if (trigger === "shutdown") yield* recovery.prepareForShutdown;
yield* recovery.reconcile(trigger);
assert.lengthOf(writes, 0);
assert.lengthOf(commits, 1);
assert.lengthOf(commits[0]!.effects, 0);
const events = commits[0]!.events;
for (const type of ["turn-item.updated", "subagent.updated"] as const)
assert.isTrue(
events.some((event) => event.type === type && event.payload.status === "cancelled"),
);
const note = events.find((event) => event.type === "run.background-work-cancelled");
assert.equal(note?.runId, runId);
assert.deepEqual(note?.payload.restartCancelledBackgroundWork, [
{ kind: "subagent", label: "Background reviewer", id: "turn-item:background-subagent" },
]);
if (status === "completed")
assert.isFalse(events.some((event) => event.type === "run.updated"));
}
}),
);

it.effect("does not cancel or resume a run that completes while shutdown intent commits", () =>
Effect.gen(function* () {
let projection = makeProjection();
Expand Down
Loading
Loading