Skip to content
Draft
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
270 changes: 267 additions & 3 deletions apps/server/src/orchestration-v2/Adapters/PiAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,18 +29,29 @@ import * as Schema from "effect/Schema";
import * as Sink from "effect/Sink";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";
import * as Deferred from "effect/Deferred";
import { ChildProcess, ChildProcessSpawner } from "effect/process";

import * as ServerConfig from "../../config.ts";
import * as McpProviderSession from "../../mcp/McpProviderSession.ts";
import * as IdAllocator from "../IdAllocator.ts";
import * as McpSessionRegistry from "../../mcp/McpSessionRegistry.ts";
import * as SqlitePersistence from "../../persistence/Sqlite.ts";
import * as EventSink from "../EventSink.ts";
import * as EventStore from "../EventStore.ts";
import * as ProjectionStore from "../ProjectionStore.ts";
import * as ProviderAdapterRegistry from "../ProviderAdapterRegistry.ts";
import * as ProviderEventIngestor from "../ProviderEventIngestor.ts";
import * as ThreadCommandExecutor from "../ThreadCommandExecutor.ts";
import * as ProviderSessionManager from "../ProviderSessionManager.ts";
import { ProviderAdapterEventStreamError } from "../ProviderAdapter.ts";
import {
ProviderAdapterV2RuntimePolicy,
type ProviderAdapterV2Event,
type ProviderAdapterV2SessionRuntime,
} from "../ProviderAdapter.ts";
import { handoffBudget } from "../ContextHandoffBudget.ts";
import { makePiAdapterV2, PI_PROVIDER } from "./PiAdapterV2.ts";
import { makePiAdapterV2, PI_PROVIDER, piLastErrorAt } from "./PiAdapterV2.ts";
import { makePiRpcConnection, type PiRpcRecord } from "./PiRpc.ts";

const layerServerConfig = ServerConfig.layerTest(process.cwd(), {
Expand Down Expand Up @@ -347,8 +358,10 @@ const openRuntime = Effect.fnUntraced(function* (
runtimePolicy,
});
const emitted = yield* Queue.unbounded<ProviderAdapterV2Event>();
const eventsEnded = yield* Deferred.make<void>();
yield* runtime.events.pipe(
Stream.runForEach((event) => Queue.offer(emitted, event)),
Effect.ensuring(Deferred.succeed(eventsEnded, undefined)),
Effect.forkScoped,
);
const takeEvent = (predicate: (event: ProviderAdapterV2Event) => boolean) =>
Expand All @@ -358,7 +371,7 @@ const openRuntime = Effect.fnUntraced(function* (
if (predicate(event)) return event;
}
});
return { runtime, takeEvent };
return { runtime, takeEvent, eventsEnded: Deferred.await(eventsEnded) };
});

const makeAppThread = Effect.fnUntraced(function* (model: string, threadId = THREAD_ID) {
Expand Down Expand Up @@ -450,7 +463,10 @@ const expectModelFailure = (errorMessage: string) =>
);
assert.isTrue(
sessionError.type === "provider_session.updated" &&
sessionError.providerSession.lastError === errorMessage,
sessionError.providerSession.lastError === errorMessage &&
sessionError.providerSession.lastErrorAt != null &&
DateTime.toEpochMillis(sessionError.providerSession.lastErrorAt) ===
DateTime.toEpochMillis(sessionError.providerSession.updatedAt),
);
const terminal = yield* takeEvent((event) => event.type === "turn.terminal");
assert.isTrue(
Expand All @@ -461,6 +477,254 @@ const expectModelFailure = (errorMessage: string) =>
}).pipe(Effect.scoped, Effect.provide(layerTest));

