diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index e8a3c1bcfc08..1008c4a2f58f 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -349,6 +349,157 @@ it.effect("memory recovery selection includes unfinished items from missing runs }).pipe(Effect.provide(ProjectionStore.layerMemory)), ); +// Project A holds an active thread, an archived thread and a fork of a run in +// an archived project B thread, so a project-scoped snapshot still has to load +// a fork source that is outside the project and archived. +const assertProjectScopedShells = Effect.fn("assertProjectScopedShells")(function* ( + suffix: string, +) { + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const projectId = ProjectId.make(`project:${suffix}:a`); + const otherProjectId = ProjectId.make(`project:${suffix}:b`); + const activeThreadId = ThreadId.make(`thread:${suffix}:a-active`); + const archivedThreadId = ThreadId.make(`thread:${suffix}:a-archived`); + const forkThreadId = ThreadId.make(`thread:${suffix}:a-fork`); + const sourceThreadId = ThreadId.make(`thread:${suffix}:b-source`); + const otherThreadId = ThreadId.make(`thread:${suffix}:b-other`); + const sourceRunId = RunId.make(`run:${suffix}:b-source`); + const sourceNodeId = NodeId.make(`node:${suffix}:b-source`); + const createThread = (input: { + readonly threadId: ThreadId; + readonly projectId: ProjectId; + readonly archived?: boolean; + readonly forkedFromRunOf?: ThreadId; + }) => + projectionStore.apply({ + id: EventId.make(`event:${input.threadId}:created`), + type: "thread.created", + threadId: input.threadId, + occurredAt: now, + payload: { + createdBy: "user", + creationSource: "web", + id: input.threadId, + projectId: input.projectId, + title: input.threadId, + providerInstanceId, + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: null, + lineage: + input.forkedFromRunOf === undefined + ? { parentThreadId: null, relationshipToParent: null, rootThreadId: input.threadId } + : { + parentThreadId: input.forkedFromRunOf, + relationshipToParent: "fork", + rootThreadId: input.forkedFromRunOf, + }, + forkedFrom: + input.forkedFromRunOf === undefined + ? null + : { type: "run", threadId: input.forkedFromRunOf, runId: sourceRunId }, + createdAt: now, + updatedAt: now, + archivedAt: input.archived === true ? now : null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + deletedAt: null, + }, + }); + + yield* createThread({ threadId: sourceThreadId, projectId: otherProjectId, archived: true }); + yield* createThread({ threadId: otherThreadId, projectId: otherProjectId }); + yield* createThread({ threadId: activeThreadId, projectId }); + yield* createThread({ threadId: archivedThreadId, projectId, archived: true }); + yield* createThread({ threadId: forkThreadId, projectId, forkedFromRunOf: sourceThreadId }); + yield* projectionStore.apply({ + id: EventId.make(`event:${suffix}:source-run`), + type: "run.created", + threadId: sourceThreadId, + runId: sourceRunId, + nodeId: sourceNodeId, + driver, + occurredAt: now, + payload: { + id: sourceRunId, + threadId: sourceThreadId, + ordinal: 1, + providerInstanceId, + modelSelection, + providerThreadId: null, + userMessageId: MessageId.make(`message:${suffix}:source-user`), + rootNodeId: sourceNodeId, + activeAttemptId: null, + status: "completed", + requestedAt: now, + startedAt: now, + completedAt: now, + checkpointId: null, + contextHandoffId: null, + }, + }); + yield* projectionStore.apply({ + id: EventId.make(`event:${suffix}:source-item`), + type: "turn-item.updated", + threadId: sourceThreadId, + runId: sourceRunId, + occurredAt: now, + payload: { + id: TurnItemId.make(`turn-item:${suffix}:source`), + threadId: sourceThreadId, + runId: sourceRunId, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "completed", + title: null, + startedAt: now, + completedAt: now, + updatedAt: now, + type: "assistant_message", + messageId: MessageId.make(`message:${suffix}:source-assistant`), + text: "Source answer", + streaming: false, + }, + }); + + const unscoped = yield* projectionStore.getShellSnapshot(); + const inProject = (thread: { readonly projectId: ProjectId }) => thread.projectId === projectId; + const active = yield* projectionStore.getShellSnapshot({ projectId, location: "active" }); + assert.deepEqual(active.threads, unscoped.threads.filter(inProject)); + assert.deepEqual(active.archivedThreads, []); + assert.sameMembers( + active.threads.map((thread) => thread.id), + [activeThreadId, forkThreadId], + ); + const fork = active.threads.find((thread) => thread.id === forkThreadId); + assert.isAbove(fork?.visibleItemCount ?? 0, 0); + + const all = yield* projectionStore.getShellSnapshot({ projectId }); + assert.deepEqual(all.threads, active.threads); + assert.deepEqual(all.archivedThreads, unscoped.archivedThreads.filter(inProject)); + assert.sameMembers( + all.archivedThreads.map((thread) => thread.id), + [archivedThreadId], + ); + return { projectId, activeThreadId, forkThreadId, otherThreadId }; +}); + +it.effect("memory shell snapshots scope to one project", () => + assertProjectScopedShells("memory-project-shell").pipe( + Effect.asVoid, + Effect.provide(ProjectionStore.layerMemory), + ), +); + it.layer(layerTest)("ProjectionStoreV2", (it) => { it.effect( "keeps restart-cancelled work through a stale run.updated", @@ -1884,6 +2035,35 @@ it.layer(layerTest)("ProjectionStoreV2", (it) => { }), ); + it.effect("reads only the requested project's threads into a scoped shell snapshot", () => + Effect.gen(function* () { + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const sql = yield* SqlClient.SqlClient; + const { projectId, activeThreadId, forkThreadId, otherThreadId } = + yield* assertProjectScopedShells("sql-project-shell"); + + // A thread in another project that cannot decode fails the unscoped + // snapshot, so a scoped snapshot that succeeds never read it. + const [stored] = yield* sql<{ readonly payload_json: string }>` + SELECT payload_json FROM orchestration_v2_projection_threads + WHERE thread_id = ${otherThreadId} + `; + const setPayload = (payload: string) => + sql`UPDATE orchestration_v2_projection_threads SET payload_json = ${payload} + WHERE thread_id = ${otherThreadId}`; + yield* setPayload("{}"); + const [unscoped, scoped] = yield* Effect.all([ + projectionStore.getShellSnapshot().pipe(Effect.flip), + projectionStore.getShellSnapshot({ projectId, location: "active" }), + ]).pipe(Effect.ensuring(Effect.orDie(setPayload(stored!.payload_json)))); + assert.strictEqual(unscoped._tag, "ProjectionStoreReadError"); + assert.sameMembers( + scoped.threads.map((thread) => thread.id), + [activeThreadId, forkThreadId], + ); + }), + ); + 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 51343fd8af91..22e87e0ec632 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -31,6 +31,7 @@ import type { OrchestrationV2ThreadShell, OrchestrationV2ThreadProjection, OrchestrationV2TurnItem, + ProjectId, ProviderInstanceId, ProviderSessionId, ProviderThreadId, @@ -325,6 +326,11 @@ export interface ProjectionTimelinePage { export interface ShellSnapshotOptions { readonly location?: "active" | "archive"; + /** + * Reads only this project's threads. Fork sources in other projects still + * load, so visible item counts match the unscoped snapshot. + */ + readonly projectId?: ProjectId; /** * For background sweeps, not clients: skips settled threads before any of * their run, item or session rows are read. @@ -5132,8 +5138,7 @@ export const layer: Layer.Layer = const selectShellThreadRows = ( threadId?: ThreadId, - location?: "active" | "archive", - unsettledOnly = false, + { location, projectId, unsettledOnly = false }: ShellSnapshotOptions = {}, ) => sql` SELECT @@ -5309,6 +5314,8 @@ export const layer: Layer.Layer = LIMIT 1 ) AND blocked.status = 'failed' WHERE t.deleted_at IS NULL${threadId === undefined ? sql`` : sql` AND t.thread_id = ${threadId}`}${ + projectId === undefined ? sql`` : sql` AND t.project_id = ${projectId}` + }${ location === "active" ? sql` AND json_extract(t.payload_json, '$.archivedAt') IS NULL` : location === "archive" @@ -5732,11 +5739,7 @@ export const layer: Layer.Layer = const readShellSnapshotRows = (options: ShellSnapshotOptions | undefined) => Effect.gen(function* () { - const targetThreadRows = yield* selectShellThreadRows( - undefined, - options?.location, - options?.unsettledOnly ?? false, - ); + const targetThreadRows = yield* selectShellThreadRows(undefined, options); const targetThreadIds = new Set( targetThreadRows.map((row) => ThreadId.make(row.thread_id)), ); @@ -5992,6 +5995,12 @@ export const layerMemory: Layer.Layer = Layer.effect( const existing = (yield* Ref.get(replayState)).projections; const selectedThreadIds = [...existing.entries()] .filter(([, projection]) => { + if ( + options?.projectId !== undefined && + projection.thread.projectId !== options.projectId + ) { + return false; + } if ( options?.unsettledOnly && (projection.thread.settledAt !== null || diff --git a/apps/server/src/orchestration-v2/ThreadManagementService.ts b/apps/server/src/orchestration-v2/ThreadManagementService.ts index a12e6680ca7e..3da3ef12438d 100644 --- a/apps/server/src/orchestration-v2/ThreadManagementService.ts +++ b/apps/server/src/orchestration-v2/ThreadManagementService.ts @@ -563,7 +563,7 @@ const make = Effect.gen(function* () { ); const listProjectThreads: ThreadManagementServiceShape["listProjectThreads"] = (input) => - orchestrator.getShellSnapshot().pipe( + orchestrator.getShellSnapshot({ projectId: input.projectId, location: "active" }).pipe( Effect.mapError( (cause) => new ThreadManagementProjectThreadsListError({ @@ -573,7 +573,6 @@ const make = Effect.gen(function* () { ), Effect.map((snapshot) => snapshot.threads - .filter((thread) => thread.projectId === input.projectId) .filter( (thread) => input.includeSubagents || thread.lineage.relationshipToParent !== "subagent",