Skip to content
Closed
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
199 changes: 176 additions & 23 deletions apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,10 @@ import * as ProjectService from "../project/ProjectService.ts";
import { ProviderAuthService } from "../provider/Services/ProviderAuthService.ts";
import * as ContextHandoffService from "./ContextHandoffService.ts";
import * as EventSink from "./EventSink.ts";
import {
ProviderAdapterEventStreamError,
type ProviderAdapterV2SessionRuntime,
} from "./ProviderAdapter.ts";
import * as IdAllocator from "./IdAllocator.ts";
import * as ProjectionStore from "./ProjectionStore.ts";
import * as ProviderSessionManager from "./ProviderSessionManager.ts";
Expand Down Expand Up @@ -158,6 +162,8 @@ function makeLocalCommandHarness(input: {
readonly interruptOpen?: boolean;
readonly interruptRunBeforeOpenFailure?: boolean;
readonly writeFailure?: unknown;
readonly postOpenFailure?: "stream" | "interrupt";
readonly interruptRunBeforePostOpenFailure?: boolean;
}) {
const now = DateTime.makeUnsafe("2026-09-04T12:00:00Z");
const threadId = ThreadId.make("thread-native-account-command");
Expand Down Expand Up @@ -195,7 +201,14 @@ function makeLocalCommandHarness(input: {
providerSessionId,
appThreadId: threadId,
ownerNodeId: null,
nativeThreadRef: null,
nativeThreadRef:
input.postOpenFailure === undefined
? null
: {
driver: ProviderDriverKind.make("pi"),
nativeId: "pi-session",
strength: "strong" as const,
},
nativeConversationHeadRef: null,
status: "not_loaded",
firstRunOrdinal: 2,
Expand Down Expand Up @@ -330,33 +343,65 @@ function makeLocalCommandHarness(input: {
updatedAt: now,
};
const events: Array<OrchestrationV2DomainEvent> = [];
const postOpenStreamError = new ProviderAdapterEventStreamError({
driver: ProviderDriverKind.make("pi"),
providerSessionId,
cause: new Error("Pi RPC read failed: pi process exited with code 1."),
});
const open = vi.fn(() =>
input.interruptOpen === true
? Effect.interrupt
: "openFailure" in input
? Effect.sync(() => {
if (input.interruptRunBeforeOpenFailure === true) {
projection = {
...projection,
runs: projection.runs.map((candidate) =>
candidate.id === runId
? { ...candidate, status: "interrupted", completedAt: now }
: candidate,
: input.postOpenFailure !== undefined
? Effect.succeed({
driver: ProviderDriverKind.make("pi"),
resumeThread: () =>
Effect.sync(() => {
if (input.interruptRunBeforePostOpenFailure === true) {
projection = {
...projection,
runs: projection.runs.map((candidate) =>
candidate.id === runId
? { ...candidate, status: "interrupted" as const, completedAt: now }
: candidate,
),
};
}
}).pipe(
Effect.andThen(
input.postOpenFailure === "interrupt"
? Effect.interrupt
: Effect.fail(postOpenStreamError),
),
};
}
}).pipe(
Effect.andThen(
Effect.fail(
new ProviderSessionManager.ProviderSessionOpenError({
instanceId: newInstanceId,
providerSessionId,
cause: input.openFailure,
}),
),
),
)
: Effect.die("A local command must not open a native session."),
ensureThread: () =>
input.postOpenFailure === "interrupt"
? Effect.interrupt
: Effect.fail(postOpenStreamError),
} as unknown as ProviderAdapterV2SessionRuntime)
: "openFailure" in input
? Effect.sync(() => {
if (input.interruptRunBeforeOpenFailure === true) {
projection = {
...projection,
runs: projection.runs.map((candidate) =>
candidate.id === runId
? { ...candidate, status: "interrupted", completedAt: now }
: candidate,
),
};
}
}).pipe(
Effect.andThen(
Effect.fail(
new ProviderSessionManager.ProviderSessionOpenError({
instanceId: newInstanceId,
providerSessionId,
cause: input.openFailure,
}),
),
),
)
: Effect.die("A local command must not open a native session."),
);
const startRootRun = vi.fn(() => Effect.die("A local command must not start a native turn."));
const tryHandlePromptCommand = vi.fn(() =>
Expand Down Expand Up @@ -560,6 +605,114 @@ effectIt.effect("does not overwrite a run interrupted while its provider session
}),
);

effectIt.effect(
"terminalizes a starting run when resume and ensure fail after the session opens",
() =>
Effect.gen(function* () {
const harness = makeLocalCommandHarness({
text: "Continue",
postOpenFailure: "stream",
});

yield* harness.start;

expect(harness.open).toHaveBeenCalledOnce();
expect(harness.startRootRun).not.toHaveBeenCalled();
expect(harness.writeIfRunCurrent).toHaveBeenCalledWith(
expect.objectContaining({
activeAttemptId: harness.attemptId,
expectedStatus: "starting",
}),
);
const projection = harness.projection();
expect(projection.runs.at(-1)).toMatchObject({ status: "failed", startedAt: null });
expect(projection.attempts[0]).toMatchObject({ status: "failed", startedAt: null });
expect(projection.nodes[0]).toMatchObject({ status: "failed", startedAt: null });
expect(projection.turnItems).toMatchObject([
{
type: "error",
title: "Provider turn failed to start",
status: "failed",
failure: {
class: "provider_error",
message: "Pi RPC read failed: pi process exited with code 1.",
},
},
]);
}),
);

effectIt.effect("leaves the run starting when a post-open failure will be retried", () =>
Effect.gen(function* () {
const harness = makeLocalCommandHarness({
text: "Continue",
postOpenFailure: "stream",
});

const error = yield* harness.startWithRetry.pipe(Effect.flip);

expect(error._tag).toBe("ProviderTurnStartError");
expect(harness.writeIfRunCurrent).not.toHaveBeenCalled();
expect(harness.startRootRun).not.toHaveBeenCalled();
expect(harness.projection().runs.at(-1)?.status).toBe("starting");
}),
);

effectIt.effect("keeps a post-open failure retryable when terminal persistence fails", () =>
Effect.gen(function* () {
const harness = makeLocalCommandHarness({
text: "Continue",
postOpenFailure: "stream",
writeFailure: new Error("database unavailable"),
});

const error = yield* harness.start.pipe(Effect.flip);

expect(error._tag).toBe("ProviderTurnStartError");
expect(harness.writeIfRunCurrent).toHaveBeenCalledOnce();
expect(harness.startRootRun).not.toHaveBeenCalled();
expect(harness.projection().runs.at(-1)?.status).toBe("starting");
expect(harness.events).toEqual([]);
}),
);

effectIt.effect("does not terminalize a post-open interruption", () =>
Effect.gen(function* () {
const harness = makeLocalCommandHarness({ text: "Continue", postOpenFailure: "interrupt" });

const exit = yield* Effect.exit(harness.start);

expect(exit._tag).toBe("Failure");
if (exit._tag === "Failure") {
expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true);
}
expect(harness.writeIfRunCurrent).not.toHaveBeenCalled();
expect(harness.startRootRun).not.toHaveBeenCalled();
expect(harness.projection().runs.at(-1)?.status).toBe("starting");
expect(harness.events).toEqual([]);
}),
);

effectIt.effect("does not overwrite a run interrupted while provider setup fails", () =>
Effect.gen(function* () {
const harness = makeLocalCommandHarness({
text: "Continue",
postOpenFailure: "stream",
interruptRunBeforePostOpenFailure: true,
});

const error = yield* harness.start.pipe(Effect.flip);

expect(error._tag).toBe("ProviderTurnStartError");
expect(harness.writeIfRunCurrent).not.toHaveBeenCalled();
expect(harness.startRootRun).not.toHaveBeenCalled();
const projection = harness.projection();
expect(projection.runs.at(-1)?.status).toBe("interrupted");
expect(projection.turnItems).toEqual([]);
expect(harness.events).toEqual([]);
}),
);

effectIt.effect(
"signs out the existing native provider before opening the newly selected provider",
() =>
Expand Down
Loading
Loading