describe("PiAdapterV2", () => {
it("keeps the occurrence identity on same-text refreshes, including legacy null", () => {
const now = DateTime.makeUnsafe("2026-09-20T12:00:00Z");
for (const previousErrorAt of [null, DateTime.makeUnsafe("2026-09-19T12:00:00Z")]) {
assert.strictEqual(
piLastErrorAt({
previousError: "capacity exhausted",
previousErrorAt,
nextError: "capacity exhausted",
now,
}),
previousErrorAt,
);
}
});

it("clears the occurrence and stamps changed or newly set errors", () => {
const now = DateTime.makeUnsafe("2026-09-20T12:00:00Z");
const previousErrorAt = DateTime.makeUnsafe("2026-09-19T12:00:00Z");
assert.isNull(
piLastErrorAt({ previousError: "capacity exhausted", previousErrorAt, nextError: null, now }),
);
for (const previousError of [null, "different failure"]) {
assert.strictEqual(
piLastErrorAt({
previousError,
previousErrorAt: null,
nextError: "capacity exhausted",
now,
}),
now,
);
}
});

it.effect("emits a new occurrence for the same failure after the next turn clears it", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
const { runtime, takeEvent } = yield* openRuntime(fake);
const providerThread = yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
runtimePolicy,
});
const occurrences: Array<number> = [];
for (const runOrdinal of [1, 2]) {
yield* startTurn(runtime, providerThread, "default", [], "Hello pi", undefined, runOrdinal);
yield* fake.takeRequest("prompt");
const running = yield* takeEvent(
(event) =>
event.type === "provider_session.updated" && event.providerSession.status === "running",
);
assert.equal(running.type, "provider_session.updated");
if (running.type !== "provider_session.updated") return;
assert.isNull(running.providerSession.lastError);
assert.isNull(running.providerSession.lastErrorAt);

yield* TestClock.adjust("1 second");
const failedAt = yield* DateTime.now;
yield* fake.emit({ type: "agent_start" });
yield* fake.emit({
type: "message_end",
message: {
role: "assistant",
content: [],
stopReason: "error",
errorMessage: "capacity exhausted",
},
});
yield* fake.emit({ type: "agent_settled" });
const failed = yield* takeEvent(
(event) =>
event.type === "provider_session.updated" && event.providerSession.status === "error",
);
assert.equal(failed.type, "provider_session.updated");
if (failed.type !== "provider_session.updated") return;
assert.equal(failed.providerSession.lastError, "capacity exhausted");
assert.isNotNull(failed.providerSession.lastErrorAt);
assert.isDefined(failed.providerSession.lastErrorAt);
const occurrence = DateTime.toEpochMillis(failed.providerSession.lastErrorAt!);
assert.equal(occurrence, DateTime.toEpochMillis(failedAt));
occurrences.push(occurrence);
yield* takeEvent((event) => event.type === "turn.terminal");
}
assert.isAbove(occurrences[1]!, occurrences[0]!);
}).pipe(Effect.scoped, Effect.provide(layerTest)),
);

it.effect("retains the occurrence when transport cleanup repeats an unsolicited-work error", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
const { runtime, takeEvent, eventsEnded } = yield* openRuntime(fake);
yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
runtimePolicy,
});
yield* fake.emit({ type: "agent_start" });
const first = yield* takeEvent((event) => event.type === "provider_session.updated");
assert.equal(first.type, "provider_session.updated");
if (first.type !== "provider_session.updated") return;
assert.equal(first.providerSession.status, "error");
assert.isNotNull(first.providerSession.lastErrorAt);
assert.isDefined(first.providerSession.lastErrorAt);
yield* eventsEnded;
const refreshed = yield* takeEvent((event) => event.type === "provider_session.updated");
assert.equal(refreshed.type, "provider_session.updated");
if (refreshed.type !== "provider_session.updated") return;
assert.equal(refreshed.providerSession.lastError, first.providerSession.lastError);
assert.deepEqual(refreshed.providerSession.lastErrorAt, first.providerSession.lastErrorAt);
}).pipe(Effect.scoped, Effect.provide(layerTest)),
);

