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
319 changes: 315 additions & 4 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import {
ProjectId,
ProviderItemId,
RuntimeRequestId,
RuntimeItemId,
type ServerSettings,
ThreadId,
TurnId,
Expand Down Expand Up @@ -60,6 +61,7 @@ import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts";
import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import {
ProviderRuntimeIngestionLive,
providerDiffSignal,
splitBufferedAssistantText,
} from "./ProviderRuntimeIngestion.ts";
import { DEFAULT_THREAD_TITLE } from "../threadTitles.ts";
Expand Down Expand Up @@ -239,6 +241,29 @@ async function waitForThread(
}

describe("ProviderRuntimeIngestion", () => {
it("retains only checkpoint identity when queuing a provider diff", () => {
const identity = {
type: "turn.diff.updated" as const,
eventId: asEventId("large-diff"),
threadId: ThreadId.make("thread-1"),
turnId: TurnId.make("turn-1"),
itemId: RuntimeItemId.make("item-1"),
createdAt: "2026-01-01T00:00:00.000Z",
};
const diff = "x".repeat(4_000_000);
const signal = providerDiffSignal({
...identity,
provider: ProviderDriverKind.make("codex"),
payload: { unifiedDiff: diff },
raw: {
source: "codex.app-server.notification",
method: "turn/diff/updated",
payload: { diff },
},
});
expect(signal).toEqual(identity);
});

let runtime: ManagedRuntime.ManagedRuntime<
OrchestrationEngineService | ProviderRuntimeIngestionService | ProjectionSnapshotQuery,
unknown
Expand Down Expand Up @@ -271,6 +296,7 @@ describe("ProviderRuntimeIngestion", () => {
threadTitle?: string;
workspaceSubdirectory?: string;
isGitRepository?: CheckpointStore.CheckpointStore["Service"]["isGitRepository"];
beforeDispatch?: (command: OrchestrationCommand) => Effect.Effect<void>;
}) {
const repositoryRoot = makeTempDir("t3-provider-project-");
NodeChildProcess.execFileSync("git", ["init", "--initial-branch=main"], {
Expand Down Expand Up @@ -322,7 +348,18 @@ describe("ProviderRuntimeIngestion", () => {
};
const layer = ProviderRuntimeIngestionLive.pipe(
Layer.provide(Layer.succeed(Clock.Clock, shiftedClock)),
Layer.provideMerge(orchestrationLayer),
Layer.provideMerge(
Layer.effect(
OrchestrationEngineService,
Effect.map(OrchestrationEngineService, (engine) => ({
...engine,
dispatch: (command, dispatchOptions) =>
(options?.beforeDispatch?.(command) ?? Effect.void).pipe(
Effect.andThen(engine.dispatch(command, dispatchOptions)),
),
})),
).pipe(Layer.provide(orchestrationLayer)),
),
Layer.provideMerge(ingestionProjectionSnapshotLayer),
// Single shared liveness instance across ingestion (writer), the
// engine, and the snapshot query (reader).
Expand Down Expand Up @@ -355,6 +392,8 @@ describe("ProviderRuntimeIngestion", () => {
const drain = () => testRuntime.runPromise(ingestion.drain);
const dispatch = (command: OrchestrationCommand) =>
testRuntime.runPromise(engine.dispatch(command));
const emitAndWaitForEnqueue = (events: ReadonlyArray<LegacyProviderRuntimeEvent>) =>
testRuntime.runPromise(provider.emitAndWaitForEnqueue(events));
const emitAndDrain = (events: ReadonlyArray<LegacyProviderRuntimeEvent>) =>
testRuntime.runPromise(
provider.emitAndWaitForEnqueue(events).pipe(Effect.andThen(ingestion.drain)),
Expand Down Expand Up @@ -433,6 +472,7 @@ describe("ProviderRuntimeIngestion", () => {
advanceClock: (ms: number) => {
clockOffsetMs += ms;
},
emitAndWaitForEnqueue,
emitAndDrain,
sqlCount: sqlCounter.count,
setProviderSession: provider.setSession,
Expand Down Expand Up @@ -4093,16 +4133,277 @@ describe("ProviderRuntimeIngestion", () => {
});
});

effectIt.effect("advances another diff key while unrelated lifecycle work is blocked", () =>
Effect.gen(function* () {
const firstDiffStarted = yield* Deferred.make<void>();
const releaseFirstDiff = yield* Deferred.make<void>();
const lifecycleBlocked = yield* Deferred.make<void>();
const releaseLifecycle = yield* Deferred.make<void>();
const secondDetection = yield* Deferred.make<void>();
let detectionCount = 0;
let blockLifecycle = false;
const harness = yield* Effect.promise(() =>
createHarness({
isGitRepository: () =>
Effect.gen(function* () {
if (++detectionCount === 2) yield* Deferred.succeed(secondDetection, undefined);
return true;
}),
beforeDispatch: (command) => {
if (command.type === "thread.turn.diff.complete" && command.threadId === "thread-1") {
return Deferred.succeed(firstDiffStarted, undefined).pipe(
Effect.andThen(Deferred.await(releaseFirstDiff)),
);
}
if (blockLifecycle && command.type === "thread.session.set") {
return Deferred.succeed(lifecycleBlocked, undefined).pipe(
Effect.andThen(Deferred.await(releaseLifecycle)),
);
}
return Effect.void;
},
}),
);
yield* Effect.addFinalizer(() =>
Deferred.succeed(releaseFirstDiff, undefined).pipe(
Effect.andThen(Deferred.succeed(releaseLifecycle, undefined)),
),
);
const createdAt = "2026-01-01T00:00:00.000Z";
yield* Effect.promise(() =>
harness.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-second-diff-thread"),
threadId: asThreadId("thread-2"),
projectId: asProjectId("project-1"),
title: "Second thread",
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex" },
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
branch: null,
worktreePath: null,
createdAt,
}),
);
const first = {
provider: ProviderDriverKind.make("codex"),
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-1"),
createdAt,
};
const second = { ...first, threadId: asThreadId("thread-2"), turnId: asTurnId("turn-2") };
yield* Effect.promise(() =>
harness.emitAndDrain([
{ ...first, type: "turn.started", eventId: asEventId("start-first") },
{ ...second, type: "turn.started", eventId: asEventId("start-second") },
]),
);
harness.emit({
...first,
type: "turn.diff.updated",
eventId: asEventId("diff-first"),
payload: { unifiedDiff: "first" },
});
yield* Deferred.await(firstDiffStarted);
blockLifecycle = true;
harness.emit({
...first,
type: "session.state.changed",
eventId: asEventId("unrelated-state"),
payload: { state: "running" },
});
harness.emit({
...second,
type: "turn.diff.updated",
eventId: asEventId("diff-second"),
payload: { unifiedDiff: "second" },
});
yield* Deferred.succeed(releaseFirstDiff, undefined);
yield* Deferred.await(lifecycleBlocked);
// The second key must reach detection before the lifecycle queue becomes idle.
yield* Deferred.await(secondDetection);
yield* Deferred.succeed(releaseLifecycle, undefined);
yield* Effect.promise(harness.drain);
const snapshot = yield* Effect.promise(harness.readModel);
expect(
snapshot.threads.find((thread) => thread.id === first.threadId)?.checkpoints,
).toHaveLength(1);
expect(
snapshot.threads.find((thread) => thread.id === second.threadId)?.checkpoints,
).toHaveLength(1);
}),
);

effectIt.effect("keeps a diff key coalesced behind blocked lifecycle work", () =>
Effect.gen(function* () {
const firstDetection = yield* Deferred.make<void>();
const lifecycleBlocked = yield* Deferred.make<void>();
const releaseLifecycle = yield* Deferred.make<void>();
let detectionCount = 0;
let blockLifecycle = false;
const dispatched: string[] = [];
const harness = yield* Effect.promise(() =>
createHarness({
isGitRepository: () =>
Effect.gen(function* () {
if (++detectionCount === 1) yield* Deferred.succeed(firstDetection, undefined);
return true;
}),
beforeDispatch: (command) => {
if (!blockLifecycle) return Effect.void;
if (
command.type === "thread.session.set" ||
command.type === "thread.message.assistant.complete" ||
command.type === "thread.turn.diff.complete"
) {
dispatched.push(`${command.type}:${command.threadId}`);
}
return command.type === "thread.session.set"
? Deferred.succeed(lifecycleBlocked, undefined).pipe(
Effect.andThen(Deferred.await(releaseLifecycle)),
)
: Effect.void;
},
}),
);
yield* Effect.addFinalizer(() => Deferred.succeed(releaseLifecycle, undefined));
const createdAt = "2026-01-01T00:00:00.000Z";
yield* Effect.promise(() =>
harness.dispatch({
type: "thread.create",
commandId: CommandId.make("cmd-backlog-thread"),
threadId: asThreadId("thread-2"),
projectId: asProjectId("project-1"),
title: "Second thread",
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex" },
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "approval-required",
branch: null,
worktreePath: null,
createdAt,
}),
);
const first = {
provider: ProviderDriverKind.make("codex"),
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-1"),
createdAt,
};
const second = { ...first, threadId: asThreadId("thread-2"), turnId: asTurnId("turn-2") };
yield* Effect.promise(() =>
harness.emitAndDrain([
{ ...first, type: "turn.started", eventId: asEventId("start-first") },
{ ...second, type: "turn.started", eventId: asEventId("start-second") },
]),
);

blockLifecycle = true;
harness.emit({
...first,
type: "session.state.changed",
eventId: asEventId("blocked-state"),
payload: { state: "running" },
});
yield* Deferred.await(lifecycleBlocked);
yield* Effect.promise(() =>
harness.emitAndWaitForEnqueue([
{
...second,
type: "item.completed",
eventId: asEventId("backlog-reply"),
itemId: asItemId("backlog-reply"),
payload: {
itemType: "assistant_message",
status: "completed",
detail: "Second reply.",
},
},
]),
);
// Lifecycle backlog is now queued, so this diff's lifecycle item lands behind it.
yield* Effect.promise(() =>
harness.emitAndWaitForEnqueue([
{
...first,
type: "turn.diff.updated",
eventId: asEventId("diff-first-0"),
payload: { unifiedDiff: "first 0" },
},
]),
);
yield* Deferred.await(firstDetection);
for (let index = 1; index <= 20; index++) {
yield* Effect.promise(() =>
harness.emitAndWaitForEnqueue([
{
...first,
type: "turn.diff.updated",
eventId: asEventId(`diff-first-${index}`),
payload: { unifiedDiff: `first ${index}` },
},
]),
);
}
yield* Effect.promise(() =>
harness.emitAndWaitForEnqueue([
{
...second,
type: "turn.diff.updated",
eventId: asEventId("diff-second"),
payload: { unifiedDiff: "second" },
},
]),
);
// The active key waits for its own lifecycle item, so repeated signals merge
// and the second key stays queued behind it.
expect(detectionCount).toBe(1);

yield* Deferred.succeed(releaseLifecycle, undefined);
yield* Effect.promise(harness.drain);
// One extra probe for the merged signal, then one for the second key.
expect(detectionCount).toBe(3);
expect(dispatched).toEqual([
"thread.session.set:thread-1",
"thread.message.assistant.complete:thread-2",
"thread.turn.diff.complete:thread-1",
"thread.turn.diff.complete:thread-2",
]);
const snapshot = yield* Effect.promise(harness.readModel);
const firstThread = snapshot.threads.find((thread) => thread.id === first.threadId);
const secondThread = snapshot.threads.find((thread) => thread.id === second.threadId);
expect(firstThread?.checkpoints).toEqual([
expect.objectContaining({
turnId: first.turnId,
checkpointRef: "provider-diff:diff-first-0",
checkpointTurnCount: 1,
}),
]);
expect(secondThread?.checkpoints).toEqual([
expect.objectContaining({
turnId: second.turnId,
checkpointRef: "provider-diff:diff-second",
checkpointTurnCount: 1,
}),
]);
expect(secondThread?.messages).toEqual(
expect.arrayContaining([expect.objectContaining({ text: "Second reply." })]),
);
}),
);

effectIt.effect("settles the turn while repository detection for a diff is blocked", () =>
Effect.gen(function* () {
let detectionCount = 0;
const detectionStarted = yield* Deferred.make<void>();
const releaseDetection = yield* Deferred.make<boolean>();
const harness = yield* Effect.promise(() =>
createHarness({
isGitRepository: () =>
Deferred.succeed(detectionStarted, undefined).pipe(
isGitRepository: () => {
detectionCount++;
return Deferred.succeed(detectionStarted, undefined).pipe(
Effect.andThen(Deferred.await(releaseDetection)),
),
);
},
}),
);
yield* Effect.addFinalizer(() => Deferred.succeed(releaseDetection, true));
Expand All @@ -4125,6 +4426,15 @@ describe("ProviderRuntimeIngestion", () => {
});
yield* Deferred.await(detectionStarted);

for (let index = 0; index < 200; index++) {
harness.emit({
...base,
type: "turn.diff.updated",
eventId: asEventId(`evt-repeated-diff-${index}`),
payload: { unifiedDiff: "x".repeat(10_000) },
});
}

const settled = yield* harness.engine.streamDomainEvents.pipe(
Stream.filter(
(event) =>
Expand Down Expand Up @@ -4179,6 +4489,7 @@ describe("ProviderRuntimeIngestion", () => {
yield* Fiber.join(nextTurnStarted);
yield* Deferred.succeed(releaseDetection, true);
yield* Effect.promise(harness.drain);
expect(detectionCount).toBeLessThanOrEqual(2);
const released = yield* Effect.promise(harness.readModel);
expect(released.threads[0]?.checkpoints).toEqual([]);
expect(released.threads[0]?.latestTurn).toMatchObject({
Expand Down
Loading
Loading