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
35 changes: 35 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionSettlement.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -433,3 +433,38 @@ it.effect(
assert.isFalse(messagePlan.some((row) => row.detail.includes("TEMP B-TREE")));
}).pipe(Effect.provide(SqlLayer)),
);

it.effect("shell failure lookups stay on the thread's own turn items", () =>
Effect.gen(function* () {
const store = yield* ProjectionStoreV2;
const sql = yield* SqlClient.SqlClient;
const failed = yield* createThread("failed-latest");
yield* createRun(failed, "failed");
const queries: Array<readonly [string, ReadonlyArray<unknown>]> = [];
const record: Statement.Transformer = (statement) =>
Effect.sync(() => {
queries.push(statement.compile());
return statement;
});
const shell = yield* store
.getShellSnapshot()
.pipe(Effect.provideService(Statement.CurrentTransformer, record));
assert.deepEqual(
shell.threads.map((thread) => thread.id),
[failed],
);
const shellQuery = queries.find(([query]) =>
query.includes("AS blocking_failure_payload_json"),
);
assert.isDefined(shellQuery);
const plan = yield* sql.unsafe<{ readonly detail: string }>(
`EXPLAIN QUERY PLAN ${shellQuery![0]}`,
shellQuery![1],
);
// A failed run's root node is often null, and every runless item shares that
// node_id, so a node_ordinal lookup walks the whole history once per thread.
const itemLookups = plan.filter((row) => row.detail.startsWith("SEARCH item "));
assert.lengthOf(itemLookups, 2);
assert.isTrue(itemLookups.every((row) => row.detail.includes("turn_items_thread_run_idx")));
}).pipe(Effect.provide(SqlLayer)),
);
6 changes: 5 additions & 1 deletion apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4831,13 +4831,15 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
(
SELECT item.payload_json
FROM orchestration_v2_projection_turn_items item
INDEXED BY orchestration_v2_projection_turn_items_thread_run_idx
INNER JOIN orchestration_v2_projection_runs r ON r.run_id = item.run_id
WHERE r.run_id = (
SELECT latest.run_id FROM orchestration_v2_projection_runs latest
WHERE latest.thread_id = t.thread_id
ORDER BY latest.ordinal DESC, latest.run_id DESC LIMIT 1
)
AND r.status = 'failed'
AND item.thread_id = t.thread_id
AND item.type = 'error' AND item.status = 'failed'
AND item.node_id IS json_extract(r.payload_json, '$.rootNodeId')
ORDER BY item.updated_at DESC, item.ordinal DESC, item.turn_item_id DESC
Expand All @@ -4850,7 +4852,9 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
(
SELECT item.payload_json
FROM orchestration_v2_projection_turn_items item
WHERE item.run_id = blocked.run_id
INDEXED BY orchestration_v2_projection_turn_items_thread_run_idx
WHERE item.thread_id = t.thread_id
AND item.run_id = blocked.run_id
AND item.type = 'error' AND item.status = 'failed'
AND item.node_id IS json_extract(blocked.payload_json, '$.rootNodeId')
ORDER BY item.updated_at DESC, item.ordinal DESC, item.turn_item_id DESC
Expand Down
Loading