it.effect.each([false, true])(
"persists Pi occurrence through session-manager stream completion (transport failure: %s)",
(transportFailure) =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
const adapter = yield* makeAdapter(fake, "");
const streamEnded = yield* Deferred.make<void>();
const finishPump = yield* Deferred.make<void>();
const released = yield* Deferred.make<void>();
const stores = Layer.merge(EventStore.layer, ProjectionStore.layer).pipe(
Layer.provide(SqlitePersistence.layerMemory),
);
const sink = EventSink.layer.pipe(
Layer.provide(Layer.merge(stores, SqlitePersistence.layerMemory)),
);
const observedSink = Layer.effect(
EventSink.EventSinkV2,
Effect.gen(function* () {
const delegate = yield* EventSink.EventSinkV2;
return EventSink.EventSinkV2.of({
...delegate,
write: (input) =>
delegate
.write(input)
.pipe(
Effect.tap(() =>
input.events.some(
(event) =>
event.type === "provider-session.updated" &&
DateTime.toEpochMillis(event.payload.updatedAt) > 0,
)
? Deferred.succeed(released, undefined)
: Effect.void,
),
),
});
}),
).pipe(Layer.provide(sink));
const registry = ProviderAdapterRegistry.layerSingle({
...adapter,
openSession: (input) =>
adapter.openSession(input).pipe(
Effect.map((runtime) => ({
...runtime,
events: Stream.concat(
runtime.events,
Stream.fromEffect(
Effect.gen(function* () {
yield* Deferred.succeed(streamEnded, undefined);
yield* Deferred.await(finishPump);
if (transportFailure)
return yield* new ProviderAdapterEventStreamError({
driver: PI_PROVIDER,
providerSessionId: input.providerSessionId,
cause: "distinct transport failure",
});
}),
).pipe(Stream.drain),
),
})),
),
});
const dependencies = Layer.mergeAll(
stores,
observedSink,
IdAllocator.layer,
ThreadCommandExecutor.layer,
);
const managerLayer = ProviderSessionManager.layerWithOptions({
configureMcp: false,
idleTimeoutMs: 60_000,
}).pipe(
Layer.provide(
Layer.mergeAll(
dependencies,
registry,
Layer.mock(McpSessionRegistry.McpSessionRegistry)({}),
ProviderEventIngestor.layer.pipe(Layer.provide(dependencies)),
),
),
);
yield* Effect.gen(function* () {
const manager = yield* ProviderSessionManager.ProviderSessionManagerV2;
const eventSink = yield* EventSink.EventSinkV2;
const store = yield* ProjectionStore.ProjectionStoreV2;
const ids = yield* IdAllocator.IdAllocatorV2;
const now = yield* DateTime.now;
yield* eventSink.write({
events: [
{
id: yield* ids.allocate.event({ threadId: THREAD_ID }),
type: "thread.created",
threadId: THREAD_ID,
occurredAt: now,
payload: yield* makeAppThread("default"),
},
],
});
const runtime = yield* manager.open({
threadId: THREAD_ID,
providerSessionId: SESSION_ID,
modelSelection: modelSelection("default"),
runtimePolicy,
});
yield* runtime.ensureThread({
threadId: THREAD_ID,
modelSelection: modelSelection("default"),
runtimePolicy,
});
yield* fake.emit({ type: "agent_start" });
yield* Deferred.await(streamEnded);
const before = (yield* store.getThreadProviderContext(THREAD_ID)).providerSessions.at(
-1,
)!;
assert.equal(before.status, "error");
assert.include(before.lastError!, "invisible tool execution");
assert.isNotNull(before.lastErrorAt);
assert.isDefined(before.lastErrorAt);
assert.equal(DateTime.toEpochMillis(before.lastErrorAt!), 0);
// Hold only the manager's observation of the already-ended adapter stream.
yield* TestClock.adjust("1 second");
yield* Deferred.succeed(finishPump, undefined);
yield* Deferred.await(released);
const after = (yield* store.getThreadProviderContext(THREAD_ID)).providerSessions.at(-1)!;
assert.equal(after.status, "error");
if (transportFailure) {
assert.include(after.lastError!, "distinct transport failure");
assert.equal(DateTime.toEpochMillis(after.lastErrorAt!), 1000);
} else {
assert.equal(after.lastError, before.lastError);
assert.deepEqual(after.lastErrorAt, before.lastErrorAt);
}
}).pipe(Effect.provide(Layer.merge(dependencies, managerLayer)));
}).pipe(Effect.scoped, Effect.provide(layerTest)),
);

