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
172 changes: 172 additions & 0 deletions apps/server/src/orchestration-v2/RunExecutionService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,12 @@ import {
ProviderTurnId,
RunAttemptId,
RunId,
ServerSettingsError,
ThreadId,
TurnItemId,
} from "@t3tools/contracts";
import * as Cause from "effect/Cause";
import * as Exit from "effect/Exit";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
Expand Down Expand Up @@ -780,6 +783,175 @@ it.effect("starts the provider when checkpoint baseline capture fails", () =>
}),
);

for (const scenario of ["failure", "interruption", "stale-attempt", "start-guard"] as const) {
it.effect(`handles ${scenario} before the provider turn starts`, () =>
Effect.gen(function* () {
const threadId = ThreadId.make("thread:run-execution-settings-failure");
const runId = RunId.make("run:run-execution-settings-failure");
const attemptId = RunAttemptId.make("attempt:run-execution-settings-failure");
const providerInstanceId = ProviderInstanceId.make("codex");
const providerSessionId = ProviderSessionId.make("session:run-execution-settings-failure");
const providerThreadId = ProviderThreadId.make(
"provider-thread:run-execution-settings-failure",
);
const rootNodeId = NodeId.make("node:run-execution-settings-failure");
const checkpointScope = {
id: CheckpointScopeId.make("checkpoint-scope:run-execution-settings-failure"),
} as OrchestrationV2CheckpointScope;
const providerStarts = yield* Ref.make(0);
const refreshes = yield* Ref.make(0);
const guardedWrites = yield* Ref.make(0);
const writes = yield* Ref.make<ReadonlyArray<ReadonlyArray<OrchestrationV2DomainEvent>>>([]);
const testLayer = runExecutionServiceLayer.pipe(
Layer.provide(
Layer.mergeAll(
Layer.mock(CheckpointServiceV2)({
captureBaseline: () =>
scenario === "start-guard" ? Effect.void : Effect.die("not reached"),
}),
Layer.mock(EventSinkV2)({
writeIfRunCurrent: (input) =>
Effect.gen(function* () {
assert.equal(input.threadId, threadId);
assert.equal(input.runId, runId);
assert.equal(input.activeAttemptId, attemptId);
assert.equal(input.expectedStatus, "running");
yield* Ref.update(guardedWrites, (count) => count + 1);
if (scenario === "stale-attempt") {
return { committed: false, storedEvents: [] };
}
yield* Ref.update(writes, (current) => [...current, input.events]);
return { committed: true, storedEvents: [] };
}),
}),
idAllocatorLayer,
Layer.mock(ProviderEventIngestorV2)({ ingestNormalized: () => Effect.succeed([]) }),
scenario === "start-guard"
? ServerSettingsService.layerTest()
: Layer.mock(ServerSettingsService)({
getSettings:
scenario === "interruption"
? Effect.interrupt
: Effect.fail(
new ServerSettingsError({
settingsPath: "<test>",
operation: "read-file",
cause: new Error("settings read failed"),
}),
),
}),
Layer.succeed(RunFinalizationObserver, {
refresh: () => Effect.void,
refreshAfterTurn: () => Ref.update(refreshes, (count) => count + 1),
}),
),
),
);

const result = yield* Effect.gen(function* () {
const runExecution = yield* RunExecutionServiceV2;
yield* runExecution.startRootRun({
commandId: CommandId.make("command:run-execution-settings-failure"),
appThread: { id: threadId } as OrchestrationV2AppThread,
providerSessionId,
session: {
events: Stream.never,
startTurn: () => Ref.update(providerStarts, (count) => count + 1),
} as unknown as ProviderAdapterV2SessionRuntime,
run: {
id: runId,
threadId,
ordinal: 1,
providerInstanceId,
status: "running",
} as OrchestrationV2Run,
rootNode: { id: rootNodeId, status: "running" } as OrchestrationV2ExecutionNode,
checkpointScope,
providerThread: {
id: providerThreadId,
driver,
} as OrchestrationV2ProviderThread,
attempt: {
id: attemptId,
providerTurnId: null,
status: "running",
} as OrchestrationV2RunAttempt,
attemptId,
providerTurnOrdinal: 1,
// A declined start is a normal exit, not a preparation failure.
...(scenario === "start-guard"
? { shouldStartProviderTurn: () => Effect.succeed(false) }
: {}),
message: {
messageId: MessageId.make("message:run-execution-settings-failure"),
text: "Start after settings fail.",
attachments: [],
createdBy: "user",
creationSource: "web",
},
modelSelection: { instanceId: providerInstanceId, model: "gpt-5.4" },
runtimePolicy: {
runtimeMode: "full-access",
interactionMode: "default",
cwd: process.cwd(),
approvalPolicy: "never",
sandboxPolicy: {
type: "readOnly",
access: { type: "fullAccess" },
networkAccess: false,
},
},
});
}).pipe(Effect.provide(testLayer), Effect.exit);

assert.equal(yield* Ref.get(providerStarts), 0);
const events = (yield* Ref.get(writes)).flat();
if (scenario === "interruption") {
assert.isTrue(Exit.isFailure(result));
if (Exit.isFailure(result)) assert.isTrue(Cause.hasInterruptsOnly(result.cause));
assert.equal(yield* Ref.get(guardedWrites), 0);
assert.equal(yield* Ref.get(refreshes), 0);
assert.isEmpty(events);
return;
}
assert.isTrue(Exit.isSuccess(result));
if (scenario === "start-guard") {
assert.equal(yield* Ref.get(guardedWrites), 0);
assert.equal(yield* Ref.get(refreshes), 0);
assert.isEmpty(events);
return;
}
assert.equal(yield* Ref.get(guardedWrites), 1);
if (scenario === "stale-attempt") {
assert.equal(yield* Ref.get(refreshes), 0);
assert.isEmpty(events);
return;
}
assert.equal(yield* Ref.get(refreshes), 1);
assert.deepEqual(
events
.filter(
(event) =>
event.type === "run.updated" ||
event.type === "run-attempt.updated" ||
event.type === "node.updated",
)
.map((event) => event.payload.status),
["failed", "failed", "failed"],
);
const errorItem = events.find(
(event) => event.type === "turn-item.updated" && event.payload.type === "error",
);
assert.isDefined(errorItem);
if (errorItem?.type === "turn-item.updated" && errorItem.payload.type === "error") {
// The persisted item carries a bounded curated message; the exact
// underlying text stays in the logged cause.
assert.equal(errorItem.payload.failure.message, "Run preparation failed.");
}
}),
);
}

