diff --git a/apps/server/src/checkpointing/CheckpointDiffQuery.test.ts b/apps/server/src/checkpointing/CheckpointDiffQuery.test.ts index 6df7d7534130..ed72f0d34a33 100644 --- a/apps/server/src/checkpointing/CheckpointDiffQuery.test.ts +++ b/apps/server/src/checkpointing/CheckpointDiffQuery.test.ts @@ -77,6 +77,7 @@ describe("CheckpointDiffQuery.layer", () => { Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("CheckpointDiffQuery should not request the command read model"), getSnapshot: () => @@ -195,6 +196,7 @@ describe("CheckpointDiffQuery.layer", () => { Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("CheckpointDiffQuery should not request the command read model"), getSnapshot: () => @@ -288,6 +290,7 @@ describe("CheckpointDiffQuery.layer", () => { Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("CheckpointDiffQuery should not request the command read model"), getSnapshot: () => @@ -366,6 +369,7 @@ describe("CheckpointDiffQuery.layer", () => { Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("CheckpointDiffQuery should not request the command read model"), getSnapshot: () => @@ -429,6 +433,7 @@ describe("CheckpointDiffQuery.layer", () => { Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("CheckpointDiffQuery should not request the command read model"), getSnapshot: () => diff --git a/apps/server/src/cli/project.ts b/apps/server/src/cli/project.ts index cef4dc879073..003bdd81e08c 100644 --- a/apps/server/src/cli/project.ts +++ b/apps/server/src/cli/project.ts @@ -337,8 +337,8 @@ const dispatchLiveOrchestrationCommand = ( const getOfflineSnapshot = Effect.fn("getOfflineSnapshot")(function* () { const projectionSnapshotQuery = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; // Project commands only read the project list, so use the lightweight - // command read model instead of hydrating every thread body in the database. - return yield* projectionSnapshotQuery.getCommandReadModel(); + // metadata snapshot instead of hydrating every thread body in the database. + return yield* projectionSnapshotQuery.getMetadataSnapshot(); }); const tryResolveLiveProjectExecutionMode = Effect.fn("tryResolveLiveProjectExecutionMode")( diff --git a/apps/server/src/orchestration/CommandReadModel.ts b/apps/server/src/orchestration/CommandReadModel.ts new file mode 100644 index 000000000000..f7d8ef1dff32 --- /dev/null +++ b/apps/server/src/orchestration/CommandReadModel.ts @@ -0,0 +1,100 @@ +import { + OrchestrationCheckpointSummary, + OrchestrationMessage, + OrchestrationProposedPlan, + OrchestrationReadModel, + OrchestrationThread, + OrchestrationThreadActivity, +} from "@t3tools/contracts"; +import * as Duration from "effect/Duration"; +import * as Predicate from "effect/Predicate"; +import * as Schema from "effect/Schema"; +import * as Struct from "effect/Struct"; + +// This process-local model is for command decisions. Content belongs in the +// SQL projections; omitting it here makes accidental content reads a type error. +export const CommandMessage = Schema.Struct( + Struct.pick(OrchestrationMessage.fields, ["id", "role", "turnId", "createdAt", "updatedAt"]), +); +export type CommandMessage = typeof CommandMessage.Type; + +export const CommandCheckpoint = Schema.Struct( + Struct.omit(OrchestrationCheckpointSummary.fields, ["files"]), +); +export const CommandProposedPlan = Schema.Struct( + Struct.pick(OrchestrationProposedPlan.fields, ["id", "turnId", "createdAt", "updatedAt"]), +); +export const CommandActivity = Schema.Struct({ + ...Struct.pick(OrchestrationThreadActivity.fields, [ + "id", + "kind", + "turnId", + "sequence", + "createdAt", + ]), + requestId: Schema.NullOr(Schema.String), + responseMode: Schema.NullOr(Schema.Literal("message")), + staleFailure: Schema.Boolean, +}); +export type CommandActivity = typeof CommandActivity.Type; + +export const CommandThread = Schema.Struct({ + ...OrchestrationThread.fields, + messages: Schema.Array(CommandMessage), + activities: Schema.Array(CommandActivity), + checkpoints: Schema.Array(CommandCheckpoint), + proposedPlans: Schema.Array(CommandProposedPlan), +}); +export type CommandThread = typeof CommandThread.Type; +export const CommandReadModel = Schema.Struct({ + ...OrchestrationReadModel.fields, + threads: Schema.Array(CommandThread), +}); +export type CommandReadModel = typeof CommandReadModel.Type; + +export const MAX_COMMAND_MESSAGES = 2_000; +export const MAX_COMMAND_ACTIVITIES = 500; +export const MAX_COMMAND_CHECKPOINTS = 500; +export const MAX_COMMAND_PLANS = 200; + +/** Startup restores decision history only for live threads updated within this + * window of the newest thread activity. Older and archived threads start empty, + * as they did before this model existed, so boot memory follows recent work + * instead of every thread ever created. */ +export const COMMAND_HISTORY_RESTORE_WINDOW = Duration.days(7); + +/** Reduces activity bodies to the request facts used by the decider and retention. */ +export function toCommandActivity( + activity: Pick< + OrchestrationThreadActivity, + "id" | "kind" | "turnId" | "sequence" | "createdAt" | "payload" + >, +): CommandActivity { + const isRequestActivity = + activity.kind === "approval.requested" || + activity.kind === "approval.resolved" || + activity.kind === "user-input.requested" || + activity.kind === "user-input.resolved" || + activity.kind === "provider.approval.respond.failed" || + activity.kind === "provider.user-input.respond.failed"; + const payload = + isRequestActivity && Predicate.isObject(activity.payload) ? activity.payload : null; + const detail = typeof payload?.detail === "string" ? payload.detail.toLowerCase() : ""; + return { + id: activity.id, + kind: activity.kind, + turnId: activity.turnId, + ...(activity.sequence !== undefined ? { sequence: activity.sequence } : {}), + createdAt: activity.createdAt, + requestId: typeof payload?.requestId === "string" ? payload.requestId : null, + responseMode: payload?.responseMode === "message" ? "message" : null, + staleFailure: + detail.includes("stale pending approval request") || + detail.includes("unknown pending approval request") || + detail.includes("unknown pending permission request") || + detail.includes("stale pending user-input request") || + detail.includes("unknown pending user-input request") || + detail.includes("unknown pending user input request") || + detail.includes("unknown pending codex user input request"), + }; +} diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 9ccf78ca1744..d74c55ec3488 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -54,6 +54,7 @@ import { } from "../Services/ProjectionPipeline.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import { ServerConfig } from "../../config.ts"; +import { createEmptyReadModel, projectEvent } from "../projector.ts"; const asProjectId = (value: string): ProjectId => ProjectId.make(value); const asMessageId = (value: string): MessageId => MessageId.make(value); @@ -107,6 +108,7 @@ async function createOrchestrationSystem( return { engine, readModel: () => runtime.runPromise(snapshotQuery.getSnapshot()), + readCommandModel: () => runtime.runPromise(snapshotQuery.getCommandReadModel()), readThread: (threadId: ThreadId) => runtime.runPromise(snapshotQuery.getThreadDetailById(threadId)), run: (effect: Effect.Effect) => runtime.runPromise(effect), @@ -130,6 +132,212 @@ const hasMetricSnapshot = ( ); describe("OrchestrationEngine", () => { + it("keeps command decisions and compact history across a restart", async () => { + const directory = await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "t3-command-state-")); + const databasePath = NodePath.join(directory, "state.sqlite"); + let system = await createOrchestrationSystem(databasePath); + const projectId = ProjectId.make("command-project"); + const threadId = ThreadId.make("command-thread"); + const importedId = ThreadId.make("import:codex:command-state"); + const turnId = TurnId.make("command-turn"); + const messageId = MessageId.make("command-message"); + let commandCount = 0; + const commandId = () => CommandId.make(`command-state-${++commandCount}`); + const dispatch = (command: OrchestrationCommand) => system.run(system.engine.dispatch(command)); + const body = "large conversation content ".repeat(4_000); + try { + await dispatch({ + type: "project.create", + commandId: commandId(), + projectId, + title: "Command state", + workspaceRoot: "/tmp/command-state", + createdAt: now(), + }); + for (const id of [threadId, importedId]) { + await dispatch({ + type: "thread.create", + commandId: commandId(), + threadId: id, + projectId, + title: "Command state", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.4" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: now(), + }); + } + await dispatch({ + type: "thread.message.user.append", + commandId: commandId(), + threadId, + message: { messageId, text: body, attachments: [] }, + createdAt: now(), + }); + await dispatch({ + type: "thread.message.assistant.delta", + commandId: commandId(), + threadId, + messageId: MessageId.make("assistant:command"), + delta: body, + createdAt: now(), + }); + await dispatch({ + type: "thread.activity.append", + commandId: commandId(), + threadId, + createdAt: now(), + activity: { + id: EventId.make("tool-command"), + kind: "tool.completed", + tone: "tool", + summary: body, + turnId, + createdAt: now(), + payload: { output: body }, + }, + }); + await dispatch({ + type: "thread.turn.diff.complete", + commandId: commandId(), + threadId, + turnId, + checkpointRef: CheckpointRef.make("refs/t3/checkpoints/command"), + checkpointTurnCount: 1, + status: "ready", + files: [{ path: "src/app.ts", kind: "modified", additions: 3, deletions: 1 }], + createdAt: now(), + completedAt: now(), + }); + await dispatch({ + type: "thread.proposed-plan.upsert", + commandId: commandId(), + threadId, + createdAt: now(), + proposedPlan: { + id: "plan-command", + turnId, + planMarkdown: body, + implementedAt: null, + implementationThreadId: null, + createdAt: now(), + updatedAt: now(), + }, + }); + await dispatch({ + type: "thread.activity.append", + commandId: commandId(), + threadId, + createdAt: now(), + activity: { + id: EventId.make("approval-command"), + kind: "approval.requested", + tone: "approval", + summary: body, + turnId, + createdAt: now(), + payload: { requestId: "approval-command", detail: body }, + }, + }); + await dispatch({ + type: "thread.history.import", + commandId: commandId(), + threadId: importedId, + messages: [ + { + messageId: MessageId.make(`${importedId}:000000`), + role: "assistant", + text: body, + createdAt: now(), + }, + ], + }); + + // Rebuild the live command projection independently of the SQL bootstrap. + const events = await system.run(Stream.runCollect(system.engine.readEvents(0))); + const live = await system.run( + Effect.gen(function* () { + let model = createEmptyReadModel(now()); + for (const event of events) model = yield* projectEvent(model, event); + return model; + }), + ); + const histories = (model: typeof live) => + model.threads.map((thread) => ({ + id: thread.id, + messages: thread.messages, + activities: thread.activities, + checkpoints: thread.checkpoints, + proposedPlans: thread.proposedPlans, + })); + const before = histories(live).toSorted((a, b) => a.id.localeCompare(b.id)); + expect(JSON.stringify(before)).not.toContain("large conversation content"); + expect(live.threads[0]?.messages).toHaveLength(2); + expect(live.threads[0]?.checkpoints[0]).not.toHaveProperty("files"); + expect(live.threads[0]?.proposedPlans[0]).not.toHaveProperty("planMarkdown"); + expect( + (await system.readModel()).threads.find((thread) => thread.id === threadId)?.messages[0] + ?.text, + ).toBe(body); + + const rejectedCommands = (): OrchestrationCommand[] => [ + { + type: "thread.message.user.append", + commandId: commandId(), + threadId, + message: { messageId, text: "duplicate", attachments: [] }, + createdAt: now(), + }, + { type: "thread.settle", commandId: commandId(), threadId }, + { + type: "thread.history.import", + commandId: commandId(), + threadId: importedId, + messages: [ + { + messageId: MessageId.make(`${importedId}:000001`), + role: "assistant", + text: "duplicate", + createdAt: now(), + }, + ], + }, + { + type: "thread.turn.diff.complete", + commandId: commandId(), + threadId, + turnId, + checkpointRef: CheckpointRef.make("provider-diff:command"), + checkpointTurnCount: 1, + status: "missing", + files: [], + createdAt: now(), + completedAt: now(), + }, + ]; + const rejectionTags = async () => { + const tags: string[] = []; + for (const command of rejectedCommands()) { + const error = await system.run(system.engine.dispatch(command).pipe(Effect.flip)); + tags.push(error._tag); + } + return tags; + }; + const beforeRejections = await rejectionTags(); + await system.dispose(); + system = await createOrchestrationSystem(databasePath); + expect( + histories(await system.readCommandModel()).toSorted((a, b) => a.id.localeCompare(b.id)), + ).toEqual(before); + expect(await rejectionTags()).toEqual(beforeRejections); + } finally { + await system.dispose(); + await NodeFSP.rm(directory, { recursive: true, force: true }); + } + }); + it.each(["running", "stopped"] as const)( "sends async answers with a %s session and rejects old duplicate replies", async (status) => { @@ -240,6 +448,14 @@ describe("OrchestrationEngine", () => { } }; await appendWork("work", "2026-01-01T00:00:01.000Z"); + const commandActivities = (await system.readCommandModel()).threads[0]?.activities; + expect(commandActivities).toHaveLength(501); + expect( + commandActivities?.find((activity) => activity.id === "async-question"), + ).toMatchObject({ + requestId, + responseMode: "message", + }); const before = await system.readModel(); expect( before.threads[0]?.activities.some((activity) => activity.id === "async-question"), @@ -309,6 +525,11 @@ describe("OrchestrationEngine", () => { ), ).rejects.toThrow("This question has already been answered."); await appendWork("later-work", "2026-01-01T00:00:03.000Z"); + const reloadedActivities = (await system.readCommandModel()).threads[0]?.activities; + expect(reloadedActivities).toHaveLength(500); + expect(reloadedActivities?.some((activity) => activity.requestId === requestId)).toBe( + false, + ); const afterEviction = Option.getOrThrow(await system.readThread(threadId)); expect( afterEviction.activities.some((activity) => activity.kind === "user-input.resolved"), @@ -421,6 +642,7 @@ describe("OrchestrationEngine", () => { Layer.succeed(ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.succeed(commandReadModel), getSnapshot: () => Effect.sync(() => { diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts index 9136d080c1c3..05c9ad4049f3 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts @@ -1,7 +1,7 @@ +import type { CommandReadModel } from "../CommandReadModel.ts"; import type { OrchestrationClientOrigin, OrchestrationEvent, - OrchestrationReadModel, ProjectId, ThreadId, } from "@t3tools/contracts"; @@ -97,9 +97,9 @@ const makeOrchestrationEngine = Effect.gen(function* () { const eventPubSub = yield* PubSub.unbounded(); const projectEventsOntoReadModel = ( - baseReadModel: OrchestrationReadModel, + baseReadModel: CommandReadModel, events: ReadonlyArray, - ): Effect.Effect => + ): Effect.Effect => Effect.gen(function* () { let nextReadModel = baseReadModel; for (const event of events) { @@ -235,8 +235,8 @@ const makeOrchestrationEngine = Effect.gen(function* () { } } - // Command snapshots omit activities at startup and cap them while running. - // Read this request's durable state before deciding how to send the answer. + // The command model only keeps request metadata in a bounded window. + // Read the durable question body before deciding how to send the answer. const userInputActivity = envelope.command.type === "thread.user-input.respond" || envelope.command.type === "thread.user-input.dismiss" diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index cfba90297633..7239ef7acea8 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -18,6 +18,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as TestClock from "effect/testing/TestClock"; import * as Schema from "effect/Schema"; import * as SqlClient from "effect/unstable/sql/SqlClient"; @@ -2134,6 +2135,121 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { }), ); + it.effect("restores command history only for live threads active in the last week", () => + Effect.gen(function* () { + const snapshotQuery = yield* ProjectionSnapshotQuery; + const sql = yield* SqlClient.SqlClient; + yield* TestClock.setTime(Date.parse("2026-04-25T00:00:00.000Z")); + + yield* sql`DELETE FROM projection_projects`; + yield* sql`DELETE FROM projection_threads`; + yield* sql`DELETE FROM projection_thread_messages`; + yield* sql`DELETE FROM projection_thread_activities`; + yield* sql`DELETE FROM projection_turns`; + + yield* sql` + INSERT INTO projection_projects (project_id, title, workspace_root, scripts_json, created_at, updated_at) + VALUES ('project', 'Project', '/tmp/project', '[]', '2026-04-01T00:00:00.000Z', '2026-04-01T00:00:00.000Z') + `; + // "recent" is the newest thread; "idle" was last active 8 days before it. + yield* sql` + INSERT INTO projection_threads ( + thread_id, project_id, title, model_selection_json, runtime_mode, interaction_mode, + created_at, updated_at, archived_at, deleted_at + ) + VALUES + ('recent', 'project', 'Recent', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-04-20T00:00:00.000Z', NULL, NULL), + ('week-old', 'project', 'Week old', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-04-13T00:00:00.000Z', NULL, NULL), + ('idle', 'project', 'Idle', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-04-12T00:00:00.000Z', NULL, NULL), + ('archived', 'project', 'Archived', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-04-19T00:00:00.000Z', '2026-04-19T00:00:00.000Z', NULL), + ('deleted', 'project', 'Deleted', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-04-19T00:00:00.000Z', NULL, '2026-04-19T00:00:00.000Z') + `; + for (const threadId of ["recent", "week-old", "idle", "archived", "deleted"]) { + yield* sql` + INSERT INTO projection_thread_messages (message_id, thread_id, turn_id, role, text, is_streaming, created_at, updated_at) + VALUES (${`${threadId}-message`}, ${threadId}, 'turn', 'user', 'hello', 0, '2026-04-01T00:00:00.000Z', '2026-04-01T00:00:00.000Z') + `; + yield* sql` + INSERT INTO projection_thread_activities (activity_id, thread_id, turn_id, tone, kind, summary, payload_json, created_at, sequence) + VALUES (${`${threadId}-activity`}, ${threadId}, 'turn', 'approval', 'approval.requested', 'Approve', + json_object('requestId', ${`${threadId}-request`}), '2026-04-01T00:00:00.000Z', 1) + `; + yield* sql` + INSERT INTO projection_turns (thread_id, turn_id, state, requested_at, completed_at, + checkpoint_turn_count, checkpoint_ref, checkpoint_status, checkpoint_files_json) + VALUES (${threadId}, 'turn', 'completed', '2026-04-01T00:00:00.000Z', '2026-04-01T00:00:01.000Z', + 1, 'ref', 'ready', '[]') + `; + } + + const commandReadModel = yield* snapshotQuery.getCommandReadModel(); + const restored = Object.fromEntries( + commandReadModel.threads.map((thread) => [ + thread.id, + [thread.messages.length, thread.activities.length, thread.checkpoints.length], + ]), + ); + assert.deepStrictEqual(restored, { + recent: [1, 1, 1], + "week-old": [1, 1, 1], + idle: [0, 0, 0], + archived: [0, 0, 0], + deleted: [0, 0, 0], + }); + }), + ); + + it.effect("restores recent history when a thread has a malformed or future timestamp", () => + Effect.gen(function* () { + const snapshotQuery = yield* ProjectionSnapshotQuery; + const sql = yield* SqlClient.SqlClient; + yield* TestClock.setTime(Date.parse("2026-04-25T00:00:00.000Z")); + + yield* sql`DELETE FROM projection_projects`; + yield* sql`DELETE FROM projection_threads`; + yield* sql`DELETE FROM projection_thread_messages`; + yield* sql` + INSERT INTO projection_projects (project_id, title, workspace_root, scripts_json, created_at, updated_at) + VALUES ('project', 'Project', '/tmp/project', '[]', '2026-04-01T00:00:00.000Z', '2026-04-01T00:00:00.000Z') + `; + yield* sql` + INSERT INTO projection_threads ( + thread_id, project_id, title, model_selection_json, runtime_mode, interaction_mode, + created_at, updated_at + ) + VALUES + ('recent', 'project', 'Recent', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-04-20T00:00:00.000Z'), + ('odd', 'project', 'Odd', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', + '2026-04-01T00:00:00.000Z', '2026-13-01T00:00:00.000Z') + `; + yield* sql` + INSERT INTO projection_thread_messages (message_id, thread_id, turn_id, role, text, is_streaming, created_at, updated_at) + VALUES ('recent-message', 'recent', 'turn', 'user', 'hello', 0, '2026-04-01T00:00:00.000Z', '2026-04-01T00:00:00.000Z') + `; + const restoredMessages = () => + snapshotQuery + .getCommandReadModel() + .pipe( + Effect.map( + (model) => model.threads.find((thread) => thread.id === "recent")?.messages.length, + ), + ); + + // A malformed timestamp sorts newest but cannot fail startup. + assert.equal(yield* restoredMessages(), 1); + + // A timestamp a year ahead cannot move the window past every thread. + yield* sql`UPDATE projection_threads SET updated_at = '2027-04-20T00:00:00.000Z' WHERE thread_id = 'odd'`; + assert.equal(yield* restoredMessages(), 1); + }), + ); + it.effect("searches active user messages and canonical assistant outputs", () => Effect.gen(function* () { const snapshotQuery = yield* ProjectionSnapshotQuery; diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index 07d30e47e27b..f45450b6a4f1 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -38,6 +38,7 @@ import { } from "@t3tools/contracts"; import { legacyLinkedPullRequestOf } from "@t3tools/shared/threadPullRequests"; import * as Arr from "effect/Array"; +import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -47,6 +48,20 @@ import * as Struct from "effect/Struct"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as SqlSchema from "effect/unstable/sql/SqlSchema"; +import { + CommandActivity, + CommandCheckpoint, + CommandMessage, + COMMAND_HISTORY_RESTORE_WINDOW, + type CommandReadModel, + MAX_COMMAND_MESSAGES, + MAX_COMMAND_ACTIVITIES, + MAX_COMMAND_CHECKPOINTS, + MAX_COMMAND_PLANS, + toCommandActivity, +} from "../CommandReadModel.ts"; +import { WORKTREE_SETUP_ACTIVITY_KIND } from "@t3tools/contracts"; + import { isPersistenceError, toPersistenceDecodeError, @@ -2424,63 +2439,158 @@ pending_approval_requests AS ( }), ); - const getCommandReadModel: ProjectionSnapshotQueryShape["getCommandReadModel"] = () => + // Select metadata in SQLite so startup never materializes conversation bodies. + // Each query reads only threads updated since `cutoff` (see + // COMMAND_HISTORY_RESTORE_WINDOW). + const listCommandMessages = SqlSchema.findAll({ + Request: Schema.String, + Result: Schema.Struct({ ...CommandMessage.fields, threadId: ThreadId }), + execute: (cutoff) => sql` + SELECT message_id AS id, thread_id AS "threadId", turn_id AS "turnId", + role, created_at AS "createdAt", updated_at AS "updatedAt" + FROM ( + SELECT rowid AS insertion_order, message_id, thread_id, turn_id, role, created_at, updated_at, + ROW_NUMBER() OVER (PARTITION BY thread_id ORDER BY rowid DESC) AS rank + FROM projection_thread_messages + WHERE thread_id IN ( + SELECT thread_id FROM projection_threads + WHERE deleted_at IS NULL AND archived_at IS NULL AND updated_at >= ${cutoff} + ) + ) + WHERE rank <= ${MAX_COMMAND_MESSAGES} + ORDER BY thread_id, insertion_order + `, + }); + const listCommandCheckpoints = SqlSchema.findAll({ + Request: Schema.String, + Result: Schema.Struct({ ...CommandCheckpoint.fields, threadId: ThreadId }), + execute: (cutoff) => sql` + SELECT thread_id AS "threadId", turn_id AS "turnId", checkpoint_turn_count AS "checkpointTurnCount", + checkpoint_ref AS "checkpointRef", checkpoint_status AS status, + assistant_message_id AS "assistantMessageId", completed_at AS "completedAt" + FROM ( + SELECT thread_id, turn_id, checkpoint_turn_count, checkpoint_ref, checkpoint_status, + assistant_message_id, completed_at, + ROW_NUMBER() OVER (PARTITION BY thread_id ORDER BY checkpoint_turn_count DESC) AS rank + FROM projection_turns + WHERE checkpoint_turn_count IS NOT NULL AND thread_id IN ( + SELECT thread_id FROM projection_threads + WHERE deleted_at IS NULL AND archived_at IS NULL AND updated_at >= ${cutoff} + ) + ) + WHERE rank <= ${MAX_COMMAND_CHECKPOINTS} + ORDER BY thread_id, checkpoint_turn_count + `, + }); + const listCommandActivities = SqlSchema.findAll({ + Request: Schema.String, + Result: Schema.Struct({ + ...Struct.pick(CommandActivity.fields, ["id", "kind", "turnId", "createdAt"]), + threadId: ThreadId, + sequence: Schema.NullOr(NonNegativeInt), + payload: Schema.fromJsonString(Schema.Unknown), + }), + execute: (cutoff) => sql` + WITH recent AS ( + SELECT rowid AS activity_rowid, thread_id, activity_id, kind, sequence, created_at, + ROW_NUMBER() OVER ( + PARTITION BY thread_id ORDER BY sequence DESC, created_at DESC, activity_id DESC + ) AS rank + FROM projection_thread_activities + WHERE thread_id IN ( + SELECT thread_id FROM projection_threads + WHERE deleted_at IS NULL AND archived_at IS NULL AND updated_at >= ${cutoff} + ) + ), async_lifecycle AS ( + SELECT activity_id, kind, + ROW_NUMBER() OVER ( + PARTITION BY thread_id, json_extract(payload_json, '$.requestId') + ORDER BY sequence DESC, created_at DESC, activity_id DESC + ) AS rank + FROM projection_thread_activities + WHERE (kind = 'user-input.resolved' OR + (kind = 'user-input.requested' AND json_extract(payload_json, '$.responseMode') = 'message')) + AND json_type(payload_json, '$.requestId') = 'text' + AND thread_id IN ( + SELECT thread_id FROM projection_threads + WHERE deleted_at IS NULL AND archived_at IS NULL AND updated_at >= ${cutoff} + ) + ) + SELECT a.activity_id AS id, a.thread_id AS "threadId", a.turn_id AS "turnId", + a.kind, a.sequence, a.created_at AS "createdAt", + CASE WHEN a.kind IN ('approval.requested', 'approval.resolved', + 'user-input.requested', 'user-input.resolved', + 'provider.approval.respond.failed', 'provider.user-input.respond.failed') + THEN json_object('requestId', json_extract(a.payload_json, '$.requestId'), + 'responseMode', json_extract(a.payload_json, '$.responseMode'), + 'detail', json_extract(a.payload_json, '$.detail')) + ELSE '{}' END AS payload + FROM recent r JOIN projection_thread_activities a ON a.rowid = r.activity_rowid + WHERE r.rank <= ${MAX_COMMAND_ACTIVITIES} OR a.kind = ${WORKTREE_SETUP_ACTIVITY_KIND} + OR a.activity_id IN ( + SELECT activity_id FROM async_lifecycle WHERE rank = 1 AND kind = 'user-input.requested' + ) + ORDER BY a.thread_id, a.sequence, a.created_at, a.activity_id + `, + }); + + const getMetadataSnapshot: ProjectionSnapshotQueryShape["getMetadataSnapshot"] = () => sql .withTransaction( Effect.all([ listProjectRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listProjects:query", - "ProjectionSnapshotQuery.getCommandReadModel:listProjects:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listProjects:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listProjects:decodeRows", ), ), ), listThreadRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listThreads:query", - "ProjectionSnapshotQuery.getCommandReadModel:listThreads:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreads:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreads:decodeRows", ), ), ), listThreadProposedPlanRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listThreadProposedPlans:query", - "ProjectionSnapshotQuery.getCommandReadModel:listThreadProposedPlans:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreadProposedPlans:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreadProposedPlans:decodeRows", ), ), ), listThreadPullRequestRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listThreadPullRequests:query", - "ProjectionSnapshotQuery.getCommandReadModel:listThreadPullRequests:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreadPullRequests:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreadPullRequests:decodeRows", ), ), ), listThreadSessionRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listThreadSessions:query", - "ProjectionSnapshotQuery.getCommandReadModel:listThreadSessions:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreadSessions:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listThreadSessions:decodeRows", ), ), ), listLatestTurnRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listLatestTurns:query", - "ProjectionSnapshotQuery.getCommandReadModel:listLatestTurns:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listLatestTurns:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listLatestTurns:decodeRows", ), ), ), listProjectionStateRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getCommandReadModel:listProjectionState:query", - "ProjectionSnapshotQuery.getCommandReadModel:listProjectionState:decodeRows", + "ProjectionSnapshotQuery.getMetadataSnapshot:listProjectionState:query", + "ProjectionSnapshotQuery.getMetadataSnapshot:listProjectionState:decodeRows", ), ), ), @@ -2661,8 +2771,84 @@ pending_approval_requests AS ( if (isPersistenceError(error)) { return error; } - return toPersistenceSqlError("ProjectionSnapshotQuery.getCommandReadModel:query")(error); + return toPersistenceSqlError("ProjectionSnapshotQuery.getMetadataSnapshot:query")(error); + }), + ); + + const getCommandReadModel: ProjectionSnapshotQueryShape["getCommandReadModel"] = () => + sql + .withTransaction( + Effect.gen(function* () { + const metadata = yield* getMetadataSnapshot(); + // Measured from the newest thread rather than the wall clock, so a + // server restarted after a long idle stretch still restores its latest + // work. Clamped to now, so a malformed or future timestamp can neither + // fail startup nor move the window past every thread. + const newest = metadata.threads.reduce( + (latest, thread) => (thread.updatedAt > latest ? thread.updatedAt : latest), + "", + ); + const now = yield* DateTime.now; + const anchor = DateTime.make(newest).pipe( + Option.map((latest) => DateTime.min(latest, now)), + Option.getOrElse(() => now), + ); + const cutoff = DateTime.formatIso( + DateTime.subtractDuration(anchor, COMMAND_HISTORY_RESTORE_WINDOW), + ); + const [messages, activities, checkpoints] = yield* Effect.all([ + listCommandMessages(cutoff), + listCommandActivities(cutoff), + listCommandCheckpoints(cutoff), + ]); + const messagesByThread = new Map>(); + const activitiesByThread = new Map>(); + const checkpointsByThread = new Map>(); + for (const { threadId, ...message } of messages) { + const entries = messagesByThread.get(threadId) ?? []; + entries.push(message); + messagesByThread.set(threadId, entries); + } + for (const { threadId, sequence, ...activity } of activities) { + const entries = activitiesByThread.get(threadId) ?? []; + entries.push( + toCommandActivity({ ...activity, ...(sequence !== null ? { sequence } : {}) }), + ); + activitiesByThread.set(threadId, entries); + } + for (const { threadId, ...checkpoint } of checkpoints) { + const entries = checkpointsByThread.get(threadId) ?? []; + entries.push(checkpoint); + checkpointsByThread.set(threadId, entries); + } + return { + ...metadata, + threads: metadata.threads.map((thread) => ({ + ...thread, + messages: messagesByThread.get(thread.id) ?? [], + activities: activitiesByThread.get(thread.id) ?? [], + checkpoints: checkpointsByThread.get(thread.id) ?? [], + proposedPlans: thread.proposedPlans + .slice(-MAX_COMMAND_PLANS) + .map(({ id, turnId, createdAt, updatedAt }) => ({ + id, + turnId, + createdAt, + updatedAt, + })), + })), + } satisfies CommandReadModel; }), + ) + .pipe( + Effect.mapError((error) => + isPersistenceError(error) + ? error + : toPersistenceSqlOrDecodeError( + "ProjectionSnapshotQuery.getCommandReadModel:query", + "ProjectionSnapshotQuery.getCommandReadModel:decodeRows", + )(error), + ), ); const getShellSnapshot: ProjectionSnapshotQueryShape["getShellSnapshot"] = (options) => { @@ -3838,6 +4024,7 @@ pending_approval_requests AS ( return { getCommandReadModel, + getMetadataSnapshot, getUserInputActivity, listActivitiesByKind, getSnapshot, diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 77022a8518d7..6b7f80ed18ec 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -1119,7 +1119,7 @@ const make = Effect.gen(function* () { }); }); const findPendingThreadTitles = Effect.fn("findPendingThreadTitles")(function* () { - const readModel = yield* projectionSnapshotQuery.getCommandReadModel(); + const readModel = yield* projectionSnapshotQuery.getMetadataSnapshot(); return { interruptedRegenerations: readModel.threads.flatMap((thread) => { const requestId = thread.titleRegeneration?.requestId; diff --git a/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts index 08da6e00940c..ea8c4b2b0aa9 100644 --- a/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts @@ -1,3 +1,4 @@ +import type { CommandReadModel } from "../CommandReadModel.ts"; /** * ProjectionSnapshotQuery - Read-model snapshot query service interface. * @@ -99,10 +100,18 @@ export interface ProjectionSnapshotQueryShape { ) => Effect.Effect, ProjectionRepositoryError>; /** - * Read the lightweight command snapshot used to bootstrap the in-memory - * orchestration engine without hydrating message/activity/checkpoint bodies. + * Read the compact command snapshot used to bootstrap the in-memory + * orchestration engine. It never hydrates message, activity or checkpoint + * bodies, and restores decision history only for live threads active within + * COMMAND_HISTORY_RESTORE_WINDOW of the newest thread. */ - readonly getCommandReadModel: () => Effect.Effect< + readonly getCommandReadModel: () => Effect.Effect; + + /** Legacy metadata response used by the project CLI and startup maintenance. + * Message, activity and checkpoint bodies are intentionally omitted. Never + * use this response to decide commands; use getCommandReadModel instead. + */ + readonly getMetadataSnapshot: () => Effect.Effect< OrchestrationReadModel, ProjectionRepositoryError >; diff --git a/apps/server/src/orchestration/commandInvariants.test.ts b/apps/server/src/orchestration/commandInvariants.test.ts index a3db108f2020..8003fba3f258 100644 --- a/apps/server/src/orchestration/commandInvariants.test.ts +++ b/apps/server/src/orchestration/commandInvariants.test.ts @@ -1,3 +1,4 @@ +import type { CommandReadModel } from "./CommandReadModel.ts"; import { describe, expect, it } from "vite-plus/test"; import { MessageId, @@ -6,7 +7,6 @@ import { ProjectId, ThreadId, type OrchestrationCommand, - type OrchestrationReadModel, ProviderInstanceId, } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; @@ -15,7 +15,7 @@ import { listThreadsByProjectId, requireThread, requireThreadAbsent } from "./co const now = "2026-01-01T00:00:00.000Z"; -const readModel: OrchestrationReadModel = { +const readModel: CommandReadModel = { snapshotSequence: 2, updatedAt: now, projects: [ @@ -198,7 +198,7 @@ describe("commandInvariants", () => { it("lets a draft retry re-create a thread id after its first attempt was deleted", async () => { const threadId = ThreadId.make("thread-1"); const firstAttempt = readModel.threads.find((thread) => thread.id === threadId)!; - const afterRollback: OrchestrationReadModel = { + const afterRollback: CommandReadModel = { ...readModel, threads: readModel.threads.map((thread) => thread.id === threadId ? { ...thread, deletedAt: now, updatedAt: now } : thread, diff --git a/apps/server/src/orchestration/commandInvariants.ts b/apps/server/src/orchestration/commandInvariants.ts index 244a6a9e9cab..e99a25d1bbe7 100644 --- a/apps/server/src/orchestration/commandInvariants.ts +++ b/apps/server/src/orchestration/commandInvariants.ts @@ -1,8 +1,7 @@ +import type { CommandReadModel, CommandThread } from "./CommandReadModel.ts"; import type { OrchestrationCommand, OrchestrationProject, - OrchestrationReadModel, - OrchestrationThread, ProjectId, ThreadId, } from "@t3tools/contracts"; @@ -19,28 +18,28 @@ function invariantError(commandType: string, detail: string): OrchestrationComma } function findThreadById( - readModel: OrchestrationReadModel, + readModel: CommandReadModel, threadId: ThreadId, -): OrchestrationThread | undefined { +): CommandThread | undefined { return readModel.threads.find((thread) => thread.id === threadId); } function findProjectById( - readModel: OrchestrationReadModel, + readModel: CommandReadModel, projectId: ProjectId, ): OrchestrationProject | undefined { return readModel.projects.find((project) => project.id === projectId); } export function listThreadsByProjectId( - readModel: OrchestrationReadModel, + readModel: CommandReadModel, projectId: ProjectId, -): ReadonlyArray { +): ReadonlyArray { return readModel.threads.filter((thread) => thread.projectId === projectId); } export function requireProject(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly projectId: ProjectId; }): Effect.Effect { @@ -57,7 +56,7 @@ export function requireProject(input: { } export function requireProjectAbsent(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly projectId: ProjectId; }): Effect.Effect { @@ -73,7 +72,7 @@ export function requireProjectAbsent(input: { } export function requireActiveProjectWorkspaceRootAbsent(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly workspaceRoot: string; readonly exceptProjectId?: ProjectId; @@ -97,10 +96,10 @@ export function requireActiveProjectWorkspaceRootAbsent(input: { } export function requireThread(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly threadId: ThreadId; -}): Effect.Effect { +}): Effect.Effect { const thread = findThreadById(input.readModel, input.threadId); if (thread) { return Effect.succeed(thread); @@ -114,10 +113,10 @@ export function requireThread(input: { } export function requireThreadArchived(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly threadId: ThreadId; -}): Effect.Effect { +}): Effect.Effect { return requireThread(input).pipe( Effect.filterOrFail( (thread) => thread.archivedAt !== null, @@ -131,10 +130,10 @@ export function requireThreadArchived(input: { } export function requireThreadNotArchived(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly threadId: ThreadId; -}): Effect.Effect { +}): Effect.Effect { return requireThread(input).pipe( Effect.filterOrFail( (thread) => thread.archivedAt === null, @@ -148,7 +147,7 @@ export function requireThreadNotArchived(input: { } export function requireThreadAbsent(input: { - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly command: OrchestrationCommand; readonly threadId: ThreadId; }): Effect.Effect { diff --git a/apps/server/src/orchestration/decider.active-order.test.ts b/apps/server/src/orchestration/decider.active-order.test.ts index 39bccba11460..3143676f0229 100644 --- a/apps/server/src/orchestration/decider.active-order.test.ts +++ b/apps/server/src/orchestration/decider.active-order.test.ts @@ -1,11 +1,10 @@ +import type { CommandReadModel, CommandThread } from "./CommandReadModel.ts"; import { CommandId, ProjectId, ProviderInstanceId, ThreadId, type OrchestrationCommand, - type OrchestrationReadModel, - type OrchestrationThread, } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; @@ -21,7 +20,7 @@ const SNOOZED_AT = "1969-12-31T00:00:00.000Z"; const FUTURE_WAKE = "1970-01-02T00:00:00.000Z"; const THREAD_ID = ThreadId.make("thread-1"); -function makeReadModel(overrides: Partial = {}): OrchestrationReadModel { +function makeReadModel(overrides: Partial = {}): CommandReadModel { return { snapshotSequence: 0, projects: [], @@ -104,7 +103,7 @@ it.layer(NodeServices.layer)("active thread ordering", (it) => { ["deleted", { deletedAt: NOW }], ["pinned", { pinnedAt: NOW }], ["settled", { settledOverride: "settled", settledAt: NOW }], - ] satisfies ReadonlyArray]>) { + ] satisfies ReadonlyArray]>) { it.effect(`rejects reordering a ${label} thread`, () => Effect.gen(function* () { const error = yield* decideOrchestrationCommand({ diff --git a/apps/server/src/orchestration/decider.autoSettleSet.test.ts b/apps/server/src/orchestration/decider.autoSettleSet.test.ts index 99657069cc74..357fb4620b35 100644 --- a/apps/server/src/orchestration/decider.autoSettleSet.test.ts +++ b/apps/server/src/orchestration/decider.autoSettleSet.test.ts @@ -1,10 +1,5 @@ -import { - CommandId, - ProjectId, - ProviderInstanceId, - ThreadId, - type OrchestrationReadModel, -} from "@t3tools/contracts"; +import type { CommandReadModel } from "./CommandReadModel.ts"; +import { CommandId, ProjectId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; @@ -17,7 +12,7 @@ const DISABLED_AT = "2025-12-30T00:00:00.000Z"; function makeReadModel(input: { readonly autoSettleDisabledAt?: string | null; readonly settledOverride?: "settled" | "active" | null; -}): OrchestrationReadModel { +}): CommandReadModel { return { snapshotSequence: 0, projects: [], diff --git a/apps/server/src/orchestration/decider.import.test.ts b/apps/server/src/orchestration/decider.import.test.ts index c809c733800a..6fc08a18d43b 100644 --- a/apps/server/src/orchestration/decider.import.test.ts +++ b/apps/server/src/orchestration/decider.import.test.ts @@ -167,9 +167,9 @@ it.layer(NodeServices.layer)("thread history import", (it) => { metadata: {}, payload: { threadId, turnCount: 0 }, }); - expect(projected.threads[0]?.messages.map((message) => message.text)).toEqual([ - "Fix the bug", - "Fixed", + expect(projected.threads[0]?.messages.map((message) => message.role)).toEqual([ + "user", + "assistant", ]); }), ); diff --git a/apps/server/src/orchestration/decider.pinned.test.ts b/apps/server/src/orchestration/decider.pinned.test.ts index 7d6f55d8f622..149c46a56b40 100644 --- a/apps/server/src/orchestration/decider.pinned.test.ts +++ b/apps/server/src/orchestration/decider.pinned.test.ts @@ -1,10 +1,5 @@ -import { - CommandId, - ProjectId, - ProviderInstanceId, - ThreadId, - type OrchestrationReadModel, -} from "@t3tools/contracts"; +import type { CommandReadModel } from "./CommandReadModel.ts"; +import { CommandId, ProjectId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; @@ -22,7 +17,7 @@ function makeReadModel(input: { readonly settledAt?: string | null; readonly snoozedUntil?: string | null; readonly snoozedAt?: string | null; -}): OrchestrationReadModel { +}): CommandReadModel { return { snapshotSequence: 0, projects: [], diff --git a/apps/server/src/orchestration/decider.pullRequests.test.ts b/apps/server/src/orchestration/decider.pullRequests.test.ts index 7a961af12d5b..7ea6eada7cd1 100644 --- a/apps/server/src/orchestration/decider.pullRequests.test.ts +++ b/apps/server/src/orchestration/decider.pullRequests.test.ts @@ -1,3 +1,4 @@ +import type { CommandReadModel } from "./CommandReadModel.ts"; import { CommandId, ProjectId, @@ -5,7 +6,6 @@ import { ThreadId, OrchestrationEvent, OrchestrationCommand, - type OrchestrationReadModel, type ThreadPullRequestLink, type ThreadPullRequestSnapshot, } from "@t3tools/contracts"; @@ -50,7 +50,7 @@ function makeLink(overrides: Partial = {}): ThreadPullReq }; } -function makeReadModel(pullRequests: ReadonlyArray): OrchestrationReadModel { +function makeReadModel(pullRequests: ReadonlyArray): CommandReadModel { return { snapshotSequence: 0, projects: [ diff --git a/apps/server/src/orchestration/decider.questionAttachments.test.ts b/apps/server/src/orchestration/decider.questionAttachments.test.ts index 53e87d7593cc..01db219e01aa 100644 --- a/apps/server/src/orchestration/decider.questionAttachments.test.ts +++ b/apps/server/src/orchestration/decider.questionAttachments.test.ts @@ -1,3 +1,4 @@ +import type { CommandReadModel } from "./CommandReadModel.ts"; import { ApprovalRequestId, EventId, @@ -5,7 +6,6 @@ import { ProjectId, ProviderInstanceId, ThreadId, - type OrchestrationReadModel, } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; @@ -15,7 +15,7 @@ import { decideOrchestrationCommand } from "./decider.ts"; const UPDATED_AT = "2026-01-01T00:00:00.000Z"; -const readModel: OrchestrationReadModel = { +const readModel: CommandReadModel = { snapshotSequence: 0, projects: [], threads: [ diff --git a/apps/server/src/orchestration/decider.settled.test.ts b/apps/server/src/orchestration/decider.settled.test.ts index bcb2408c21e1..37636996fa4d 100644 --- a/apps/server/src/orchestration/decider.settled.test.ts +++ b/apps/server/src/orchestration/decider.settled.test.ts @@ -1,3 +1,9 @@ +import type { OrchestrationThread } from "@t3tools/contracts"; +import { + toCommandActivity, + type CommandReadModel, + type CommandThread, +} from "./CommandReadModel.ts"; import { CommandId, EventId, @@ -6,9 +12,7 @@ import { ProviderInstanceId, ThreadId, type OrchestrationEvent, - type OrchestrationReadModel, type OrchestrationSession, - type OrchestrationThread, } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; @@ -23,7 +27,7 @@ const SETTLE_BLOCKED_MESSAGE = "This thread still needs attention. Resolve or interrupt it first, then try again."; function makeReadModel( - settledOverride: OrchestrationThread["settledOverride"], + settledOverride: CommandThread["settledOverride"], archivedAt: string | null = null, session: OrchestrationSession | null = null, activities: OrchestrationThread["activities"] = [], @@ -33,7 +37,7 @@ function makeReadModel( readonly snoozedUntil?: string | null; readonly snoozedAt?: string | null; } = {}, -): OrchestrationReadModel { +): CommandReadModel { return { snapshotSequence: 0, projects: [], @@ -60,7 +64,7 @@ function makeReadModel( deletedAt: null, messages, proposedPlans: [], - activities, + activities: activities.map(toCommandActivity), checkpoints: [], session, }, diff --git a/apps/server/src/orchestration/decider.snoozed.test.ts b/apps/server/src/orchestration/decider.snoozed.test.ts index 505bb79df034..0146dcc47c2c 100644 --- a/apps/server/src/orchestration/decider.snoozed.test.ts +++ b/apps/server/src/orchestration/decider.snoozed.test.ts @@ -1,3 +1,5 @@ +import type { OrchestrationThread } from "@t3tools/contracts"; +import { toCommandActivity, type CommandReadModel } from "./CommandReadModel.ts"; import { CommandId, EventId, @@ -5,8 +7,6 @@ import { ProjectId, ProviderInstanceId, ThreadId, - type OrchestrationReadModel, - type OrchestrationThread, } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; @@ -27,7 +27,7 @@ function makeReadModel(input: { readonly archivedAt?: string | null; readonly activities?: OrchestrationThread["activities"]; readonly messages?: OrchestrationThread["messages"]; -}): OrchestrationReadModel { +}): CommandReadModel { return { snapshotSequence: 0, projects: [], @@ -53,7 +53,7 @@ function makeReadModel(input: { deletedAt: null, messages: input.messages ?? [], proposedPlans: [], - activities: input.activities ?? [], + activities: (input.activities ?? []).map(toCommandActivity), checkpoints: [], session: null, }, diff --git a/apps/server/src/orchestration/decider.titleRegeneration.test.ts b/apps/server/src/orchestration/decider.titleRegeneration.test.ts index e93580c15f3b..c7bb658d2418 100644 --- a/apps/server/src/orchestration/decider.titleRegeneration.test.ts +++ b/apps/server/src/orchestration/decider.titleRegeneration.test.ts @@ -1,10 +1,5 @@ -import { - CommandId, - ProjectId, - ProviderInstanceId, - ThreadId, - type OrchestrationReadModel, -} from "@t3tools/contracts"; +import type { CommandReadModel } from "./CommandReadModel.ts"; +import { CommandId, ProjectId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; @@ -13,7 +8,7 @@ import { decideOrchestrationCommand } from "./decider.ts"; const UPDATED_AT = "2026-01-01T00:00:00.000Z"; -const readModel: OrchestrationReadModel = { +const readModel: CommandReadModel = { snapshotSequence: 0, projects: [], threads: [ diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 21df39a162e6..1a0be41770ef 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1,3 +1,4 @@ +import type { CommandReadModel, CommandThread } from "./CommandReadModel.ts"; import { EventId, MAX_SCRIPT_ID_LENGTH, @@ -8,8 +9,6 @@ import { isImportedAgentSessionMessageId, type OrchestrationCommand, type OrchestrationEvent, - type OrchestrationReadModel, - type OrchestrationThread, type ThreadPullRequestKey, type ThreadPullRequestLink, type OrchestrationThreadActivity, @@ -64,32 +63,12 @@ const threadPullRequestLinksEqual = Schema.toEquivalence(Schema.NullOr(ThreadLin * resolved activities always clear, respond.failed clears only when the * failure detail marks the request stale/unknown — or settle would be * rejected on threads whose shell flags read as clear. + * Activities retain the most recent 500 plus pending async questions. */ -function isStaleRequestFailureDetail(payload: Record | null): boolean { - const detail = typeof payload?.detail === "string" ? payload.detail.toLowerCase() : null; - if (detail === null) return false; - return ( - detail.includes("stale pending approval request") || - detail.includes("unknown pending approval request") || - detail.includes("unknown pending permission request") || - detail.includes("stale pending user-input request") || - detail.includes("unknown pending user-input request") || - detail.includes("unknown pending user input request") || - detail.includes("unknown pending codex user input request") - ); -} - -// Scans the read model's activities, which the projector caps at the most -// recent 500 plus pending async questions. Async questions remain actionable -// while the agent works, so they must not expire with the activity window. -function openRequests(thread: Pick) { - const requests = new Map(); +function openRequests(thread: Pick) { + const requests = new Map(); for (const activity of thread.activities) { - const payload = - typeof activity.payload === "object" && activity.payload !== null - ? (activity.payload as Record) - : null; - const requestId = typeof payload?.requestId === "string" ? payload.requestId : null; + const { requestId } = activity; if (requestId === null) continue; if (activity.kind === "approval.requested" || activity.kind === "user-input.requested") { requests.set(requestId, activity); @@ -98,7 +77,7 @@ function openRequests(thread: Pick) { } else if ( (activity.kind === "provider.approval.respond.failed" || activity.kind === "provider.user-input.respond.failed") && - isStaleRequestFailureDetail(payload) + activity.staleFailure ) { requests.delete(requestId); } @@ -108,7 +87,7 @@ function openRequests(thread: Pick) { /** Apply the shared shell-level rule to the detailed command read model. */ function hasQueuedTurnStartForThread( - thread: Pick, + thread: Pick, now: string, ): boolean { let latestUserMessageAt: string | null = null; @@ -132,7 +111,7 @@ function hasQueuedTurnStartForThread( } function findPullRequestLink( - thread: Pick, + thread: Pick, key: ThreadPullRequestKey, ): ThreadPullRequestLink | undefined { return thread.pullRequests.find((link) => threadPullRequestKeysEqual(link, key)); @@ -179,7 +158,7 @@ const decideCommandSequence = Effect.fn("decideCommandSequence")(function* ({ readModel, }: { readonly commands: ReadonlyArray; - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; }): Effect.fn.Return< ReadonlyArray, OrchestrationCommandRejection | PlatformError.PlatformError, @@ -214,7 +193,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" userInputActivity, }: { readonly command: OrchestrationCommand; - readonly readModel: OrchestrationReadModel; + readonly readModel: CommandReadModel; readonly userInputActivity?: OrchestrationThreadActivity; }): Effect.fn.Return< DecideOrchestrationCommandResult, @@ -506,8 +485,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" (activity) => command.type === "thread.auto-settle" || activity.kind !== "user-input.requested" || - !Predicate.isObject(activity.payload) || - activity.payload.responseMode !== "message", + activity.responseMode !== "message", ) ) { return yield* new OrchestrationThreadSettleBlockedError({ threadId: command.threadId }); diff --git a/apps/server/src/orchestration/decider.turnDiffComplete.test.ts b/apps/server/src/orchestration/decider.turnDiffComplete.test.ts index fa52c4eaa044..425cf3d1d88d 100644 --- a/apps/server/src/orchestration/decider.turnDiffComplete.test.ts +++ b/apps/server/src/orchestration/decider.turnDiffComplete.test.ts @@ -1,3 +1,4 @@ +import type { CommandReadModel } from "./CommandReadModel.ts"; import { CheckpointRef, CommandId, @@ -7,7 +8,6 @@ import { ThreadId, TurnId, type OrchestrationCheckpointSummary, - type OrchestrationReadModel, } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; @@ -54,7 +54,7 @@ function makeReadModel(checkpoints: ReadonlyArray, -): OrchestrationReadModel { +function makeReadModel(activities: ReadonlyArray): CommandReadModel { return { snapshotSequence: 0, projects: [], @@ -64,7 +62,7 @@ function makeReadModel( deletedAt: null, messages: [], proposedPlans: [], - activities: [...activities], + activities: activities.map(toCommandActivity), checkpoints: [], session: null, }, diff --git a/apps/server/src/orchestration/http.ts b/apps/server/src/orchestration/http.ts index 82f16fc74ce6..6122390d52b1 100644 --- a/apps/server/src/orchestration/http.ts +++ b/apps/server/src/orchestration/http.ts @@ -34,13 +34,13 @@ export const orchestrationHttpApiLayer = HttpApiBuilder.group( Effect.fn("environment.orchestration.snapshot")(function* (args) { yield* annotateEnvironmentRequest(args.endpoint.name); yield* requireEnvironmentScope(AuthOrchestrationReadScope); - // Serve the lightweight command read model (thread bodies empty) + // Serve the lightweight metadata snapshot (thread bodies empty) // instead of the fully hydrated snapshot. Hydrating every message // and activity payload in the database has OOM-killed servers, and // the route's only consumer (the project CLI) reads projects alone — // UI clients load the shell and per-thread snapshots instead. return yield* projectionSnapshotQuery - .getCommandReadModel() + .getMetadataSnapshot() .pipe( Effect.catch((cause) => failEnvironmentInternal("orchestration_snapshot_failed", cause), diff --git a/apps/server/src/orchestration/messageContext.test.ts b/apps/server/src/orchestration/messageContext.test.ts index 0dcbbab292cd..81247859b42a 100644 --- a/apps/server/src/orchestration/messageContext.test.ts +++ b/apps/server/src/orchestration/messageContext.test.ts @@ -1,21 +1,17 @@ +import type { CommandReadModel } from "./CommandReadModel.ts"; import { CommandId, - EventId, MessageId, ProjectId, - ProviderDriverKind, ProviderInstanceId, ThreadId, - type OrchestrationEvent, type OrchestrationMessageContext, - type OrchestrationReadModel, } from "@t3tools/contracts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { expect, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import { decideOrchestrationCommand } from "./decider.ts"; -import { createEmptyReadModel, projectEvent } from "./projector.ts"; const NOW = "2026-01-01T00:00:00.000Z"; @@ -32,7 +28,7 @@ const context: OrchestrationMessageContext = { ], }; -function makeReadModel(): OrchestrationReadModel { +function makeReadModel(): CommandReadModel { return { snapshotSequence: 0, projects: [], @@ -68,20 +64,6 @@ function makeReadModel(): OrchestrationReadModel { }; } -function makeEvent(sequence: number, type: OrchestrationEvent["type"], payload: unknown) { - return { - sequence, - eventId: EventId.make(`event-${sequence}`), - type, - aggregateKind: "thread", - aggregateId: ThreadId.make("thread-1"), - occurredAt: NOW, - commandId: CommandId.make(`cmd-${sequence}`), - causationEventId: null, - payload, - } as OrchestrationEvent; -} - it.layer(NodeServices.layer)("message context plumbing", (it) => { it.effect("carries context records from turn start into the message-sent event", () => Effect.gen(function* () { @@ -110,56 +92,4 @@ it.layer(NodeServices.layer)("message context plumbing", (it) => { ); }), ); - - it.effect("projects context records onto the read-model message", () => - Effect.gen(function* () { - const afterCreate = yield* projectEvent( - createEmptyReadModel(NOW), - makeEvent(1, "thread.created", { - threadId: "thread-1", - projectId: "project-1", - title: "demo", - modelSelection: { provider: ProviderDriverKind.make("codex"), model: "gpt-5.4" }, - runtimeMode: "full-access", - branch: null, - worktreePath: null, - createdAt: NOW, - updatedAt: NOW, - }), - ); - const afterMessage = yield* projectEvent( - afterCreate, - makeEvent(2, "thread.message-sent", { - threadId: "thread-1", - messageId: "message-1", - role: "user", - text: "Use [$pinchtab](t3-context://v1/skill/ctx_1)", - attachments: [], - context, - turnId: null, - streaming: false, - createdAt: NOW, - updatedAt: NOW, - }), - ); - const message = afterMessage.threads[0]?.messages[0]; - expect(message?.context).toEqual(context); - - // A later non-streaming update without context keeps the original records. - const afterUpdate = yield* projectEvent( - afterMessage, - makeEvent(3, "thread.message-sent", { - threadId: "thread-1", - messageId: "message-1", - role: "user", - text: "edited", - turnId: null, - streaming: false, - createdAt: NOW, - updatedAt: NOW, - }), - ); - expect(afterUpdate.threads[0]?.messages[0]?.context).toEqual(context); - }), - ); }); diff --git a/apps/server/src/orchestration/projector.pullRequests.test.ts b/apps/server/src/orchestration/projector.pullRequests.test.ts index 6876cdde4f0c..6ee032d79066 100644 --- a/apps/server/src/orchestration/projector.pullRequests.test.ts +++ b/apps/server/src/orchestration/projector.pullRequests.test.ts @@ -1,10 +1,10 @@ +import type { CommandReadModel } from "./CommandReadModel.ts"; import { CommandId, EventId, ProjectId, ThreadId, type OrchestrationEvent, - type OrchestrationReadModel, type RepositoryIdentity, type ThreadPullRequestLink, type ThreadPullRequestSnapshot, @@ -63,7 +63,7 @@ const snapshot: ThreadPullRequestSnapshot = { syncedAt: LATER, }; -const createThread = (model: OrchestrationReadModel) => +const createThread = (model: CommandReadModel) => projectEvent( model, makeEvent({ @@ -84,7 +84,7 @@ const createThread = (model: OrchestrationReadModel) => }), ); -const createProject = (model: OrchestrationReadModel, repositoryIdentity: RepositoryIdentity) => +const createProject = (model: CommandReadModel, repositoryIdentity: RepositoryIdentity) => projectEvent(model, { ...makeEvent({ sequence: model.snapshotSequence + 1, diff --git a/apps/server/src/orchestration/projector.test.ts b/apps/server/src/orchestration/projector.test.ts index 516b85fbdae8..02f8747996e3 100644 --- a/apps/server/src/orchestration/projector.test.ts +++ b/apps/server/src/orchestration/projector.test.ts @@ -698,8 +698,8 @@ describe("orchestration projector", () => { const message = afterComplete.threads[0]?.messages[0]; expect(message?.id).toBe("assistant:msg-1"); - expect(message?.text).toBe("hello"); - expect(message?.streaming).toBe(false); + expect(message).not.toHaveProperty("text"); + expect(message).not.toHaveProperty("streaming"); expect(message?.updatedAt).toBe(completeAt); }); @@ -905,12 +905,7 @@ describe("orchestration projector", () => { ); const thread = afterRevert.threads[0]; - expect(thread?.messages.map((message) => ({ role: message.role, text: message.text }))).toEqual( - [ - { role: "user", text: "First edit" }, - { role: "assistant", text: "Updated README to v2.\n" }, - ], - ); + expect(thread?.messages.map((message) => message.role)).toEqual(["user", "assistant"]); expect( thread?.activities.map((activity) => ({ id: activity.id, turnId: activity.turnId })), ).toEqual([{ id: "activity-1", turnId: "turn-1" }]); diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index dd418b95b9bc..b45d5b223016 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -1,7 +1,6 @@ import type { OrchestrationEvent, OrchestrationProject, - OrchestrationReadModel, ThreadId, ThreadLinkedPullRequest, ThreadPullRequestKey, @@ -9,10 +8,7 @@ import type { } from "@t3tools/contracts"; import { isImportedAgentSessionMessageId, - OrchestrationCheckpointSummary, - OrchestrationMessage, OrchestrationSession, - OrchestrationThread, WORKTREE_SETUP_ACTIVITY_KIND, } from "@t3tools/contracts"; import { @@ -23,7 +19,19 @@ import { import { compareDateTimeStrings } from "@t3tools/shared/dateTime"; import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; -import * as Predicate from "effect/Predicate"; + +import { + CommandReadModel, + CommandThread, + CommandMessage, + CommandCheckpoint, + CommandProposedPlan, + MAX_COMMAND_MESSAGES, + MAX_COMMAND_ACTIVITIES, + MAX_COMMAND_CHECKPOINTS, + MAX_COMMAND_PLANS, + toCommandActivity, +} from "./CommandReadModel.ts"; import { toProjectorDecodeError, type OrchestrationProjectorDecodeError } from "./Errors.ts"; import { @@ -56,21 +64,18 @@ import { ThreadTurnDiffCompletedPayload, } from "./Schemas.ts"; -type ThreadPatch = Partial>; -const MAX_THREAD_MESSAGES = 2_000; -const MAX_THREAD_CHECKPOINTS = 500; +type ThreadPatch = Partial>; // Async questions can stay open while the agent produces more activity. // Match the database snapshot's pending-question retention. -function retainThreadActivities(activities: OrchestrationThread["activities"]) { - const recentStart = activities.length - 500; +function retainThreadActivities(activities: CommandThread["activities"]) { + const recentStart = activities.length - MAX_COMMAND_ACTIVITIES; if (recentStart <= 0) return activities; - const pending = new Map(); + const pending = new Map(); for (const activity of activities) { - if (!Predicate.isObject(activity.payload)) continue; - const requestId = activity.payload.requestId; - if (typeof requestId !== "string") continue; - if (activity.kind === "user-input.requested" && activity.payload.responseMode === "message") { + const { requestId } = activity; + if (requestId === null) continue; + if (activity.kind === "user-input.requested" && activity.responseMode === "message") { pending.set(requestId, activity); } else if (activity.kind === "user-input.resolved") { pending.delete(requestId); @@ -81,9 +86,8 @@ function retainThreadActivities(activities: OrchestrationThread["activities"]) { (activity, index) => index >= recentStart || pendingActivities.has(activity) || - // The worktree setup record is upserted under one id for the thread's - // whole life and is the only durable copy of a running setup; an async - // setup script can outlast a chatty first turn. + // Preserve the existing exception for the thread's setup record. + // Its body stays in SQL; only its activity metadata lives here. activity.kind === WORKTREE_SETUP_ACTIVITY_KIND, ); } @@ -120,20 +124,20 @@ function settledTurnStateForSessionStatus( // Runs for every thread event (including streaming deltas) against every // thread the server has ever seen, so copy the array rather than map it. function updateThread( - threads: ReadonlyArray, + threads: ReadonlyArray, threadId: ThreadId, patch: ThreadPatch, -): ReadonlyArray { +): ReadonlyArray { const index = threads.findIndex((thread) => thread.id === threadId); return index === -1 ? threads : patchThreadAt(threads, index, patch); } /** For callers that already located the thread and must not scan again. */ function patchThreadAt( - threads: ReadonlyArray, + threads: ReadonlyArray, index: number, patch: ThreadPatch, -): ReadonlyArray { +): ReadonlyArray { const next = threads.slice(); next[index] = { ...threads[index]!, ...patch }; return next; @@ -141,10 +145,10 @@ function patchThreadAt( /** Patch that swaps a thread's links and re-derives the legacy single-PR field from them. */ function pullRequestsPatch( - thread: Pick, + thread: Pick, pullRequests: ReadonlyArray, - projects: OrchestrationReadModel["projects"], -): Pick { + projects: CommandReadModel["projects"], +): Pick { return { pullRequests, linkedPullRequest: legacyLinkedPullRequestOf( @@ -191,7 +195,7 @@ function legacyPullRequestHost( } function legacyLinkToPullRequests( - thread: Pick, + thread: Pick, project: OrchestrationProject | undefined, linked: ThreadLinkedPullRequest | null, linkedAt: string, @@ -222,10 +226,10 @@ function decodeForEvent( } function retainThreadMessagesAfterRevert( - messages: ReadonlyArray, + messages: ReadonlyArray, retainedTurnIds: ReadonlySet, turnCount: number, -): ReadonlyArray { +): ReadonlyArray { const retainedMessageIds = new Set(); for (const message of messages) { if (message.role === "system" || isImportedAgentSessionMessageId(message.id)) { @@ -293,26 +297,26 @@ function retainThreadMessagesAfterRevert( } function retainThreadActivitiesAfterRevert( - activities: ReadonlyArray, + activities: ReadonlyArray, retainedTurnIds: ReadonlySet, -): ReadonlyArray { +): ReadonlyArray { return activities.filter( (activity) => activity.turnId === null || retainedTurnIds.has(activity.turnId), ); } function retainThreadProposedPlansAfterRevert( - proposedPlans: ReadonlyArray, + proposedPlans: ReadonlyArray, retainedTurnIds: ReadonlySet, -): ReadonlyArray { +): ReadonlyArray { return proposedPlans.filter( (proposedPlan) => proposedPlan.turnId === null || retainedTurnIds.has(proposedPlan.turnId), ); } function compareThreadActivities( - left: OrchestrationThread["activities"][number], - right: OrchestrationThread["activities"][number], + left: CommandThread["activities"][number], + right: CommandThread["activities"][number], ): number { if (left.sequence !== undefined && right.sequence !== undefined) { if (left.sequence !== right.sequence) { @@ -327,7 +331,7 @@ function compareThreadActivities( return left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id); } -export function createEmptyReadModel(nowIso: string): OrchestrationReadModel { +export function createEmptyReadModel(nowIso: string): CommandReadModel { return { snapshotSequence: 0, projects: [], @@ -337,10 +341,10 @@ export function createEmptyReadModel(nowIso: string): OrchestrationReadModel { } export function projectEvent( - model: OrchestrationReadModel, + model: CommandReadModel, event: OrchestrationEvent, -): Effect.Effect { - const nextBase: OrchestrationReadModel = { +): Effect.Effect { + const nextBase: CommandReadModel = { ...model, snapshotSequence: event.sequence, updatedAt: event.occurredAt, @@ -434,8 +438,8 @@ export function projectEvent( event.type, "payload", ); - const thread: OrchestrationThread = yield* decodeForEvent( - OrchestrationThread, + const thread: CommandThread = yield* decodeForEvent( + CommandThread, { id: payload.threadId, projectId: payload.projectId, @@ -460,6 +464,7 @@ export function projectEvent( snoozedAt: null, deletedAt: null, messages: [], + proposedPlans: [], activities: [], checkpoints: [], session: null, @@ -787,16 +792,12 @@ export function projectEvent( return nextBase; } - const message: OrchestrationMessage = yield* decodeForEvent( - OrchestrationMessage, + const message: CommandMessage = yield* decodeForEvent( + CommandMessage, { id: payload.messageId, role: payload.role, - text: payload.text, - ...(payload.attachments !== undefined ? { attachments: payload.attachments } : {}), - ...(payload.context !== undefined ? { context: payload.context } : {}), turnId: payload.turnId, - streaming: payload.streaming, createdAt: payload.createdAt, updatedAt: payload.updatedAt, }, @@ -810,23 +811,13 @@ export function projectEvent( entry.id === message.id ? { ...entry, - text: message.streaming - ? `${entry.text}${message.text}` - : message.text.length > 0 - ? message.text - : entry.text, - streaming: message.streaming, updatedAt: message.updatedAt, turnId: message.turnId, - ...(message.attachments !== undefined - ? { attachments: message.attachments } - : {}), - ...(message.context !== undefined ? { context: message.context } : {}), } : entry, ) : [...thread.messages, message]; - const cappedMessages = messages.slice(-MAX_THREAD_MESSAGES); + const cappedMessages = messages.slice(-MAX_COMMAND_MESSAGES); return { ...nextBase, @@ -915,13 +906,18 @@ export function projectEvent( const proposedPlans = [ ...thread.proposedPlans.filter((entry) => entry.id !== payload.proposedPlan.id), - payload.proposedPlan, + yield* decodeForEvent( + CommandProposedPlan, + payload.proposedPlan, + event.type, + "proposedPlan", + ), ] .toSorted( (left, right) => left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id), ) - .slice(-200); + .slice(-MAX_COMMAND_PLANS); return { ...nextBase, @@ -946,13 +942,12 @@ export function projectEvent( } const checkpoint = yield* decodeForEvent( - OrchestrationCheckpointSummary, + CommandCheckpoint, { turnId: payload.turnId, checkpointTurnCount: payload.checkpointTurnCount, checkpointRef: payload.checkpointRef, status: payload.status, - files: payload.files, assistantMessageId: payload.assistantMessageId, completedAt: payload.completedAt, }, @@ -975,7 +970,7 @@ export function projectEvent( checkpoint, ] .toSorted((left, right) => left.checkpointTurnCount - right.checkpointTurnCount) - .slice(-MAX_THREAD_CHECKPOINTS); + .slice(-MAX_COMMAND_CHECKPOINTS); // Mid-turn diff updates produce placeholder checkpoints; record the // checkpoint, but don't settle a turn its session is still running. @@ -1022,17 +1017,17 @@ export function projectEvent( const checkpoints = thread.checkpoints .filter((entry) => entry.checkpointTurnCount <= payload.turnCount) .toSorted((left, right) => left.checkpointTurnCount - right.checkpointTurnCount) - .slice(-MAX_THREAD_CHECKPOINTS); + .slice(-MAX_COMMAND_CHECKPOINTS); const retainedTurnIds = new Set(checkpoints.map((checkpoint) => checkpoint.turnId)); const messages = retainThreadMessagesAfterRevert( thread.messages, retainedTurnIds, payload.turnCount, - ).slice(-MAX_THREAD_MESSAGES); + ).slice(-MAX_COMMAND_MESSAGES); const proposedPlans = retainThreadProposedPlansAfterRevert( thread.proposedPlans, retainedTurnIds, - ).slice(-200); + ).slice(-MAX_COMMAND_PLANS); const activities = retainThreadActivitiesAfterRevert(thread.activities, retainedTurnIds); const latestCheckpoint = checkpoints.at(-1) ?? null; @@ -1079,7 +1074,7 @@ export function projectEvent( const activities = retainThreadActivities( [ ...thread.activities.filter((entry) => entry.id !== payload.activity.id), - payload.activity, + toCommandActivity(payload.activity), ].toSorted(compareThreadActivities), ); diff --git a/apps/server/src/project/AgentSessionScanner.test.ts b/apps/server/src/project/AgentSessionScanner.test.ts index 2dde0643ed60..7b354043fc66 100644 --- a/apps/server/src/project/AgentSessionScanner.test.ts +++ b/apps/server/src/project/AgentSessionScanner.test.ts @@ -36,6 +36,7 @@ const makeProjectShell = (workspaceRoot: string): OrchestrationProjectShell => ( /** Only `getShellSnapshot` is exercised; the rest must not be called. */ const makeProjectionSnapshotQueryLayer = (importedWorkspaceRoots: ReadonlyArray) => Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), diff --git a/apps/server/src/project/ProjectSetupScriptRunner.test.ts b/apps/server/src/project/ProjectSetupScriptRunner.test.ts index 69c553a4deaf..2896f764f8f5 100644 --- a/apps/server/src/project/ProjectSetupScriptRunner.test.ts +++ b/apps/server/src/project/ProjectSetupScriptRunner.test.ts @@ -30,6 +30,7 @@ const makeProjectionSnapshotQueryLayer = (project: OrchestrationProject) => Layer.succeed(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 7f0465b9c401..904dbff7af81 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -5088,6 +5088,7 @@ describe("agent browser access", () => { getImportedAgentSessionSources: () => Effect.die("unused"), getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 7b1fec90f867..2ffd04d51016 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -235,6 +235,7 @@ describe("ProviderSessionReaper", () => { Layer.succeed(ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.die("unused"), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 43595dbafce1..bc1bd9cea911 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -1018,7 +1018,7 @@ const buildAppUnderTest = (options?: { Layer.provide( Layer.mock(ProjectionSnapshotQuery.ProjectionSnapshotQuery)({ getUserInputActivity: () => Effect.die("unused"), - getCommandReadModel: () => Effect.succeed(makeDefaultOrchestrationReadModel()), + getMetadataSnapshot: () => Effect.succeed(makeDefaultOrchestrationReadModel()), getSnapshot: () => Effect.succeed(makeDefaultOrchestrationReadModel()), getShellSnapshot: () => Effect.succeed({ diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 050650b11515..2a315f3f887a 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -75,7 +75,7 @@ const makeProviderService = (liveThreadIds: ReadonlyArray = []) => const queryWithThreads = (threads: ReadonlyArray>) => ({ getUserInputActivity: () => Effect.die("unused"), - getCommandReadModel: () => Effect.succeed({ threads } as never), + getMetadataSnapshot: () => Effect.succeed({ threads } as never), }) as unknown as ProjectionSnapshotQuery.ProjectionSnapshotQuery["Service"]; const runReconciliation = (input: { @@ -690,7 +690,7 @@ it.effect("does not fail startup when the live provider session inventory cannot return ServerRuntimeStartup.reconcileProviderSessions.pipe( Effect.provideService(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), - getCommandReadModel: () => + getMetadataSnapshot: () => Effect.sync(() => { queried = true; return { threads: [] } as never; diff --git a/apps/server/src/serverRuntimeStartup.test.ts b/apps/server/src/serverRuntimeStartup.test.ts index 4bbdb2e9c8b5..74ecdd14377d 100644 --- a/apps/server/src/serverRuntimeStartup.test.ts +++ b/apps/server/src/serverRuntimeStartup.test.ts @@ -165,6 +165,7 @@ it.effect("resolveAutoBootstrapWelcomeTargets returns existing project and threa Effect.provideService(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.succeed([]), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), @@ -295,6 +296,7 @@ it.effect.each([ Effect.provideService(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.succeed([]), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), @@ -383,6 +385,7 @@ it.effect( Effect.provideService(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.succeed([]), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), @@ -449,6 +452,7 @@ it.effect("resolveAutoBootstrapWelcomeTargets preserves typed UUID generation fa Effect.provideService(ProjectionSnapshotQuery.ProjectionSnapshotQuery, { getUserInputActivity: () => Effect.die("unused"), listActivitiesByKind: () => Effect.succeed([]), + getMetadataSnapshot: () => Effect.die("unused"), getCommandReadModel: () => Effect.die("unused"), getSnapshot: () => Effect.die("unused"), getShellSnapshot: () => Effect.die("unused"), diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 138dbcedfa85..1bb5bdd00f52 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -405,7 +405,7 @@ const toServerUpdateThreadContinuationError = (cause: unknown) => export const markRunningProviderSessionsForContinuation = Effect.gen(function* () { const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory; const query = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; - const { threads } = yield* query.getCommandReadModel(); + const { threads } = yield* query.getMetadataSnapshot(); const running = threads.filter( (thread) => thread.archivedAt === null && @@ -502,7 +502,7 @@ export const reconcileProviderSessions = Effect.gen(function* () { const liveThreadIds = new Set( (yield* providerService.listSessions()).map((session) => session.threadId), ); - const { threads } = yield* query.getCommandReadModel(); + const { threads } = yield* query.getMetadataSnapshot(); // Provider startup can report ready before the continuation is submitted. // Find those markers in one read rather than querying every idle thread. const preparedThreadIds = new Set( @@ -759,7 +759,7 @@ export const reconcileWorktreeSetups = Effect.gen(function* () { const crypto = yield* Crypto.Crypto; const orchestrationEngine = yield* OrchestrationEngine.OrchestrationEngineService; const query = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; - // The command read model carries no activity bodies; read the setup + // The metadata snapshot carries no activity bodies; read the setup // records directly, live threads only. const recordedSetups = yield* query.listActivitiesByKind(WORKTREE_SETUP_ACTIVITY_KIND); const interruptedAt = DateTime.formatIso(yield* DateTime.now);