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
44 changes: 44 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Schema from "effect/Schema";
import * as Tracer from "effect/Tracer";
import * as SqlClient from "effect/sql/SqlClient";
import { projectThreadAwarenessV2 } from "@t3tools/shared/agentAwareness";

Expand Down Expand Up @@ -588,12 +589,26 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => {
});
yield* sql`INSERT INTO orchestration_v2_projection_turn_items ${sql.insert(rows)}`;
}
const statements: Array<string> = [];
const tracer = Tracer.make({
span(options) {
const span = new Tracer.NativeSpan(options);
const end = span.end.bind(span);
span.end = (endTime, exit) => {
end(endTime, exit);
const query = span.attributes.get("db.query.text");
if (typeof query === "string") statements.push(query);
};
return span;
},
});
const initial = yield* Effect.acquireUseRelease(
Effect.sync(() => vi.spyOn(JSON, "parse")),
(parse) =>
projectionStore
.getThreadSnapshotWindow(threadId, { rowLimit: 77, userTurnLimit: 10 })
.pipe(
Effect.withTracer(tracer),
Effect.tap(() =>
Effect.sync(() => {
// Tool outputs must not be allocated again just to collect cohort IDs.
Expand All @@ -606,6 +621,35 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => {
),
(parse) => Effect.sync(() => parse.mockRestore()),
);
// Large threads blocked the server for over a second here. The boundary
// must come from the user-message index rather than a scan of every tool
// row, and payloads must not go through a sort after they are fetched.
const windowStatement = statements.find((statement) => statement.includes("turn_anchors"));
assert.isDefined(windowStatement);
const windowPlan = yield* sql.unsafe<{ readonly parent: number; readonly detail: string }>(
`EXPLAIN QUERY PLAN ${windowStatement}`,
);
assert.match(
windowPlan.map((row) => row.detail).join("\n"),
/SEARCH item USING INDEX orchestration_v2_projection_turn_items_user_message_idx \(thread_id=\? AND ordinal<\?\)/,
);
const topLevel = windowPlan.filter((row) => row.parent === 0).map((row) => row.detail);
assert.include(
topLevel,
"SEARCH item USING INDEX sqlite_autoindex_orchestration_v2_projection_turn_items_1 (turn_item_id=?)",
);
assert.notInclude(topLevel, "USE TEMP B-TREE FOR ORDER BY");
const nodeStatement = statements.find((statement) =>
statement.includes("WITH RECURSIVE retained"),
);
assert.isDefined(nodeStatement);
const nodePlan = yield* sql.unsafe<{ readonly detail: string }>(
`EXPLAIN QUERY PLAN ${nodeStatement}`,
);
assert.include(
nodePlan.map((row) => row.detail),
"SEARCH orchestration_v2_projection_nodes USING INDEX orchestration_v2_projection_nodes_live_idx (thread_id=?)",
);
// Only the selected turn cohort and two lookahead anchors are decoded.
assert.lengthOf(initial.projection.turnItems, 12 * 102);
const bounded = buildBoundedThreadProjection({
Expand Down
24 changes: 16 additions & 8 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2762,7 +2762,10 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
END AS ordinal
FROM user_anchors
), selected AS (
SELECT payload_json, ordinal, turn_item_id, run_id, type
-- Choose and sort rows by ID, then fetch payloads in that
-- order. A turn window can hold megabytes of tool output, and
-- carrying it through the union and sort cost about a second.
SELECT ordinal, turn_item_id, run_id, type
FROM eligible
WHERE ordinal >= (SELECT ordinal FROM boundary)
ORDER BY ordinal DESC, turn_item_id DESC
Expand All @@ -2771,31 +2774,36 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
WHEN (SELECT COUNT(*) FROM user_anchors) > 0 THEN -1
ELSE ${window.rowLimit}
END
), retained AS (
SELECT payload_json, ordinal, turn_item_id FROM selected
), retained AS MATERIALIZED (
SELECT ordinal, turn_item_id FROM selected
UNION
SELECT request.payload_json, request.ordinal, request.turn_item_id
SELECT request.ordinal, request.turn_item_id
FROM orchestration_v2_projection_turn_items AS request
WHERE request.run_id IN (
SELECT run_id FROM selected
WHERE type = 'run_interrupt_result' AND run_id IS NOT NULL
)
AND request.type = 'run_interrupt_request'
UNION
SELECT latest.payload_json, latest.ordinal, latest.turn_item_id
SELECT latest.ordinal, latest.turn_item_id
FROM (
SELECT payload_json, ordinal, turn_item_id
SELECT ordinal, turn_item_id
FROM orchestration_v2_projection_turn_items
WHERE thread_id = ${threadId}
AND ${window.anchorItemId ?? null} IS NULL
AND ${window.requiredRunId ?? null} IS NULL
ORDER BY ordinal DESC, turn_item_id DESC
LIMIT 1
) AS latest
ORDER BY ordinal ASC, turn_item_id ASC
)
SELECT payload_json
-- CROSS JOIN keeps the sorted IDs as the outer loop, so SQLite
-- skips sorting again once the payloads are attached.
SELECT item.payload_json
FROM retained
ORDER BY ordinal ASC, turn_item_id ASC
CROSS JOIN orchestration_v2_projection_turn_items AS item
ON item.turn_item_id = retained.turn_item_id
ORDER BY retained.ordinal ASC, retained.turn_item_id ASC
`;
// Reuse the decoded items for cohort IDs and the resulting projection.
// Parsing these rows separately duplicates every retained tool output.
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/persistence/Migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ import Migration0056 from "./Migrations/056_RemoveRedundantProjectionIndexes.ts"
import Migration0057 from "./Migrations/057_ScheduledTaskWebhooks.ts";
import Migration0058 from "./Migrations/058_WebhookRelayDeliveries.ts";
import Migration0059 from "./Migrations/059_McpAppModelContext.ts";
import Migration0060 from "./Migrations/060_ThreadSnapshotWindowIndexes.ts";

/**
* Migration loader with all migrations defined inline.
Expand Down Expand Up @@ -146,6 +147,7 @@ export const migrationEntries = [
[57, "ScheduledTaskWebhooks", Migration0057],
[58, "WebhookRelayDeliveries", Migration0058],
[59, "McpAppModelContext", Migration0059],
[60, "ThreadSnapshotWindowIndexes", Migration0060],
] as const;

export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ layer("055_OrchestrationV2", (it) => {
Effect.sync(() => {
assert.deepStrictEqual(
migrationEntries.map(([id]) => id),
Array.from({ length: 59 }, (_, index) => index + 1),
Array.from({ length: 60 }, (_, index) => index + 1),
);
}),
);
Expand All @@ -31,6 +31,7 @@ layer("055_OrchestrationV2", (it) => {
[57, "ScheduledTaskWebhooks"],
[58, "WebhookRelayDeliveries"],
[59, "McpAppModelContext"],
[60, "ThreadSnapshotWindowIndexes"],
]);
assert.deepStrictEqual(yield* runMigrations(), []);

Expand All @@ -56,6 +57,7 @@ layer("055_OrchestrationV2", (it) => {
{ migration_id: 57, name: "ScheduledTaskWebhooks" },
{ migration_id: 58, name: "WebhookRelayDeliveries" },
{ migration_id: 59, name: "McpAppModelContext" },
{ migration_id: 60, name: "ThreadSnapshotWindowIndexes" },
]);

const tables = yield* sql<{ readonly name: string }>`
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
import * as Effect from "effect/Effect";
import * as SqlClient from "effect/sql/SqlClient";

export default Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;

// A bounded thread snapshot finds its turn boundary from the newest user
// messages. Without this index it reads every turn item in the thread to
// find a few dozen of them. ProjectionStore's turn_anchors CTE relies on it.
yield* sql`
CREATE INDEX IF NOT EXISTS orchestration_v2_projection_turn_items_user_message_idx
ON orchestration_v2_projection_turn_items(thread_id, ordinal, turn_item_id)
WHERE type = 'user_message'
`;
// The same snapshot keeps every unfinished node. A long thread has thousands
// of finished ones, and reading them all to find a few live ones was more
// than half of the node query's time.
yield* sql`
CREATE INDEX IF NOT EXISTS orchestration_v2_projection_nodes_live_idx
ON orchestration_v2_projection_nodes(thread_id)
WHERE status IN ('pending', 'starting', 'running', 'waiting')
`;
});
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ describe("V2 preview upgrade", () => {
[57, "ScheduledTaskWebhooks"],
[58, "WebhookRelayDeliveries"],
[59, "McpAppModelContext"],
[60, "ThreadSnapshotWindowIndexes"],
]);
assert.deepStrictEqual(yield* runMigrations(), []);
assert.deepStrictEqual(yield* sql`SELECT * FROM orchestration_v2_legacy_imports`, imports);
Expand Down Expand Up @@ -122,6 +123,7 @@ describe("V2 preview upgrade", () => {
[57, "ScheduledTaskWebhooks"],
[58, "WebhookRelayDeliveries"],
[59, "McpAppModelContext"],
[60, "ThreadSnapshotWindowIndexes"],
]);
}).pipe(Effect.provide(NodeSqliteClient.layer({ filename: ":memory:" }))),
);
Expand Down
Loading