it.effect("keeps ingesting owned child events after the root turn terminalizes", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("thread:run-execution-late-child");
Expand Down
146 changes: 99 additions & 47 deletions apps/server/src/orchestration-v2/RunExecutionService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -560,6 +560,10 @@ export const layer: Layer.Layer<
readonly terminal: ProviderTerminalEvent;
readonly failureItemPersisted: boolean;
readonly refreshAfterTurn: Effect.Effect<void>;
readonly writeIfRunCurrent?: {
readonly activeAttemptId: RunAttemptId;
readonly expectedStatus: OrchestrationV2Run["status"];
};
}) =>
Effect.gen(function* () {
const completedAt = yield* DateTime.now;
Expand Down Expand Up @@ -652,7 +656,7 @@ export const layer: Layer.Layer<
const checkpointCaptureCommandId = CommandId.make(
`command:effect:checkpoint.capture:${input.run.id}`,
);
yield* eventSink.writeWithEffects({
const finalization = {
effects:
input.terminal.status === "completed"
? [
Expand Down Expand Up @@ -760,57 +764,27 @@ export const layer: Layer.Layer<
payload: finalizedProviderThread,
},
],
});
} satisfies Parameters<typeof eventSink.writeWithEffects>[0];
if (input.writeIfRunCurrent !== undefined) {
const result = yield* eventSink.writeIfRunCurrent({
threadId: input.run.threadId,
runId: input.run.id,
activeAttemptId: input.writeIfRunCurrent.activeAttemptId,
expectedStatus: input.writeIfRunCurrent.expectedStatus,
events: finalization.events,
});
if (!result.committed) {
return;
}
} else {
yield* eventSink.writeWithEffects(finalization);
}
yield* input.refreshAfterTurn;
});