it.effect("stops provider-initiated work that has no T3 turn owner", () =>
Effect.gen(function* () {
const fake = yield* makeFakePi;
Expand Down
20 changes: 19 additions & 1 deletion apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,18 @@ import {
import { PI_FILE_CHANGE_TOOLS } from "./piT3McpExtensionSource.ts";

export const PI_PROVIDER = ProviderDriverKind.make("pi");
export function piLastErrorAt(input: {
readonly previousError: string | null;
readonly previousErrorAt: DateTime.Utc | null;
readonly nextError: string | null;
readonly now: DateTime.Utc;
}): DateTime.Utc | null {
if (input.nextError === null) return null;
// A status refresh is not a new occurrence, even for untimestamped legacy errors.
if (input.nextError === input.previousError) return input.previousErrorAt;
return input.now;
}

const PI_DRIVER_KIND = PI_PROVIDER;
const PI_DEFAULT_INSTANCE_ID = defaultInstanceIdForDriver(PI_DRIVER_KIND);
const DEFAULT_PI_SETTINGS = Schema.decodeSync(PiSettings)({});
Expand Down Expand Up @@ -533,7 +545,13 @@ export function makePiAdapterV2(
) =>
Effect.gen(function* () {
const updatedAt = yield* DateTime.now;
sessionEntity = { ...sessionEntity, status, lastError, updatedAt };
const lastErrorAt = piLastErrorAt({
previousError: sessionEntity.lastError,
previousErrorAt: sessionEntity.lastErrorAt ?? null,
nextError: lastError,
now: updatedAt,
});
sessionEntity = { ...sessionEntity, status, lastError, lastErrorAt, updatedAt };
yield* emit({
type: "provider_session.updated",
driver: PI_PROVIDER,
Expand Down
17 changes: 16 additions & 1 deletion apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2261,13 +2261,15 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => {
const assertSummary = Effect.fnUntraced(function* (
lastError: string | null,
lastErrorClass: string | null,
lastErrorAt: string | null = lastError === null ? null : DateTime.formatIso(now),
) {
const projection = yield* store.getThreadProjection(threadId);
const memoryShell = ProjectionStore.threadShellFromProjection(projection);
const shells = yield* store.getShellSnapshot();
const sqlShell = shells.threads.find((row) => row.id === threadId)!;
for (const shell of [memoryShell, sqlShell]) {
assert.equal(shell.lastError, lastError);
assert.equal(shell.lastErrorAt, lastErrorAt);
assert.equal(shell.lastErrorClass, lastErrorClass);
assert.equal(
shell.usageLimitResetAt,
Expand Down Expand Up @@ -2533,9 +2535,22 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => {
type: "provider-session.updated",
threadId,
occurredAt: now,
payload: {
...session,
status: "error",
lastError: "Provider process exited.",
lastErrorAt: DateTime.makeUnsafe("2026-09-20T02:00:00.000Z"),
},
});
yield* assertSummary("Provider process exited.", null, "2026-09-20T02:00:00.000Z");
yield* store.apply({
id: EventId.make("event:limit-shell:legacy-error"),
type: "provider-session.updated",
threadId,
occurredAt: now,
payload: { ...session, status: "error", lastError: "Provider process exited." },
});
yield* assertSummary("Provider process exited.", null);
yield* assertSummary("Provider process exited.", null, null);
yield* store.apply({
id: EventId.make("event:limit-shell:session-recovered"),
type: "provider-session.updated",
Expand Down
Loading