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
2 changes: 2 additions & 0 deletions apps/server/integration/transferBudgetV2.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
});
Expand Down
41 changes: 28 additions & 13 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<OrchestrationV2ThreadShellSnapshot, OrchestratorV2Error>;
readonly getShellSnapshot: (
options?: ShellSnapshotOptions,
) => Effect.Effect<OrchestrationV2ThreadShellSnapshot, OrchestratorV2Error>;
/** See `ProjectionStoreV2Shape.readShellSnapshot`. */
readonly readShellSnapshot: (
options?: ShellSnapshotOptions,
) => Effect.Effect<
Effect.Effect<OrchestrationV2ThreadShellSnapshot, OrchestratorV2Error>,
OrchestratorV2Error
>;
readonly getThreadShell: (
threadId: ThreadId,
) => Effect.Effect<OrchestrationV2ThreadShell | null, OrchestratorV2Error>;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -10982,6 +10990,13 @@ const layerUnavailable: Layer.Layer<OrchestratorV2> = 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({
Expand Down
75 changes: 75 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import {
MessageId,
type ModelSelection,
type OrchestrationV2ProviderThread,
type OrchestrationV2ThreadShellSnapshot,
NodeId,
ProjectId,
ProviderDriverKind,
Expand Down Expand Up @@ -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;
Expand Down
Loading
Loading