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
180 changes: 180 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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;
Expand Down
23 changes: 16 additions & 7 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import type {
OrchestrationV2ThreadShell,
OrchestrationV2ThreadProjection,
OrchestrationV2TurnItem,
ProjectId,
ProviderInstanceId,
ProviderSessionId,
ProviderThreadId,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -5132,8 +5138,7 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =

const selectShellThreadRows = (
threadId?: ThreadId,
location?: "active" | "archive",
unsettledOnly = false,
{ location, projectId, unsettledOnly = false }: ShellSnapshotOptions = {},
) =>
sql<ShellThreadRow>`
SELECT
Expand Down Expand Up @@ -5309,6 +5314,8 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
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"
Expand Down Expand Up @@ -5732,11 +5739,7 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =

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)),
);
Expand Down Expand Up @@ -5992,6 +5995,12 @@ export const layerMemory: Layer.Layer<ProjectionStoreV2> = 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 ||
Expand Down
3 changes: 1 addition & 2 deletions apps/server/src/orchestration-v2/ThreadManagementService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand All @@ -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",
Expand Down
Loading