return RunExecutionServiceV2.of({
startRootRun: (input) =>
Effect.gen(function* () {
const responseStreamingMode = yield* serverSettings.getSettings.pipe(
Effect.map(
(settings) =>
resolveProjectSettings(settings, input.appThread.projectId).settings
.responseStreamingMode,
),
Effect.mapError(
(cause) =>
new RunExecutionStartError({
commandId: input.commandId,
runId: input.run.id,
cause,
}),
),
);
yield* checkpointService
.captureBaseline({
scope: input.checkpointScope,
ordinalWithinScope: Math.max(0, input.run.ordinal - 1),
})
.pipe(
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause)
? Effect.failCause(cause)
: Effect.logWarning(
"orchestration V2 checkpoint baseline capture failed; starting provider without a baseline",
{ runId: input.run.id },
),
),
Effect.mapError(
(cause) =>
new RunExecutionStartError({
commandId: input.commandId,
runId: input.run.id,
cause,
}),
),
);
if (
input.shouldStartProviderTurn !== undefined &&
!(yield* input.shouldStartProviderTurn())
) {
return;
}
// Startup failure and stream shutdown can report the same attempt.
const refreshAfterTurn = yield* Effect.cached(
finalizationObserver.refreshAfterTurn(input.appThread.projectId).pipe(
Expand All @@ -823,7 +797,6 @@ export const layer: Layer.Layer<
),
),
);
const terminalEvent = yield* Ref.make<ProviderTerminalEvent | null>(null);
const makeFailedTerminalEvent = (
failure: OrchestrationV2ProviderFailure,
failureItemOrdinal: number,
Expand All @@ -843,6 +816,85 @@ export const layer: Layer.Layer<
failure,
threadDisposition: "reusable",
});
const responseStreamingMode = yield* Effect.gen(function* () {
const responseStreamingMode = yield* serverSettings.getSettings.pipe(
Effect.map(
(settings) =>
resolveProjectSettings(settings, input.appThread.projectId).settings
.responseStreamingMode,
),
);
yield* checkpointService
.captureBaseline({
scope: input.checkpointScope,
ordinalWithinScope: Math.max(0, input.run.ordinal - 1),
})
.pipe(
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause)
? Effect.failCause(cause)
: Effect.logWarning(
"orchestration V2 checkpoint baseline capture failed; starting provider without a baseline",
{ runId: input.run.id },
),
),
);
if (
input.shouldStartProviderTurn !== undefined &&
!(yield* input.shouldStartProviderTurn())
) {
return null;
}
return responseStreamingMode;
}).pipe(
Effect.catchCause((cause) =>
Effect.gen(function* () {
if (Cause.hasInterruptsOnly(cause)) {
return yield* Effect.failCause(cause);
}
yield* Effect.logError("orchestration V2 run preparation failed", {
runId: input.run.id,
cause,
});
yield* writeFinalRunEvents({
run: input.run,
rootNode: input.rootNode,
checkpointScope: input.checkpointScope,
providerThread: input.providerThread,
attempt: input.attempt,
terminal: makeFailedTerminalEvent(
makeProviderFailure({
cause: Cause.squash(cause),
// Keep exact underlying text in the logged cause only;
// the persisted turn item gets a bounded curated message.
message: "Run preparation failed.",
class: "unknown",
}),
input.providerTurnOrdinal * 100 + 1,
),
failureItemPersisted: false,
refreshAfterTurn,
writeIfRunCurrent: {
activeAttemptId: input.attemptId,
expectedStatus: "running",
},
});
return null;
}),
),
Effect.mapError(
(cause) =>
new RunExecutionStartError({
commandId: input.commandId,
runId: input.run.id,
cause,
}),
),
);
if (responseStreamingMode === null) {
return;
}
const terminalEvent = yield* Ref.make<ProviderTerminalEvent | null>(null);
const latestTurnItemOrdinal = yield* Ref.make(input.providerTurnOrdinal * 100);
const latestProviderThread = yield* Ref.make(input.providerThread);
const routeIdentity: ProviderEventRouteIdentity = {
Expand Down
Loading