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
288 changes: 288 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts";
import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts";
import {
isTurnItemAtOrBeforeRun,
layerMemory as projectionStoreMemoryLayer,
ProjectionStoreV2,
ProjectionStoreThreadNotFoundError,
layer as projectionStoreLayer,
Expand All @@ -51,6 +52,190 @@ const driver = ProviderDriverKind.make("codex");
const providerInstanceId = modelSelection.instanceId;
const encodeUnknownJsonString = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));

const addRolledBackRecoveryCandidate = Effect.fn("addRolledBackRecoveryCandidate")(function* (
suffix: string,
) {
const projectionStore = yield* ProjectionStoreV2;
const now = yield* DateTime.now;
const threadId = ThreadId.make(`thread:${suffix}:rolled-back`);
const runId = RunId.make(`run:${suffix}:rolled-back`);
const rootNodeId = NodeId.make(`node:${suffix}:rolled-back`);
const run = {
id: runId,
threadId,
ordinal: 1,
providerInstanceId,
modelSelection,
providerThreadId: null,
userMessageId: MessageId.make(`message:${suffix}:rolled-back`),
rootNodeId,
activeAttemptId: null,
status: "running" as const,
requestedAt: now,
startedAt: now,
completedAt: null,
checkpointId: null,
contextHandoffId: null,
};

yield* projectionStore.apply({
id: EventId.make(`event:${suffix}:thread-created`),
type: "thread.created",
threadId,
occurredAt: now,
payload: {
createdBy: "user",
creationSource: "web",
id: threadId,
projectId: ProjectId.make(`project:${suffix}`),
title: "Rolled-back recovery candidate",
providerInstanceId,
modelSelection,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
activeProviderThreadId: null,
lineage: {
parentThreadId: null,
relationshipToParent: null,
rootThreadId: threadId,
},
forkedFrom: null,
createdAt: now,
updatedAt: now,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
},
});
yield* projectionStore.apply({
id: EventId.make(`event:${suffix}:run-created`),
type: "run.created",
threadId,
runId,
nodeId: rootNodeId,
driver,
providerInstanceId,
occurredAt: now,
payload: run,
});
yield* projectionStore.apply({
id: EventId.make(`event:${suffix}:item-running`),
type: "turn-item.updated",
threadId,
runId,
nodeId: rootNodeId,
driver,
occurredAt: now,
payload: {
id: TurnItemId.make(`item:${suffix}:rolled-back`),
threadId,
runId,
nodeId: rootNodeId,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: 1,
status: "running",
title: "abandoned command",
startedAt: now,
completedAt: null,
updatedAt: now,
type: "command_execution",
input: "sleep 60",
},
});
yield* projectionStore.apply({
id: EventId.make(`event:${suffix}:run-rolled-back`),
type: "run.updated",
threadId,
runId,
nodeId: rootNodeId,
driver,
occurredAt: now,
payload: { ...run, status: "rolled_back", completedAt: now },
});

return threadId;
});

const addOrphanedRecoveryCandidate = Effect.fn("addOrphanedRecoveryCandidate")(function* (
suffix: string,
) {
const projectionStore = yield* ProjectionStoreV2;
const now = yield* DateTime.now;
const threadId = ThreadId.make(`thread:${suffix}:orphaned`);
const runId = RunId.make(`run:${suffix}:missing`);
const rootNodeId = NodeId.make(`node:${suffix}:orphaned`);

yield* projectionStore.apply({
id: EventId.make(`event:${suffix}:thread-created`),
type: "thread.created",
threadId,
occurredAt: now,
payload: {
createdBy: "user",
creationSource: "web",
id: threadId,
projectId: ProjectId.make(`project:${suffix}`),
title: "Orphaned recovery candidate",
providerInstanceId,
modelSelection,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
activeProviderThreadId: null,
lineage: {
parentThreadId: null,
relationshipToParent: null,
rootThreadId: threadId,
},
forkedFrom: null,
createdAt: now,
updatedAt: now,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
},
});
yield* projectionStore.apply({
id: EventId.make(`event:${suffix}:item-running`),
type: "turn-item.updated",
threadId,
runId,
nodeId: rootNodeId,
driver,
occurredAt: now,
payload: {
id: TurnItemId.make(`item:${suffix}:orphaned`),
threadId,
runId,
nodeId: rootNodeId,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: 1,
status: "running",
title: "orphaned command",
startedAt: now,
completedAt: null,
updatedAt: now,
type: "command_execution",
input: "sleep 60",
},
});

return threadId;
});

it("includes imported runless history when selecting fork context through a run", () => {
const firstRunId = RunId.make("run:projection-imported-fork:1");
const secondRunId = RunId.make("run:projection-imported-fork:2");
Expand Down Expand Up @@ -93,6 +278,24 @@ it("includes imported runless history when selecting fork context through a run"
);
});

