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
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import type {
ProviderAdapterV2TurnInput,
} from "./ProviderAdapter.ts";
import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts";
import * as ProviderSessionManager from "./ProviderSessionManager.ts";
import * as ProviderReplayHarness from "./testkit/ProviderReplayHarness.ts";
import { checkpointWorkspace } from "./testkit/ReplayFixtureWorkspace.ts";

Expand All @@ -45,10 +46,17 @@ const stopEarlierBackgroundWork = ({
failedStart = false,
stopWithQueue,
olderStart = false,
stalledRun,
}: {
readonly failedStart?: boolean;
readonly stopWithQueue?: "thread.stop" | "run.interrupt";
readonly olderStart?: boolean;
readonly stalledRun?:
| "missing-session"
| "missing-session-terminal"
| "returned-interrupt"
| "returned-interrupt-terminal"
| "superseded-attempt";
}) =>
Effect.scoped(
Effect.gen(function* () {
Expand Down Expand Up @@ -210,6 +218,181 @@ const stopEarlierBackgroundWork = ({
},
],
});
if (stalledRun !== undefined) {
const before = yield* orchestrator.getThreadProjection(threadId);
const run = before.runs[0]!;
const node = before.nodes.find((candidate) => candidate.id === run.rootNodeId)!;
const attempt = before.attempts[0]!;
if (stalledRun === "missing-session" || stalledRun === "missing-session-terminal") {
const sessions = yield* ProviderSessionManager.ProviderSessionManagerV2;
const failed = yield* watch(
(event) => event.type === "run.updated" && event.payload.status === "failed",
);
yield* sessions.release({
providerSessionId: first.providerThread.providerSessionId!,
reason: "runtime_error",
});
yield* Fiber.join(failed);
// Disk exhaustion can lose these terminal writes. Restore that stale state.
yield* sink.write({
events: [
{
id: EventId.make("stale-run"),
type: "run.updated",
threadId,
occurredAt: now,
payload: run,
},
{
id: EventId.make("stale-node"),
type: "node.updated",
threadId,
occurredAt: now,
payload: node,
},
{
id: EventId.make("stale-attempt"),
type: "run-attempt.updated",
threadId,
occurredAt: now,
payload: attempt,
},
],
});
}
const messageId = MessageId.make("partial-output");
const item = before.turnItems.find((candidate) => candidate.id === devServerId)!;
assert.ok(item.type === "command_execution");
const terminalProviderTurn = stalledRun.endsWith("-terminal");
if (terminalProviderTurn) {
yield* sink.write({
events: [
{
id: EventId.make("terminal-provider-turn"),
type: "provider-turn.updated",
threadId,
occurredAt: now,
payload: { ...codexTurn, status: "completed", completedAt: now },
},
...(stalledRun === "missing-session-terminal"
? [
{
id: EventId.make("completed-dev-server"),
type: "turn-item.updated" as const,
threadId,
occurredAt: now,
payload: { ...item, status: "completed" as const, completedAt: now },
},
]
: []),
],
});
}
yield* sink.write({
events: [
{
id: EventId.make("partial-message"),
type: "message.updated",
threadId,
runId: run.id,
occurredAt: now,
payload: {
id: messageId,
threadId,
runId: run.id,
nodeId: node.id,
role: "assistant",
text: "Partial output",
attachments: [],
streaming: true,
createdBy: "agent",
creationSource: "provider",
createdAt: now,
updatedAt: now,
},
},
{
id: EventId.make("partial-item"),
type: "turn-item.updated",
threadId,
runId: run.id,
occurredAt: now,
payload: {
...item,
providerThreadId: codexTurn.providerThreadId,
providerTurnId: codexTurn.id,
id: TurnItemId.make("partial-output"),
type: "assistant_message",
messageId,
text: "Partial output",
streaming: true,
},
},
],
});
if (stalledRun === "superseded-attempt") {
yield* sink.write({
events: [
{
id: EventId.make("new-attempt"),
type: "run.updated",
threadId,
occurredAt: now,
payload: { ...run, activeAttemptId: RunAttemptId.make("new-attempt") },
},
],
});
}
yield* orchestrator.dispatch(
stalledRun === "superseded-attempt"
? {
type: "thread.background-work.settle",
commandId: CommandId.make("late-settle"),
threadId,
providerThreadId: codexTurn.providerThreadId,
providerTurnId: codexTurn.id,
}
: {
type: "thread.stop",
commandId: CommandId.make("stop-stalled-run"),
threadId,
},
);
yield* worker.drain();
const after = yield* orchestrator.getThreadProjection(threadId);
const interrupted = stalledRun !== "superseded-attempt";
assert.equal(after.runs[0]?.status, interrupted ? "interrupted" : "running");
assert.equal(after.attempts[0]?.status, interrupted ? "interrupted" : "running");
assert.equal(
after.providerTurns[0]?.status,
terminalProviderTurn ? "completed" : interrupted ? "interrupted" : "running",
);
assert.equal(
after.nodes.find((candidate) => candidate.id === node.id)?.status,
interrupted ? "interrupted" : "running",
);
const output = after.messages.find((message) => message.id === messageId)!;
assert.equal(output.streaming, !interrupted);
assert.equal(output.text, "Partial output");
assert.equal(
after.turnItems.find((candidate) => candidate.id === devServerId)?.status,
stalledRun === "missing-session-terminal"
? "completed"
: interrupted
? "interrupted"
: "running",
);
assert.equal(
after.turnItems.find((candidate) => candidate.type === "assistant_message")?.status,
interrupted ? "interrupted" : "running",
);
assert.equal(
after.turnItems.filter((candidate) => candidate.type === "run_interrupt_result").length,
interrupted ? 1 : 0,
);
assert.isEmpty(after.runs.filter((candidate) => candidate.status === "waiting"));
return;
}
const settled = yield* watch(
(event) =>
event.type === "run.updated" &&
Expand Down Expand Up @@ -599,3 +782,13 @@ it.effect.each(["thread.stop", "run.interrupt"] as const)(
"%s reaches later background work when an older run is starting",
(stopType) => stopEarlierBackgroundWork({ stopWithQueue: stopType, olderStart: true }),
);

it.effect.each([
"missing-session",
"missing-session-terminal",
"returned-interrupt",
"returned-interrupt-terminal",
"superseded-attempt",
] as const)("Stop recovers a stalled run after %s without changing a newer attempt", (stalledRun) =>
stopEarlierBackgroundWork({ stalledRun }),
);
39 changes: 38 additions & 1 deletion apps/server/src/orchestration-v2/EffectWorker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -81,13 +81,14 @@ function layerExecutorFor(input: {
readonly failFirstStart?: Ref.Ref<boolean>;
readonly threads?: Partial<ThreadManagementService.ThreadManagementService["Service"]>;
readonly continueAfterRestart?: boolean;
readonly interrupt?: ProviderTurnControlService.ProviderTurnControlServiceV2Shape["interrupt"];
}) {
const record = (event: string) => Ref.update(input.events, (events) => [...events, event]);
const layerDependencies = Layer.mergeAll(
Layer.succeed(
ProviderTurnControlService.ProviderTurnControlServiceV2,
ProviderTurnControlService.ProviderTurnControlServiceV2.of({
interrupt: () => Effect.void,
interrupt: input.interrupt ?? (() => Effect.void),
steer: () => Effect.void,
interruptAndAwaitTerminal: (request) =>
record(
Expand Down Expand Up @@ -193,6 +194,42 @@ it("does not retry pure interrupt races where the turn is already gone", () => {
);
});

it.effect("settles a stopped run when its adapter has already lost the native turn", () =>
Effect.gen(function* () {
const now = yield* DateTime.now;
const events = yield* Ref.make<ReadonlyArray<string>>([]);
const layer = layerExecutorFor({
events,
interrupt: () =>
new ProviderTurnControlService.ProviderTurnControlError({
threadId,
operation: "interrupt",
providerTurnId,
cause: "Provider turn is not active.",
}),
threads: {
dispatch: (command) =>
Ref.update(events, (current) => [...current, command.type]).pipe(
Effect.as({ sequence: 1, storedEvents: [] }),
),
},
});
yield* Effect.gen(function* () {
const executor = yield* EffectWorker.OrchestrationEffectExecutorV2;
yield* executor.execute({
...restartEffect(now, { type: "detach" }),
request: {
type: "provider-turn.interrupt",
providerSessionId: oldSessionId,
providerThreadId,
providerTurnId,
},
});
}).pipe(Effect.provide(layer));
assert.deepEqual(yield* Ref.get(events), ["thread.background-work.settle"]);
}),
);

it.effect("requeues a claim when a pre-execution worker check fails", () =>
Effect.gen(function* () {
const now = DateTime.formatIso(yield* DateTime.now);
Expand Down
8 changes: 8 additions & 0 deletions apps/server/src/orchestration-v2/EffectWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,14 @@ export const layerExecutor: Layer.Layer<
providerTurnId: effect.request.providerTurnId,
})
.pipe(
Effect.catch((cause) =>
isNonRetryableProviderTurnControlFailure(
effect.request.type,
Cause.pretty(Cause.fail(cause)),
)
? Effect.void
: Effect.fail(cause),
),
// The provider has stopped what it still ran and reported it.
// Whatever the thread still shows on that provider thread is
// work no process will report on, so the Stop ends it too.
Expand Down
11 changes: 10 additions & 1 deletion apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ export interface EventSinkV2Shape {
readonly activeAttemptId: RunAttemptId;
readonly expectedStatus: OrchestrationV2Run["status"];
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects?: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
}) => Effect.Effect<
{
readonly committed: boolean;
Expand Down Expand Up @@ -437,9 +438,17 @@ const layerBase: Layer.Layer<
events: normalized,
});
yield* applyStoredEvents(storedEvents);
yield* effectOutbox.enqueue(input.effects ?? []);
return { committed: true as const, storedEvents };
}),
(result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void),
(result) =>
Effect.gen(function* () {
if (!result.committed) return;
if (input.effects !== undefined && input.effects.length > 0) {
yield* effectOutbox.notifyAvailable(input.effects.length);
}
yield* publishStoredEvents(result.storedEvents);
}),
);
},
);
Expand Down
29 changes: 28 additions & 1 deletion apps/server/src/orchestration-v2/FoundationPersistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1525,6 +1525,7 @@ it.layer(layerTest)("orchestration V2 foundation persistence", (it) => {
Effect.gen(function* () {
const eventSink = yield* EventSink.EventSinkV2;
const projectionStore = yield* ProjectionStore.ProjectionStoreV2;
const outbox = yield* EffectOutbox.EffectOutboxV2;
const now = yield* DateTime.now;
const threadId = ThreadId.make("thread:foundation-stale-provider-start");
const runId = RunId.make("run:foundation-stale-provider-start");
Expand Down Expand Up @@ -1567,6 +1568,30 @@ it.layer(layerTest)("orchestration V2 foundation persistence", (it) => {
],
});

const captureEffect = {
id: "effect:foundation-current-capture",
commandId: CommandId.make("command:foundation-current-capture"),
threadId,
request: {
type: "checkpoint.capture" as const,
runId,
scopeId: CheckpointScopeId.make("scope:foundation-current-capture"),
},
};
assert.isTrue(
(yield* eventSink.writeIfRunCurrent({
threadId,
runId,
activeAttemptId: attemptId,
expectedStatus: "starting",
events: [],
effects: [captureEffect],
})).committed,
);
assert.isTrue(Option.isSome(yield* outbox.get(captureEffect.id)));
yield* outbox.awaitAvailable;
const staleCaptureEffect = { ...captureEffect, id: "effect:foundation-stale-capture" };

const reachedPrecommitGap = yield* Deferred.make<void>();
const releaseStaleStart = yield* Deferred.make<void>();
const providerStartCount = yield* Ref.make(0);
Expand All @@ -1578,6 +1603,7 @@ it.layer(layerTest)("orchestration V2 foundation persistence", (it) => {
runId,
activeAttemptId: attemptId,
expectedStatus: "starting",
effects: [staleCaptureEffect],
events: [
{
id: EventId.make("event:foundation-stale-provider-start:running"),
Expand Down Expand Up @@ -1619,11 +1645,12 @@ it.layer(layerTest)("orchestration V2 foundation persistence", (it) => {

const staleResult = yield* Fiber.join(staleStartFiber);
assert.isFalse(staleResult.committed);
assert.isTrue(Option.isNone(yield* outbox.get(staleCaptureEffect.id)));
assert.deepEqual(staleResult.storedEvents, []);
assert.equal(yield* Ref.get(providerStartCount), 0);
const projection = yield* projectionStore.getThreadProjection(threadId);
assert.equal(projection.runs[0]?.status, "cancelled");
}),
}).pipe(Effect.provide(Layer.fresh(layerTest))),
);

it.effect("guards post-terminal provider-thread writes by attempt and run ordinal", () =>
Expand Down
Loading
Loading