diff --git a/apps/server/integration/transferBudgetV2.integration.test.ts b/apps/server/integration/transferBudgetV2.integration.test.ts index bf1b6caf5092..e940628e8453 100644 --- a/apps/server/integration/transferBudgetV2.integration.test.ts +++ b/apps/server/integration/transferBudgetV2.integration.test.ts @@ -88,6 +88,8 @@ const layerManagement = Layer.unwrap( projections.getThreadSnapshotWindow(id, options).pipe(Effect.orDie), getThreadShell: (id) => projections.getThreadShell(id).pipe(Effect.orDie), getShellSnapshot: (options) => projections.getShellSnapshot(options).pipe(Effect.orDie), + readShellSnapshot: (options) => + projections.readShellSnapshot(options).pipe(Effect.map(Effect.orDie), Effect.orDie), streamStoredEventsFrom: (input) => sink.stream({ ...input, bounded: true }).pipe(Stream.orDie), }); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index e0d288282b1a..d35426eb2cb2 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -100,6 +100,7 @@ import { type ProjectionRecordFilter, type ProjectionRecords, type ProjectionCheckpointContext, + type ShellSnapshotOptions, } from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import { ProviderAdapterRegistryV2 } from "./ProviderAdapterRegistry.ts"; @@ -319,11 +320,16 @@ export interface OrchestratorV2Shape { }, OrchestratorV2Error >; - readonly getShellSnapshot: (options?: { - readonly location?: "active" | "archive"; - /** Background sweeps only: skips settled threads. */ - readonly unsettledOnly?: boolean; - }) => Effect.Effect; + readonly getShellSnapshot: ( + options?: ShellSnapshotOptions, + ) => Effect.Effect; + /** See `ProjectionStoreV2Shape.readShellSnapshot`. */ + readonly readShellSnapshot: ( + options?: ShellSnapshotOptions, + ) => Effect.Effect< + Effect.Effect, + OrchestratorV2Error + >; readonly getThreadShell: ( threadId: ThreadId, ) => Effect.Effect; @@ -10802,6 +10808,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ); + const shellProjectionError = (cause: unknown) => + new OrchestratorProjectionError({ threadId: ThreadId.make("thread:shell"), cause }); + return OrchestratorV2.of({ resumeQueuedRuns, recoverDelegatedTasks, @@ -10845,15 +10854,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio .getThreadSnapshotWindow(threadId, options) .pipe(Effect.mapError((cause) => new OrchestratorProjectionError({ threadId, cause }))), getShellSnapshot: (options) => - projectionStore.getShellSnapshot(options).pipe( - Effect.mapError( - (cause) => - new OrchestratorProjectionError({ - threadId: ThreadId.make("thread:shell"), - cause, - }), + projectionStore.getShellSnapshot(options).pipe(Effect.mapError(shellProjectionError)), + readShellSnapshot: (options) => + projectionStore + .readShellSnapshot(options) + .pipe( + Effect.mapError(shellProjectionError), + Effect.map(Effect.mapError(shellProjectionError)), ), - ), getThreadShell: (threadId) => projectionStore .getThreadShell(threadId) @@ -10982,6 +10990,13 @@ const layerUnavailable: Layer.Layer = Layer.succeed( cause: "Orchestration V2 live runtime is not configured.", }), ), + readShellSnapshot: () => + Effect.fail( + new OrchestratorProjectionError({ + threadId: ThreadId.make("thread:shell"), + cause: "Orchestration V2 live runtime is not configured.", + }), + ), getThreadShell: (threadId) => Effect.fail( new OrchestratorProjectionError({ diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index 194d7507f2e6..9b48ecd1ac36 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -9,6 +9,7 @@ import { MessageId, type ModelSelection, type OrchestrationV2ProviderThread, + type OrchestrationV2ThreadShellSnapshot, NodeId, ProjectId, ProviderDriverKind, @@ -1728,6 +1729,80 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => { }), ); + it.effect("reads the shell snapshot first and decodes it in a separate step", () => + Effect.gen(function* () { + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const sql = yield* SqlClient.SqlClient; + const now = yield* DateTime.now; + const createThread = (threadId: ThreadId) => + projectionStore.apply({ + id: EventId.make(`event:${threadId}:created`), + type: "thread.created", + threadId, + occurredAt: now, + payload: { + createdBy: "user", + creationSource: "web", + id: threadId, + projectId: ProjectId.make("project:shell-read"), + title: "Shell read", + 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, + }, + }); + const readThreadId = ThreadId.make("thread:shell-read:before"); + const laterThreadId = ThreadId.make("thread:shell-read:after"); + const shellIds = (snapshot: OrchestrationV2ThreadShellSnapshot) => + snapshot.threads + .map((thread) => thread.id) + .filter((id) => id === readThreadId || id === laterThreadId); + + yield* createThread(readThreadId); + const decode = yield* sql.withTransaction(projectionStore.readShellSnapshot()); + // Commits between the read and the decode, as another request can while + // the HTTP and WebSocket loaders decode outside their transaction. + yield* createThread(laterThreadId); + + assert.deepEqual(shellIds(yield* decode), [readThreadId]); + assert.sameMembers(shellIds(yield* projectionStore.getShellSnapshot()), [ + readThreadId, + laterThreadId, + ]); + + // No decoding under the read transaction: a payload that cannot decode + // fails the returned step, not the read. + const [stored] = yield* sql<{ readonly payload_json: string }>` + SELECT payload_json FROM orchestration_v2_projection_threads + WHERE thread_id = ${laterThreadId} + `; + const setPayload = (payload: string) => + sql`UPDATE orchestration_v2_projection_threads SET payload_json = ${payload} + WHERE thread_id = ${laterThreadId}`; + yield* setPayload("{}"); + const failure = yield* sql + .withTransaction(projectionStore.readShellSnapshot()) + .pipe( + Effect.flatMap(Effect.flip), + Effect.ensuring(Effect.orDie(setPayload(stored!.payload_json))), + ); + assert.strictEqual(failure._tag, "ProjectionStoreReadError"); + }), + ); + it.effect("does not treat visited or marked-unread state as thread activity", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStore.ProjectionStoreV2; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 11f98f2330b5..df123755cfc4 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -308,6 +308,15 @@ export interface ProjectionTimelinePage { readonly hasMore: boolean; } +export interface ShellSnapshotOptions { + readonly location?: "active" | "archive"; + /** + * For background sweeps, not clients: skips settled threads before any of + * their run, item or session rows are read. + */ + readonly unsettledOnly?: boolean; +} + export interface ProjectionStoreV2Shape { readonly getThreadAttachmentIds: ( threadId: ThreadId, @@ -334,14 +343,21 @@ export interface ProjectionStoreV2Shape { readonly apply: ( event: OrchestrationV2DomainEvent, ) => Effect.Effect; + readonly getShellSnapshot: ( + options?: ShellSnapshotOptions, + ) => Effect.Effect; /** - * `unsettledOnly` is for background sweeps, not clients: it skips settled - * threads before any of their run, item or session rows are read. + * Runs the shell snapshot's SQL reads and returns the step that decodes them. + * Callers that combine the shell with other reads run this inside their own + * transaction and the returned decode after it commits, so the shared + * connection is not held while thousands of rows are converted. */ - readonly getShellSnapshot: (options?: { - readonly location?: "active" | "archive"; - readonly unsettledOnly?: boolean; - }) => Effect.Effect; + readonly readShellSnapshot: ( + options?: ShellSnapshotOptions, + ) => Effect.Effect< + Effect.Effect, + ProjectionStoreV2Error + >; readonly getThreadShell: ( threadId: ThreadId, ) => Effect.Effect; @@ -5516,103 +5532,125 @@ export const layer: Layer.Layer = } satisfies ShellThreadState; }); - const getShellSnapshot: ProjectionStoreV2Shape["getShellSnapshot"] = (options) => - sql - .withTransaction( - Effect.gen(function* () { - const targetThreadRows = yield* selectShellThreadRows( - undefined, - options?.location, - options?.unsettledOnly ?? false, - ); - const targetThreadIds = new Set( - targetThreadRows.map((row) => ThreadId.make(row.thread_id)), - ); - const rowsByThreadId = new Map( - targetThreadRows.map((row) => [ThreadId.make(row.thread_id), row] as const), - ); - const pendingSourceIds = targetThreadRows.flatMap((row) => - row.forked_from_run_source_thread_id === null - ? [] - : [ThreadId.make(row.forked_from_run_source_thread_id)], - ); - while (pendingSourceIds.length > 0) { - const sourceId = pendingSourceIds.pop(); - if (sourceId === undefined || rowsByThreadId.has(sourceId)) continue; - const source = (yield* selectShellThreadRows(sourceId))[0]; - if (source === undefined) continue; - rowsByThreadId.set(sourceId, source); - if (source.forked_from_run_source_thread_id !== null) { - pendingSourceIds.push(ThreadId.make(source.forked_from_run_source_thread_id)); - } - } - const threadRows = [...rowsByThreadId.values()]; - const threadIds = [...rowsByThreadId.keys()]; - const forkSourceIds = shellForkSourceIds(threadRows); - const readForThreadIds = ( - read: (ids: ReadonlyArray) => Effect.Effect, unknown>, - ids: ReadonlyArray = threadIds, - ) => (ids.length === 0 ? Effect.succeed([] as ReadonlyArray) : read(ids)); - const [runRows, itemCountRows, sequenceRows, providerThreadRows, pendingTurnItemRows] = - yield* Effect.all([ - readForThreadIds(selectShellRunRows, forkSourceIds), - readForThreadIds(selectShellRunItemCounts, forkSourceIds), - sql<{ readonly snapshot_sequence: number | null }>` - SELECT MAX(sequence) AS snapshot_sequence - FROM orchestration_events - WHERE application_event_version = 2 - AND aggregate_kind = 'thread' - `, - readForThreadIds(selectShellProviderThreadRows), - readForThreadIds(selectShellPendingTurnItemRows), - ]); - - const { runOrdinalsByThreadId, itemCountsByThreadId } = runMapsByThreadId({ - runRows, - itemCountRows, - }); - const { providerThreadsByThreadId, pendingTurnItemsByThreadId } = - yield* pendingBackgroundDataByThreadId({ providerThreadRows, pendingTurnItemRows }); - const states = yield* Effect.forEach(threadRows, (row) => - shellThreadStateFromRow({ - row, - runOrdinalsByThreadId, - itemCountsByThreadId, - providerThreadsByThreadId, - pendingTurnItemsByThreadId, - }), - ); - const statesByThreadId = new Map(states.map((state) => [state.thread.id, state])); + const shellSnapshotReadError = (cause: unknown) => + new ProjectionStoreReadError({ threadId: ThreadId.make("thread:shell"), cause }); - const shells = states - .filter((state) => targetThreadIds.has(state.thread.id)) - .map((state) => - shellFromState({ - state, - visibleItemCount: visibleItemCountForShell({ - threadId: state.thread.id, - statesByThreadId, - }), - }), - ); + const readShellSnapshotRows = (options: ShellSnapshotOptions | undefined) => + Effect.gen(function* () { + const targetThreadRows = yield* selectShellThreadRows( + undefined, + options?.location, + options?.unsettledOnly ?? false, + ); + const targetThreadIds = new Set( + targetThreadRows.map((row) => ThreadId.make(row.thread_id)), + ); + const rowsByThreadId = new Map( + targetThreadRows.map((row) => [ThreadId.make(row.thread_id), row] as const), + ); + const pendingSourceIds = targetThreadRows.flatMap((row) => + row.forked_from_run_source_thread_id === null + ? [] + : [ThreadId.make(row.forked_from_run_source_thread_id)], + ); + while (pendingSourceIds.length > 0) { + const sourceId = pendingSourceIds.pop(); + if (sourceId === undefined || rowsByThreadId.has(sourceId)) continue; + const source = (yield* selectShellThreadRows(sourceId))[0]; + if (source === undefined) continue; + rowsByThreadId.set(sourceId, source); + if (source.forked_from_run_source_thread_id !== null) { + pendingSourceIds.push(ThreadId.make(source.forked_from_run_source_thread_id)); + } + } + const threadRows = [...rowsByThreadId.values()]; + const threadIds = [...rowsByThreadId.keys()]; + const forkSourceIds = shellForkSourceIds(threadRows); + const readForThreadIds = ( + read: (ids: ReadonlyArray) => Effect.Effect, unknown>, + ids: ReadonlyArray = threadIds, + ) => (ids.length === 0 ? Effect.succeed([] as ReadonlyArray) : read(ids)); + const [runRows, itemCountRows, sequenceRows, providerThreadRows, pendingTurnItemRows] = + yield* Effect.all([ + readForThreadIds(selectShellRunRows, forkSourceIds), + readForThreadIds(selectShellRunItemCounts, forkSourceIds), + sql<{ readonly snapshot_sequence: number | null }>` + SELECT MAX(sequence) AS snapshot_sequence + FROM orchestration_events + WHERE application_event_version = 2 + AND aggregate_kind = 'thread' + `, + readForThreadIds(selectShellProviderThreadRows), + readForThreadIds(selectShellPendingTurnItemRows), + ]); + return { + targetThreadIds, + threadRows, + runRows, + itemCountRows, + sequenceRows, + providerThreadRows, + pendingTurnItemRows, + }; + }); - return { - schemaVersion: ORCHESTRATION_V2_PROJECTION_SCHEMA_VERSION, - snapshotSequence: sequenceRows[0]?.snapshot_sequence ?? 0, - threads: shells.filter((thread) => thread.archivedAt === null), - archivedThreads: shells.filter((thread) => thread.archivedAt !== null), - }; + const decodeShellSnapshot = ({ + targetThreadIds, + threadRows, + runRows, + itemCountRows, + sequenceRows, + providerThreadRows, + pendingTurnItemRows, + }: Effect.Success>) => + Effect.gen(function* () { + const { runOrdinalsByThreadId, itemCountsByThreadId } = runMapsByThreadId({ + runRows, + itemCountRows, + }); + const { providerThreadsByThreadId, pendingTurnItemsByThreadId } = + yield* pendingBackgroundDataByThreadId({ providerThreadRows, pendingTurnItemRows }); + const states = yield* Effect.forEach(threadRows, (row) => + shellThreadStateFromRow({ + row, + runOrdinalsByThreadId, + itemCountsByThreadId, + providerThreadsByThreadId, + pendingTurnItemsByThreadId, }), - ) - .pipe( - Effect.mapError( - (cause) => - new ProjectionStoreReadError({ - threadId: ThreadId.make("thread:shell"), - cause, - }), - ), ); + const statesByThreadId = new Map(states.map((state) => [state.thread.id, state])); + + const shells = states + .filter((state) => targetThreadIds.has(state.thread.id)) + .map((state) => + shellFromState({ + state, + visibleItemCount: visibleItemCountForShell({ + threadId: state.thread.id, + statesByThreadId, + }), + }), + ); + + return { + schemaVersion: ORCHESTRATION_V2_PROJECTION_SCHEMA_VERSION, + snapshotSequence: sequenceRows[0]?.snapshot_sequence ?? 0, + threads: shells.filter((thread) => thread.archivedAt === null), + archivedThreads: shells.filter((thread) => thread.archivedAt !== null), + }; + }).pipe(Effect.mapError(shellSnapshotReadError)); + + const readShellSnapshot: ProjectionStoreV2Shape["readShellSnapshot"] = (options) => + readShellSnapshotRows(options).pipe( + Effect.mapError(shellSnapshotReadError), + Effect.map(decodeShellSnapshot), + ); + + const getShellSnapshot: ProjectionStoreV2Shape["getShellSnapshot"] = (options) => + sql + .withTransaction(readShellSnapshotRows(options)) + .pipe(Effect.mapError(shellSnapshotReadError), Effect.flatMap(decodeShellSnapshot)); // Per-thread shell for the live shell streams: reads only the target thread // plus its fork-source chain instead of materializing every thread. Returns @@ -5687,6 +5725,7 @@ export const layer: Layer.Layer = return { apply, getShellSnapshot, + readShellSnapshot, getThreadShell, getThread, getSettlementCandidates, @@ -5780,6 +5819,8 @@ export const layerMemory: Layer.Layer = Layer.effect( archivedThreads: visible.filter((thread) => thread.archivedAt !== null), }; }), + readShellSnapshot: (options) => + service.getShellSnapshot(options).pipe(Effect.map(Effect.succeed)), getThreadShell: (threadId) => Effect.gen(function* () { const existing = (yield* Ref.get(replayState)).projections; diff --git a/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts index 7918f819eee2..09c80f66fa2c 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts @@ -213,6 +213,7 @@ it.effect( apply: () => Effect.void, getLimitRecoveryCandidates: () => Effect.die("unused getLimitRecoveryCandidates"), getShellSnapshot: () => Effect.die("unused getShellSnapshot"), + readShellSnapshot: () => Effect.die("unused readShellSnapshot"), getThreadShell: () => Effect.die("unused getThreadShell"), getThread: () => Ref.get(projection).pipe(Effect.map((state) => state.thread)), getSettlementCandidates: () => Effect.die("unused getSettlementCandidates"), diff --git a/apps/server/src/orchestration-v2/ShellStream.test.ts b/apps/server/src/orchestration-v2/ShellStream.test.ts index b9c818c3f4c7..fccff41a5d46 100644 --- a/apps/server/src/orchestration-v2/ShellStream.test.ts +++ b/apps/server/src/orchestration-v2/ShellStream.test.ts @@ -7,9 +7,12 @@ import type { OrchestrationV2ThreadShell, } from "@t3tools/contracts"; import { ProjectId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; +import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient"; import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; +import * as Option from "effect/Option"; +import * as SqlClient from "effect/sql/SqlClient"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; @@ -20,6 +23,7 @@ import { coalesceStoredThreadEvents, composeShellStreamWithEnrichment, dedupeShellEnrichment, + loadShellSnapshotParts, shellStreamItemFromEnrichmentRefresh, shellStreamItemFromThreadShell, shellStreamItemsFromInitialSnapshot, @@ -101,6 +105,47 @@ function storedThreadEvent( return { sequence, event: { threadId, ...event } } as OrchestrationV2StoredEvent; } +describe("loadShellSnapshotParts", () => { + it.effect("decodes the threads after the read transaction releases the connection", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const steps: Array = []; + const step = (name: string, value: A) => + Effect.serviceOption(sql.transactionService).pipe( + Effect.map((transaction) => { + steps.push(`${name}:${Option.isSome(transaction) ? "in" : "out of"} transaction`); + return value; + }), + ); + const thread = shellFixture({}); + + const snapshot = yield* loadShellSnapshotParts({ + sql, + readThreads: step( + "read threads", + step("decode threads", { + schemaVersion: 1, + snapshotSequence: 3, + threads: [thread], + archivedThreads: [], + }), + ), + listProjects: step("list projects", []), + latestSequence: step("latest sequence", 7), + }); + + expect(steps).toEqual([ + "read threads:in transaction", + "list projects:in transaction", + "latest sequence:in transaction", + "decode threads:out of transaction", + ]); + expect(snapshot.snapshotSequence).toBe(7); + expect(snapshot.threads.threads).toEqual([thread]); + }).pipe(Effect.provide(NodeSqliteClient.layer({ filename: ":memory:" }))), + ); +}); + function shellFixture(overrides: Partial): OrchestrationV2ThreadShell { return { id: "thread-a", archivedAt: null, ...overrides } as OrchestrationV2ThreadShell; } diff --git a/apps/server/src/orchestration-v2/ShellStream.ts b/apps/server/src/orchestration-v2/ShellStream.ts index 1464f7a211a9..cec809f45ad4 100644 --- a/apps/server/src/orchestration-v2/ShellStream.ts +++ b/apps/server/src/orchestration-v2/ShellStream.ts @@ -15,6 +15,7 @@ import { import * as Clock from "effect/Clock"; import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; +import type * as SqlClient from "effect/sql/SqlClient"; import * as Stream from "effect/Stream"; /** Build the regular navigation shell without duplicating the archive dataset. */ @@ -32,6 +33,33 @@ export function buildActiveShellSnapshot(input: { }; } +/** + * Loads the thread shell, projects and sequence for the shell snapshots. + * One transaction covers the reads so they agree; the threads are decoded + * after it commits, because the server shares one SQLite connection and + * decoding a large shell takes longer than reading it. + */ +export const loadShellSnapshotParts = (input: { + readonly sql: SqlClient.SqlClient; + readonly readThreads: Effect.Effect, E1>; + readonly listProjects: Effect.Effect, E3>; + readonly latestSequence: Effect.Effect; +}) => + Effect.gen(function* () { + const read = yield* input.sql.withTransaction( + Effect.all({ + decodeThreads: input.readThreads, + projects: input.listProjects, + snapshotSequence: input.latestSequence, + }), + ); + return { + projects: read.projects, + threads: yield* read.decodeThreads, + snapshotSequence: read.snapshotSequence, + }; + }); + export type ShellApplicationEvent = | Pick< Extract, diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.ts b/apps/server/src/orchestration-v2/ThreadManagementService.ts index a450f34df5f0..b48beb08875c 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.ts @@ -309,6 +309,12 @@ export interface ThreadManagementServiceShape { readonly getShellSnapshot: (options?: { readonly location?: "active" | "archive"; }) => Effect.Effect; + readonly readShellSnapshot: (options?: { + readonly location?: "active" | "archive"; + }) => Effect.Effect< + Effect.Effect, + Orchestrator.OrchestratorV2Error + >; readonly getThreadShell: Orchestrator.OrchestratorV2["Service"]["getThreadShell"]; readonly listProjectThreads: (input: { readonly projectId: ProjectId; @@ -818,6 +824,7 @@ const make = Effect.gen(function* () { getProjectThreadRecords, getProjectThread, getShellSnapshot: orchestrator.getShellSnapshot, + readShellSnapshot: orchestrator.readShellSnapshot, getThreadShell: orchestrator.getThreadShell, listProjectThreads, sendToThread, diff --git a/apps/server/src/orchestration-v2/http.ts b/apps/server/src/orchestration-v2/http.ts index e6b769d98c40..8f8c30034c44 100644 --- a/apps/server/src/orchestration-v2/http.ts +++ b/apps/server/src/orchestration-v2/http.ts @@ -31,7 +31,7 @@ import { } from "./threadHistoryPaging.ts"; import * as ThreadManagementService from "./ThreadManagementService.ts"; import * as ProjectStore from "./ProjectStore.ts"; -import { buildActiveShellSnapshot } from "./ShellStream.ts"; +import { buildActiveShellSnapshot, loadShellSnapshotParts } from "./ShellStream.ts"; import { boundedSnapshotResponseFields } from "./ThreadStream.ts"; import { projectThreadProjectionForWire } from "./WireProjection.ts"; @@ -95,14 +95,12 @@ export const layer = HttpApiBuilder.group( ); const loadShellSnapshot = Effect.fn("http.orchestration.loadShellSnapshot")(function* () { - const base = yield* sql.withTransaction( - Effect.gen(function* () { - const threads = yield* threadManagement.getShellSnapshot({ location: "active" }); - return buildActiveShellSnapshot({ - projects: yield* projectStore.listShells(), - threads, - snapshotSequence: yield* applicationEvents.latestApplicationSequence, - }); + const base = buildActiveShellSnapshot( + yield* loadShellSnapshotParts({ + sql, + readThreads: threadManagement.readShellSnapshot({ location: "active" }), + listProjects: projectStore.listShells(), + latestSequence: applicationEvents.latestApplicationSequence, }), ); const projects = yield* enrichProjectShells(base.projects); diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index 79f100fb46e9..ee37b9d5a5ae 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -190,6 +190,7 @@ const makeTestRelay = Effect.fnUntraced(function* ( catchUp.shellSnapshotReads += 1; return { schemaVersion: 2, snapshotSequence: 1, threads: [], archivedThreads: [] }; }), + readShellSnapshot: unused, ensureLegacyTranscript: unused, dispatch: unused, getTimelinePage: () => Effect.die("Unused timeline read"), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 53bf27800583..70238f10a882 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -130,6 +130,7 @@ import { coalesceStoredThreadEvents, composeShellStreamWithEnrichment, dedupeShellEnrichment, + loadShellSnapshotParts, shellStreamItemFromEnrichmentRefresh, shellStreamItemFromThreadShell, shellStreamItemsFromInitialSnapshot, @@ -967,14 +968,12 @@ export const subscribeOrchestrationV2Shell = Effect.fn("ws.orchestrationV2.subsc }, ); const loadSnapshot = Effect.fn("ws.orchestrationV2.loadShellSnapshot")(function* () { - const base = yield* sql.withTransaction( - Effect.gen(function* () { - const threads = yield* threadManagement.getShellSnapshot({ location: "active" }); - return buildActiveShellSnapshot({ - projects: yield* projects.listShells(), - threads, - snapshotSequence: yield* applicationEvents.latestApplicationSequence, - }); + const base = buildActiveShellSnapshot( + yield* loadShellSnapshotParts({ + sql, + readThreads: threadManagement.readShellSnapshot({ location: "active" }), + listProjects: projects.listShells(), + latestSequence: applicationEvents.latestApplicationSequence, }), ); const enriched = yield* enrichProjectShells(base.projects); @@ -1741,32 +1740,33 @@ const layerWsRpc = ( .refreshStatus(cwd) .pipe(Effect.ignoreCause({ log: true }), Effect.forkDetach, Effect.asVoid); - const getOrchestrationV2ArchivedShellSnapshot = sql - .withTransaction( - Effect.gen(function* () { - const threads = yield* threadManagement.getShellSnapshot({ location: "archive" }); - return { - schemaVersion: threads.schemaVersion, - snapshotSequence: yield* applicationEvents.latestApplicationSequence, - projects: yield* projectStore.listShells(), - threads: threads.archivedThreads, - } as const; - }), - ) - .pipe( - Effect.flatMap((snapshot) => - enrichProjectShells(snapshot.projects).pipe( - Effect.map(({ projects }) => ({ ...snapshot, projects })), - ), - ), - Effect.mapError( - (cause) => - new OrchestrationV2GetShellSnapshotError({ - message: "Failed to load archived thread snapshot", - cause, - }), + const getOrchestrationV2ArchivedShellSnapshot = Effect.gen(function* () { + const { threads, projects, snapshotSequence } = yield* loadShellSnapshotParts({ + sql, + readThreads: threadManagement.readShellSnapshot({ location: "archive" }), + listProjects: projectStore.listShells(), + latestSequence: applicationEvents.latestApplicationSequence, + }); + return { + schemaVersion: threads.schemaVersion, + snapshotSequence, + projects, + threads: threads.archivedThreads, + } as const; + }).pipe( + Effect.flatMap((snapshot) => + enrichProjectShells(snapshot.projects).pipe( + Effect.map(({ projects }) => ({ ...snapshot, projects })), ), - ); + ), + Effect.mapError( + (cause) => + new OrchestrationV2GetShellSnapshotError({ + message: "Failed to load archived thread snapshot", + cause, + }), + ), + ); const subscribeOrchestrationV2ArchivedShell = Effect.fn( "ws.orchestrationV2.subscribeArchivedShell",