it.effect("memory recovery selection ignores unfinished items from rolled-back runs", () =>
Effect.gen(function* () {
const projectionStore = yield* ProjectionStoreV2;
const threadId = yield* addRolledBackRecoveryCandidate("memory-recovery-candidates");

assert.notInclude(yield* projectionStore.getRecoveryThreadIds("runtime"), threadId);
}).pipe(Effect.provide(projectionStoreMemoryLayer)),
);

it.effect("memory recovery selection includes unfinished items from missing runs", () =>
Effect.gen(function* () {
const projectionStore = yield* ProjectionStoreV2;
const threadId = yield* addOrphanedRecoveryCandidate("memory-recovery-candidates");

assert.include(yield* projectionStore.getRecoveryThreadIds("runtime"), threadId);
}).pipe(Effect.provide(projectionStoreMemoryLayer)),
);

it.layer(TestLayer)("ProjectionStoreV2", (it) => {
it.effect("preserves stored provider usage when a terminal update omits it", () =>
Effect.gen(function* () {
Expand Down Expand Up @@ -1267,6 +1470,91 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => {
}),
);

it.effect("selects only threads with runtime state that needs recovery", () =>
Effect.gen(function* () {
const projectionStore = yield* ProjectionStoreV2;
const now = yield* DateTime.now;
const settledThreadId = ThreadId.make("thread:recovery-candidates:settled");
const runningThreadId = ThreadId.make("thread:recovery-candidates:running");
const rolledBackThreadId = yield* addRolledBackRecoveryCandidate("recovery-candidates");
const orphanedThreadId = yield* addOrphanedRecoveryCandidate("recovery-candidates");
const projectId = ProjectId.make("project:recovery-candidates");
const makeThread = (threadId: ThreadId) => ({
createdBy: "user" as const,
creationSource: "web" as const,
id: threadId,
projectId,
title: "Recovery candidate",
providerInstanceId,
modelSelection,
runtimeMode: "full-access" as const,
interactionMode: "default" as const,
branch: null,
worktreePath: null,
activeProviderThreadId: null,
lineage: {
parentThreadId: null,
relationshipToParent: null,
rootThreadId: threadId,
},
forkedFrom: null,
createdAt: now,
updatedAt: now,
archivedAt: null,
settledOverride: null,
settledAt: null,
lastVisitedAt: null,
deletedAt: null,
});

for (const threadId of [settledThreadId, runningThreadId]) {
yield* projectionStore.apply({
id: EventId.make(`event:recovery-candidates:${threadId}:created`),
type: "thread.created",
threadId,
occurredAt: now,
payload: makeThread(threadId),
});
}

const runId = RunId.make("run:recovery-candidates:running");
const rootNodeId = NodeId.make("node:recovery-candidates:running");
yield* projectionStore.apply({
id: EventId.make("event:recovery-candidates:run-created"),
type: "run.created",
threadId: runningThreadId,
runId,
nodeId: rootNodeId,
driver,
providerInstanceId,
occurredAt: now,
payload: {
id: runId,
threadId: runningThreadId,
ordinal: 1,
providerInstanceId,
modelSelection,
providerThreadId: null,
userMessageId: MessageId.make("message:recovery-candidates:running"),
rootNodeId,
activeAttemptId: null,
status: "running",
requestedAt: now,
startedAt: now,
completedAt: null,
checkpointId: null,
contextHandoffId: null,
},
});

const recoveryThreadIds = yield* projectionStore.getRecoveryThreadIds("runtime");
assert.include(recoveryThreadIds, runningThreadId);
assert.include(recoveryThreadIds, orphanedThreadId);
assert.notInclude(recoveryThreadIds, settledThreadId);
assert.notInclude(recoveryThreadIds, rolledBackThreadId);
}),
);

it.effect("projects one shared provider session into multiple thread bindings", () =>
Effect.gen(function* () {
const projectionStore = yield* ProjectionStoreV2;
Expand Down
11 changes: 8 additions & 3 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -359,7 +359,8 @@ function needsRecovery(
projection.turnItems.some(
(item) =>
["command_execution", "dynamic_tool", "subagent"].includes(item.type) &&
["pending", "running", "waiting"].includes(item.status),
["pending", "running", "waiting"].includes(item.status) &&
!projection.runs.some((run) => run.id === item.runId && run.status === "rolled_back"),
)
);
}
Expand Down Expand Up @@ -2950,8 +2951,12 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
CROSS JOIN orchestration_v2_projection_subagents AS subagents
ON subagents.provider_thread_id = pending_provider_threads.provider_thread_id
UNION
SELECT thread_id FROM orchestration_v2_projection_turn_items
WHERE type IN ('command_execution', 'dynamic_tool', 'subagent')
SELECT item.thread_id FROM orchestration_v2_projection_turn_items AS item
WHERE NOT EXISTS (
SELECT 1 FROM orchestration_v2_projection_runs AS run
WHERE run.run_id = item.run_id AND run.status = 'rolled_back'
)
AND type IN ('command_execution', 'dynamic_tool', 'subagent')
AND status IN ('pending', 'running', 'waiting')
UNION
SELECT thread_id FROM orchestration_v2_effect_outbox
Expand Down
Loading
Loading