diff --git a/apps/server/src/observability/Metrics.ts b/apps/server/src/observability/Metrics.ts index b659ac09ae12..b75ace399ccc 100644 --- a/apps/server/src/observability/Metrics.ts +++ b/apps/server/src/observability/Metrics.ts @@ -19,22 +19,6 @@ export const rpcRequestDuration = Metric.timer("t3_rpc_request_duration", { description: "RPC request handling duration.", }); -export const orchestrationCommandsTotal = Metric.counter("t3_orchestration_commands_total", { - description: "Total orchestration commands dispatched.", -}); - -export const orchestrationCommandDuration = Metric.timer("t3_orchestration_command_duration", { - description: "Orchestration command dispatch duration.", -}); - -export const orchestrationCommandAckDuration = Metric.timer( - "t3_orchestration_command_ack_duration", - { - description: - "Time from orchestration command dispatch to the first committed domain event emitted for that command.", - }, -); - const orchestrationEventsProcessedTotal = Metric.counter( "t3_orchestration_events_processed_total", { diff --git a/apps/server/src/orchestration-v2/CommandReceiptStore.ts b/apps/server/src/orchestration-v2/CommandReceiptStore.ts index 515565e1d8f6..19ebfef0b593 100644 --- a/apps/server/src/orchestration-v2/CommandReceiptStore.ts +++ b/apps/server/src/orchestration-v2/CommandReceiptStore.ts @@ -1,4 +1,4 @@ -import { CommandId, NonNegativeInt, ThreadId } from "@t3tools/contracts"; +import { CommandId, NonNegativeInt, ProjectId, ThreadId } from "@t3tools/contracts"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -62,14 +62,35 @@ export const CommandReceiptV2 = Schema.Struct({ }); export type CommandReceiptV2 = typeof CommandReceiptV2.Type; +/** Receipt for a project command; shares the receipt table with thread commands. */ +export const ProjectCommandReceiptV2 = Schema.Struct({ + commandId: CommandId, + projectId: ProjectId, + commandType: Schema.String, + acceptedAt: Schema.DateTimeUtc, + resultSequence: NonNegativeInt, + status: CommandReceiptV2Status, + error: Schema.NullOr(Schema.String), +}); +export type ProjectCommandReceiptV2 = typeof ProjectCommandReceiptV2.Type; + +type AnyCommandReceiptV2 = CommandReceiptV2 | ProjectCommandReceiptV2; + export interface CommandReceiptStoreV2Shape { readonly insertIfAbsent: ( - receipt: CommandReceiptV2, + receipt: AnyCommandReceiptV2, ) => Effect.Effect; - readonly upsert: (receipt: CommandReceiptV2) => Effect.Effect; + readonly upsert: ( + receipt: AnyCommandReceiptV2, + ) => Effect.Effect; + /** A thread command's receipt; none when the command id belongs to a project command. */ readonly getByCommandId: ( commandId: CommandId, ) => Effect.Effect, CommandReceiptStoreV2Error>; + /** A project command's receipt; none when the command id belongs to a thread command. */ + readonly getProjectByCommandId: ( + commandId: CommandId, + ) => Effect.Effect, CommandReceiptStoreV2Error>; } export class CommandReceiptStoreV2 extends Context.Service< @@ -86,6 +107,12 @@ const decodeReceipt = Schema.decodeUnknownEffect( acceptedAt: Schema.DateTimeUtcFromString, })), ); +const decodeProjectReceipt = Schema.decodeUnknownEffect( + ProjectCommandReceiptV2.mapFields((fields) => ({ + ...fields, + acceptedAt: Schema.DateTimeUtcFromString, + })), +); function fromApplicationReceipt(receipt: OrchestrationCommandReceipt) { return decodeReceipt({ @@ -99,11 +126,12 @@ function fromApplicationReceipt(receipt: OrchestrationCommandReceipt) { }); } -function toApplicationReceipt(receipt: CommandReceiptV2): OrchestrationCommandReceipt { +function toApplicationReceipt(receipt: AnyCommandReceiptV2): OrchestrationCommandReceipt { return { commandId: receipt.commandId, - aggregateKind: "thread", - aggregateId: receipt.threadId, + ...("projectId" in receipt + ? { aggregateKind: "project" as const, aggregateId: receipt.projectId } + : { aggregateKind: "thread" as const, aggregateId: receipt.threadId }), commandType: receipt.commandType, acceptedAt: DateTime.formatIso(receipt.acceptedAt), resultSequence: receipt.resultSequence, @@ -158,6 +186,23 @@ const baseLayer: Layer.Layer + receipts.getByCommandId({ commandId }).pipe( + Effect.flatMap((receipt) => + Option.isNone(receipt) || receipt.value.aggregateKind !== "project" + ? Effect.succeed(Option.none()) + : decodeProjectReceipt({ + commandId: receipt.value.commandId, + projectId: receipt.value.aggregateId, + commandType: receipt.value.commandType, + acceptedAt: receipt.value.acceptedAt, + resultSequence: receipt.value.resultSequence, + status: receipt.value.status, + error: receipt.value.error, + }).pipe(Effect.map(Option.some)), + ), + Effect.mapError((cause) => new CommandReceiptStoreReadError({ commandId, cause })), + ), } satisfies CommandReceiptStoreV2Shape); }), ); diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 8da8378a8b6e..e658d30e6cb0 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -22,20 +22,24 @@ import * as Stream from "effect/Stream"; import * as CheckpointStore from "../checkpointing/CheckpointStore.ts"; import { ServerConfig } from "../config.ts"; import { layer as mcpSessionRegistryTestLayer } from "../mcp/McpSessionRegistry.testkit.ts"; -import { OrchestrationEngineService } from "../orchestration/Services/OrchestrationEngine.ts"; -import { OrchestrationLayerLive } from "../orchestration/runtimeLayer.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import { ProjectEnrichmentService } from "../project/ProjectEnrichmentService.ts"; +import { ProjectService } from "../project/ProjectService.ts"; import type { ProviderInstance } from "../provider/ProviderDriver.ts"; import { ProviderInstanceRegistry } from "../provider/Services/ProviderInstanceRegistry.ts"; import { ServerSettingsService } from "../serverSettings.ts"; import * as VcsDriverRegistry from "../vcs/VcsDriverRegistry.ts"; import * as VcsProcess from "../vcs/VcsProcess.ts"; +import { WorkspacePaths } from "../workspace/WorkspacePaths.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import { EventSinkV2 } from "./EventSink.ts"; import { OrchestratorV2 } from "./Orchestrator.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; -import { OrchestrationV2EventSinkLayerLive, OrchestrationV2LayerLive } from "./runtimeLayer.ts"; +import { + OrchestrationV2EventSinkLayerLive, + OrchestrationV2LayerLive, + ProjectServiceLayerLive, +} from "./runtimeLayer.ts"; import { worktreeRepairDependenciesTestLayer } from "./ProviderTurnStartService.testkit.ts"; const PlatformTestLayer = Layer.merge( @@ -94,11 +98,13 @@ const TestProviderInstanceRegistry = Layer.succeed(ProviderInstanceRegistry, { subscribeChanges: Effect.never, }); -const TestLayer = Layer.mergeAll( - OrchestrationLayerLive, - OrchestrationV2LayerLive, - OrchestrationV2EventSinkLayerLive, -).pipe( +const TestLayer = Layer.mergeAll(OrchestrationV2LayerLive, OrchestrationV2EventSinkLayerLive).pipe( + Layer.provideMerge(ProjectServiceLayerLive), + Layer.provide( + Layer.mock(WorkspacePaths)({ + normalizeWorkspaceRoot: (workspaceRoot) => Effect.succeed(workspaceRoot), + }), + ), Layer.provide(worktreeRepairDependenciesTestLayer), Layer.provide( Layer.succeed(ProjectEnrichmentService, { @@ -140,22 +146,18 @@ const seedParentWithTerminalTask = (input: { readonly now: DateTime.Utc; }) => Effect.gen(function* () { - const applicationEngine = yield* OrchestrationEngineService; + const projects = yield* ProjectService; const orchestrator = yield* OrchestratorV2; const eventSink = yield* EventSinkV2; const providerThreadId = ProviderThreadId.make( `provider-thread:${String(input.threadId).replace("thread:", "")}`, ); - yield* applicationEngine.dispatch({ - type: "project.create", + yield* projects.create({ commandId: CommandId.make(`command:seed-project:${input.threadId}`), projectId: input.projectId, title: "Delegated completion delivery", workspaceRoot: `/workspace/${input.projectId}`, - defaultModelSelection: modelSelection, - scripts: [], - createdAt: DateTime.formatIso(input.now), }); yield* orchestrator.dispatch({ diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index e9306e144980..4fd6800d03e3 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -8,6 +8,7 @@ import { RunId, RuntimeRequestId, NodeId, + type ProjectId, ThreadId, } from "@t3tools/contracts"; import * as Context from "effect/Context"; @@ -21,11 +22,13 @@ import * as Stream from "effect/Stream"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { replayAndBufferProjectedLiveEvents } from "../orchestration/LiveStreamBudget.ts"; +import type { UnsequencedProjectEvent } from "../persistence/Services/OrchestrationEventStore.ts"; import { projectDomainEventForWire } from "./WireProjection.ts"; import { CommandReceiptStoreV2, type CommandReceiptV2, + type ProjectCommandReceiptV2, layer as commandReceiptStoreLayer, } from "./CommandReceiptStore.ts"; import { @@ -39,6 +42,7 @@ import { ORCHESTRATION_V2_PROJECTION_SCHEMA_VERSION, ProjectionStoreV2, } from "./ProjectionStore.ts"; +import * as ProjectStore from "./ProjectStore.ts"; import { TurnItemPositionStoreV2, layer as turnItemPositionStoreLayer, @@ -156,6 +160,28 @@ export interface EventSinkV2Shape { readonly rejectedAt: DateTime.Utc; readonly error: string; }) => Effect.Effect; + /** + * Append a project event, fold it into its row and record the receipt in one + * transaction. A reused command id commits nothing and returns its receipt. + */ + readonly commitProjectCommand: (input: { + readonly commandId: CommandId; + readonly projectId: ProjectId; + readonly commandType: string; + readonly acceptedAt: DateTime.Utc; + readonly event: UnsequencedProjectEvent; + }) => Effect.Effect< + { readonly receipt: ProjectCommandReceiptV2; readonly committed: boolean }, + EventSinkV2Error + >; + /** Record a rejected project command, or return the receipt its command id already has. */ + readonly commitRejectedProjectCommand: (input: { + readonly commandId: CommandId; + readonly projectId: ProjectId; + readonly commandType: string; + readonly rejectedAt: DateTime.Utc; + readonly error: string; + }) => Effect.Effect; readonly stream: (input?: { readonly threadId?: ThreadId; readonly afterSequence?: number; @@ -186,6 +212,7 @@ const baseLayer: Layer.Layer< | EffectOutboxV2 | EventStoreV2 | ProjectionStoreV2 + | ProjectStore.ProjectStoreV2 | SqlClient.SqlClient | TurnItemPositionStoreV2 > = Layer.effect( @@ -196,6 +223,7 @@ const baseLayer: Layer.Layer< const effectOutbox = yield* EffectOutboxV2; const eventStore = yield* EventStoreV2; const projectionStore = yield* ProjectionStoreV2; + const projectStore = yield* ProjectStore.ProjectStoreV2; const turnItemPositions = yield* TurnItemPositionStoreV2; const liveEvents = yield* PubSub.unbounded(); const liveEventsByType = new Map< @@ -568,6 +596,68 @@ const baseLayer: Layer.Layer< ); }); + const existingProjectReceipt = (commandId: CommandId) => + commandReceipts.getProjectByCommandId(commandId).pipe( + Effect.flatMap( + Option.match({ + onNone: () => + Effect.fail(`Command ${commandId} was already used by a thread command.` as const), + onSome: Effect.succeed, + }), + ), + ); + + const commitProjectCommandEffect = Effect.fn("orchestrationV2.EventSink.commitProjectCommand")( + function* (input: Parameters[0]) { + const result = yield* sql.withTransaction( + Effect.gen(function* () { + const reserved: ProjectCommandReceiptV2 = { + commandId: input.commandId, + projectId: input.projectId, + commandType: input.commandType, + acceptedAt: input.acceptedAt, + resultSequence: 0, + status: "accepted", + error: null, + }; + if (!(yield* commandReceipts.insertIfAbsent(reserved))) { + return { receipt: yield* existingProjectReceipt(input.commandId), event: undefined }; + } + const event = yield* eventStore.appendProjectEvent(input.event); + yield* projectStore.apply(event); + const receipt = { ...reserved, resultSequence: event.sequence }; + yield* commandReceipts.upsert(receipt); + return { receipt, event }; + }), + ); + if (result.event !== undefined) { + yield* eventStore.publishCommitted([result.event]); + } + return { receipt: result.receipt, committed: result.event !== undefined }; + }, + ); + + const commitRejectedProjectCommandEffect = Effect.fn( + "orchestrationV2.EventSink.commitRejectedProjectCommand", + )(function* (input: Parameters[0]) { + return yield* sql.withTransaction( + Effect.gen(function* () { + const receipt: ProjectCommandReceiptV2 = { + commandId: input.commandId, + projectId: input.projectId, + commandType: input.commandType, + acceptedAt: input.rejectedAt, + resultSequence: yield* eventStore.latestApplicationSequence, + status: "rejected", + error: input.error, + }; + return (yield* commandReceipts.insertIfAbsent(receipt)) + ? receipt + : yield* existingProjectReceipt(input.commandId); + }), + ); + }); + const catchUp = (input: { readonly afterSequence: number; readonly throughSequence: number; @@ -716,6 +806,20 @@ const baseLayer: Layer.Layer< }), ), ), + commitProjectCommand: (input) => + commitProjectCommandEffect(input).pipe( + Effect.mapError( + (cause) => + new EventSinkWriteError({ commandId: input.commandId, eventCount: 1, cause }), + ), + ), + commitRejectedProjectCommand: (input) => + commitRejectedProjectCommandEffect(input).pipe( + Effect.mapError( + (cause) => + new EventSinkWriteError({ commandId: input.commandId, eventCount: 0, cause }), + ), + ), stream: (input) => stream(input).pipe( Stream.mapError( @@ -766,6 +870,11 @@ export const layer: Layer.Layer< EventStoreV2 | ProjectionStoreV2 | SqlClient.SqlClient > = baseLayer.pipe( Layer.provide( - Layer.mergeAll(commandReceiptStoreLayer, effectOutboxLayer, turnItemPositionStoreLayer), + Layer.mergeAll( + commandReceiptStoreLayer, + effectOutboxLayer, + ProjectStore.layer, + turnItemPositionStoreLayer, + ), ), ); diff --git a/apps/server/src/orchestration-v2/EventStore.ts b/apps/server/src/orchestration-v2/EventStore.ts index 34c4c1f80051..a646152c6b80 100644 --- a/apps/server/src/orchestration-v2/EventStore.ts +++ b/apps/server/src/orchestration-v2/EventStore.ts @@ -1,4 +1,6 @@ import { + type ApplicationProjectEvent, + type ApplicationStoredEvent, CommandId, OrchestrationV2DomainEvent, OrchestrationV2StoredEvent, @@ -12,7 +14,10 @@ import * as Stream from "effect/Stream"; import type * as SqlClient from "effect/unstable/sql/SqlClient"; import { OrchestrationEventStoreLive } from "../persistence/Layers/OrchestrationEventStore.ts"; -import { OrchestrationEventStore } from "../persistence/Services/OrchestrationEventStore.ts"; +import { + OrchestrationEventStore, + type UnsequencedProjectEvent, +} from "../persistence/Services/OrchestrationEventStore.ts"; export class EventStoreAppendEventsError extends Schema.TaggedError()( "EventStoreAppendEventsError", @@ -52,6 +57,9 @@ export interface EventStoreV2Shape { readonly commandId?: CommandId; readonly events: ReadonlyArray; }) => Effect.Effect, EventStoreV2Error>; + readonly appendProjectEvent: ( + event: UnsequencedProjectEvent, + ) => Effect.Effect; readonly read: (input?: { readonly afterSequence?: number; readonly throughSequence?: number; @@ -65,9 +73,9 @@ export interface EventStoreV2Shape { readonly latestSequence: (input?: { readonly threadId?: ThreadId; }) => Effect.Effect; - readonly publishCommitted: ( - events: ReadonlyArray, - ) => Effect.Effect; + /** Latest sequence across project and V2 thread events. */ + readonly latestApplicationSequence: Effect.Effect; + readonly publishCommitted: (events: ReadonlyArray) => Effect.Effect; } export class EventStoreV2 extends Context.Service()( @@ -114,6 +122,12 @@ const baseLayer: Layer.Layer = Lay }), ), ), + appendProjectEvent: (event) => + applicationEvents + .appendProjectEvent(event) + .pipe( + Effect.mapError((cause) => new EventStoreAppendEventsError({ eventCount: 1, cause })), + ), read, readByCommandId: ({ commandId }) => applicationEvents @@ -129,6 +143,9 @@ const baseLayer: Layer.Layer = Lay }), ), ), + latestApplicationSequence: applicationEvents.latestApplicationSequence.pipe( + Effect.mapError((cause) => new EventStoreReadEventsError({ cause })), + ), publishCommitted: applicationEvents.publishCommitted, }); }), diff --git a/apps/server/src/orchestration-v2/ProjectCommands.test.ts b/apps/server/src/orchestration-v2/ProjectCommands.test.ts new file mode 100644 index 000000000000..55f2218fbd01 --- /dev/null +++ b/apps/server/src/orchestration-v2/ProjectCommands.test.ts @@ -0,0 +1,254 @@ +import { assert, describe, it } from "@effect/vitest"; +import { + CommandId, + EventId, + ProjectId, + ProviderInstanceId, + type ModelSelection, + type ProjectScript, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Result from "effect/Result"; + +import { + planProjectCommand, + type ProjectCommand, + type ProjectCommandState, +} from "./ProjectCommands.ts"; +import type { ProjectRow } from "./ProjectStore.ts"; + +const now = DateTime.makeUnsafe("2026-01-01T00:00:00.000Z"); +const projectId = ProjectId.make("project-scripts"); + +const row = (overrides: Partial = {}): ProjectRow => ({ + projectId, + title: "Scripts", + workspaceRoot: "/tmp/scripts", + defaultModelSelection: null, + defaultThreadEnvMode: null, + autoPull: false, + faviconPath: null, + projectIcon: null, + scripts: [], + createdAt: "2026-01-01T00:00:00.000Z", + updatedAt: "2026-01-01T00:00:00.000Z", + deletedAt: null, + ...overrides, +}); + +const script = (id: string): ProjectScript => ({ + id, + name: "Install dependencies", + command: "vp i", + icon: "configure", + runOnWorktreeCreate: false, +}); + +const plan = (command: ProjectCommand, state: Partial = {}) => + planProjectCommand({ + command, + state: { project: undefined, workspaceOwner: undefined, ...state }, + eventId: EventId.make("event:planned"), + now, + }); + +const update = ( + fields: Omit< + Extract, + "type" | "commandId" | "projectId" + >, +) => + plan( + { + type: "project.meta.update", + commandId: CommandId.make("cmd-update"), + projectId, + ...fields, + }, + { project: row() }, + ); + +const payloadOf = (result: ReturnType) => { + assert.isTrue(Result.isSuccess(result)); + return Result.getOrThrow(result).payload as Record; +}; + +const failureOf = (result: ReturnType) => { + assert.isTrue(Result.isFailure(result)); + return Result.isFailure(result) ? result.failure : assert.fail("expected a rejection"); +}; + +describe("planProjectCommand", () => { + it("creates projects with empty scripts and no model default", () => { + const result = plan({ + type: "project.create", + commandId: CommandId.make("cmd-create"), + projectId, + title: "Scripts", + workspaceRoot: "/tmp/scripts", + }); + const event = Result.getOrThrow(result); + assert.equal(event.type, "project.created"); + assert.equal(event.occurredAt, "2026-01-01T00:00:00.000Z"); + assert.deepInclude(event.payload, { scripts: [], defaultModelSelection: null }); + }); + + it("only treats metadata updates as explicit model defaults", () => { + const selection: ModelSelection = { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5.6-sol", + options: [{ id: "reasoningEffort", value: "high" }], + }; + assert.deepEqual( + payloadOf(update({ defaultModelSelection: selection })).defaultModelSelection, + selection, + ); + }); + + it("carries every edited field and omits the rest", () => { + const scripts = [script("lint")]; + const payload = payloadOf( + update({ + scripts, + defaultThreadEnvMode: "worktree", + autoPull: true, + faviconPath: "brand/icon.svg", + projectIcon: { kind: "lucide", name: "alarm-clock", color: "violet" }, + }), + ); + assert.deepInclude(payload, { + scripts, + defaultThreadEnvMode: "worktree", + autoPull: true, + faviconPath: "brand/icon.svg", + projectIcon: { kind: "lucide", name: "alarm-clock", color: "violet" }, + }); + const renamed = payloadOf(update({ title: "Renamed" })); + assert.isFalse("defaultThreadEnvMode" in renamed); + assert.isNull(payloadOf(update({ defaultThreadEnvMode: null })).defaultThreadEnvMode); + }); + + for (const id of ["install-javascript-dependencies", "A", "a.b", "a b", "-a", "a".repeat(25)]) { + it(`rejects a new script ID that cannot have a shortcut: ${id}`, () => { + const failure = failureOf(update({ scripts: [script("lint"), script(id)] })); + assert.equal(failure._tag, "ProjectCommandInvariantError"); + assert.include(failure.message, "Script ID"); + assert.include(failure.message, "24"); + // The detail is persisted in the rejected receipt, so it omits the raw ID. + assert.notInclude(failure.message, `'${id}'`); + }); + } + + it("accepts a script ID at the shortcut length limit", () => { + const scripts = [script("a".repeat(24))]; + assert.deepEqual(payloadOf(update({ scripts })).scripts, scripts); + }); + + it("keeps legacy scripts editable and removable while rejecting new invalid ones", () => { + const legacy = script("install-javascript-dependencies"); + const withLegacy = (scripts: ReadonlyArray) => + plan( + { + type: "project.meta.update", + commandId: CommandId.make("cmd-repair-script"), + projectId, + scripts, + }, + { project: row({ scripts: [legacy] }) }, + ); + for (const scripts of [[{ ...legacy, command: "vp install" }, script("lint")], []]) { + assert.deepEqual(payloadOf(withLegacy(scripts)).scripts, scripts); + } + assert.equal( + failureOf(withLegacy([legacy, script("another.invalid.id")]))._tag, + "ProjectCommandInvariantError", + ); + }); + + it("limits monograms to two graphemes", () => { + for (const text of ["T3", "é", "किखि", "क्ष्म", "각"]) { + const monogram = { kind: "monogram", text, color: "violet" } as const; + assert.deepEqual(payloadOf(update({ projectIcon: monogram })).projectIcon, monogram); + } + for (const text of ["ABC", "किखिगि"]) { + assert.equal( + failureOf(update({ projectIcon: { kind: "monogram", text, color: "violet" } }))._tag, + "ProjectCommandInvariantError", + ); + } + }); + + it("rejects a workspace root held by another active project", () => { + const owner = row({ + projectId: ProjectId.make("project-existing"), + workspaceRoot: "/tmp/project", + }); + const create = failureOf( + plan( + { + type: "project.create", + commandId: CommandId.make("cmd-duplicate-root"), + projectId: ProjectId.make("project-duplicate-root"), + title: "Duplicate", + workspaceRoot: "/tmp/project", + }, + { workspaceOwner: owner }, + ), + ); + assert.equal(create._tag, "ProjectWorkspaceConflictError"); + assert.equal( + create.message, + "Active project 'project-existing' already exists for workspace root '/tmp/project'.", + ); + const move = failureOf( + plan( + { + type: "project.meta.update", + commandId: CommandId.make("cmd-move-root"), + projectId, + workspaceRoot: "/tmp/project", + }, + { project: row(), workspaceOwner: owner }, + ), + ); + assert.equal(move._tag, "ProjectWorkspaceConflictError"); + }); + + it("requires the project to exist, and to be absent on create", () => { + const create: ProjectCommand = { + type: "project.create", + commandId: CommandId.make("cmd-create-twice"), + projectId, + title: "Twice", + workspaceRoot: "/tmp/twice", + }; + assert.include(failureOf(plan(create, { project: row() })).message, "cannot be created twice"); + const deleted = row({ deletedAt: "2026-01-01T00:00:00.000Z" }); + assert.include( + failureOf(plan(create, { project: deleted })).message, + "cannot be created twice", + ); + for (const command of [ + { type: "project.meta.update", commandId: CommandId.make("cmd-missing"), projectId }, + { type: "project.delete", commandId: CommandId.make("cmd-missing"), projectId }, + ] as const) { + for (const project of [undefined, deleted]) { + assert.equal( + failureOf(plan(command, { project }))._tag, + "ProjectCommandMissingProjectError", + ); + } + } + }); + + it("deletes with a single project.deleted event", () => { + const event = Result.getOrThrow( + plan( + { type: "project.delete", commandId: CommandId.make("cmd-delete"), projectId }, + { project: row() }, + ), + ); + assert.equal(event.type, "project.deleted"); + assert.deepEqual(event.payload, { projectId, deletedAt: "2026-01-01T00:00:00.000Z" }); + }); +}); diff --git a/apps/server/src/orchestration-v2/ProjectCommands.ts b/apps/server/src/orchestration-v2/ProjectCommands.ts new file mode 100644 index 000000000000..3de592283072 --- /dev/null +++ b/apps/server/src/orchestration-v2/ProjectCommands.ts @@ -0,0 +1,241 @@ +import { + type CommandId, + type EventId, + MAX_SCRIPT_ID_LENGTH, + type ModelSelection, + type ProjectIconOverride, + ProjectId, + type ProjectScript, + SCRIPT_RUN_COMMAND_PATTERN, + type ThreadEnvMode, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Result from "effect/Result"; +import * as Schema from "effect/Schema"; + +import type { UnsequencedProjectEvent } from "../persistence/Services/OrchestrationEventStore.ts"; +import type { ProjectRow } from "./ProjectStore.ts"; + +export interface ProjectCreateCommand { + readonly type: "project.create"; + readonly commandId: CommandId; + readonly projectId: ProjectId; + readonly title: string; + readonly workspaceRoot: string; + readonly scripts?: ReadonlyArray; +} + +export interface ProjectMetaUpdateCommand { + readonly type: "project.meta.update"; + readonly commandId: CommandId; + readonly projectId: ProjectId; + readonly title?: string; + readonly workspaceRoot?: string; + readonly defaultModelSelection?: ModelSelection | null; + readonly defaultThreadEnvMode?: ThreadEnvMode | null; + readonly autoPull?: boolean; + readonly faviconPath?: string | null; + readonly projectIcon?: ProjectIconOverride | null; + readonly scripts?: ReadonlyArray; +} + +export interface ProjectDeleteCommand { + readonly type: "project.delete"; + readonly commandId: CommandId; + readonly projectId: ProjectId; +} + +export type ProjectCommand = ProjectCreateCommand | ProjectMetaUpdateCommand | ProjectDeleteCommand; + +export class ProjectCommandInvariantError extends Schema.TaggedError()( + "ProjectCommandInvariantError", + { + commandType: Schema.String, + detail: Schema.String, + }, +) { + override get message(): string { + return `Project command invariant failed (${this.commandType}): ${this.detail}`; + } +} + +/** The command targets a project that does not exist or was deleted. */ +export class ProjectCommandMissingProjectError extends Schema.TaggedError()( + "ProjectCommandMissingProjectError", + { + commandType: Schema.String, + projectId: ProjectId, + }, +) { + override get message(): string { + return `Project '${this.projectId}' does not exist for command '${this.commandType}'.`; + } +} + +export class ProjectWorkspaceConflictError extends Schema.TaggedError()( + "ProjectWorkspaceConflictError", + { + workspaceRoot: Schema.String, + conflictingProjectId: ProjectId, + }, +) { + override get message(): string { + return `Active project '${this.conflictingProjectId}' already exists for workspace root '${this.workspaceRoot}'.`; + } +} + +export const ProjectCommandRejection = Schema.Union([ + ProjectCommandInvariantError, + ProjectCommandMissingProjectError, + ProjectWorkspaceConflictError, +]); +export type ProjectCommandRejection = typeof ProjectCommandRejection.Type; + +const ProjectCommandRejectionJson = Schema.fromJsonString(ProjectCommandRejection); +/** + * A rejected receipt stores its rejection as JSON, so a retried command id + * replays the same typed error even when a fresh plan would now succeed. + */ +export const encodeProjectCommandRejection = Schema.encodeSync(ProjectCommandRejectionJson); +/** None for receipts that predate structured rejections. */ +export const decodeProjectCommandRejection = Schema.decodeUnknownOption( + ProjectCommandRejectionJson, +); + +export interface ProjectCommandState { + /** The target project's row, including a soft-deleted one; only create sees deleted rows as taken. */ + readonly project: ProjectRow | undefined; + /** The active project that holds the command's requested workspace root, if any. */ + readonly workspaceOwner: ProjectRow | undefined; +} + +const monogramSegmenter = new Intl.Segmenter(undefined, { granularity: "grapheme" }); +const isScriptRunCommand = Schema.is(SCRIPT_RUN_COMMAND_PATTERN); + +/** + * Decide one project command against the rows it touches. The caller reads + * `state` under the project's lock and commits the planned event. + */ +export function planProjectCommand(input: { + readonly command: ProjectCommand; + readonly state: ProjectCommandState; + readonly eventId: EventId; + readonly now: DateTime.Utc; +}): Result.Result { + const { command, state } = input; + const invariant = (detail: string) => + Result.fail(new ProjectCommandInvariantError({ commandType: command.type, detail })); + const missingProject = () => + Result.fail( + new ProjectCommandMissingProjectError({ + commandType: command.type, + projectId: command.projectId, + }), + ); + const activeProject = state.project?.deletedAt === null ? state.project : undefined; + const requireWorkspaceAvailable = (workspaceRoot: string) => + state.workspaceOwner === undefined || state.workspaceOwner.projectId === command.projectId + ? undefined + : new ProjectWorkspaceConflictError({ + workspaceRoot, + conflictingProjectId: state.workspaceOwner.projectId, + }); + const occurredAt = DateTime.formatIso(input.now); + const base = { + eventId: input.eventId, + aggregateKind: "project" as const, + aggregateId: command.projectId, + occurredAt, + commandId: command.commandId, + causationEventId: null, + correlationId: command.commandId, + metadata: {}, + }; + + switch (command.type) { + case "project.create": { + if (state.project !== undefined) { + return invariant( + `Project '${command.projectId}' already exists and cannot be created twice.`, + ); + } + const conflict = requireWorkspaceAvailable(command.workspaceRoot); + if (conflict !== undefined) return Result.fail(conflict); + return Result.succeed({ + ...base, + type: "project.created", + payload: { + projectId: command.projectId, + title: command.title, + workspaceRoot: command.workspaceRoot, + // Project creation has no user model choice. Older clients sent an + // automatic seed, but only a metadata update records an explicit default. + defaultModelSelection: null, + faviconPath: null, + projectIcon: null, + scripts: command.scripts ?? [], + createdAt: occurredAt, + updatedAt: occurredAt, + }, + }); + } + + case "project.meta.update": { + const project = activeProject; + if (project === undefined) return missingProject(); + if ( + command.projectIcon?.kind === "monogram" && + Array.from(monogramSegmenter.segment(command.projectIcon.text)).length > 2 + ) { + return invariant("Project monograms must contain at most two characters."); + } + if (command.scripts !== undefined) { + // Persisted IDs predate shortcut validation. Let users edit or remove them + // without allowing another invalid ID to enter the project. + const existingIds = new Set(project.scripts.map((script) => script.id)); + for (const script of command.scripts) { + if (!existingIds.has(script.id) && !isScriptRunCommand(`script.${script.id}.run`)) { + // The raw ID is unbounded user input and this detail is persisted. + return invariant( + `Script IDs must be 1-${MAX_SCRIPT_ID_LENGTH} lowercase letters, digits or hyphens, starting with a letter or digit (got ${script.id.length} characters).`, + ); + } + } + } + if (command.workspaceRoot !== undefined) { + const conflict = requireWorkspaceAvailable(command.workspaceRoot); + if (conflict !== undefined) return Result.fail(conflict); + } + return Result.succeed({ + ...base, + type: "project.meta-updated", + payload: { + projectId: command.projectId, + ...(command.title === undefined ? {} : { title: command.title }), + ...(command.workspaceRoot === undefined ? {} : { workspaceRoot: command.workspaceRoot }), + ...(command.defaultModelSelection === undefined + ? {} + : { defaultModelSelection: command.defaultModelSelection }), + ...(command.defaultThreadEnvMode === undefined + ? {} + : { defaultThreadEnvMode: command.defaultThreadEnvMode }), + ...(command.autoPull === undefined ? {} : { autoPull: command.autoPull }), + ...(command.faviconPath === undefined ? {} : { faviconPath: command.faviconPath }), + ...(command.projectIcon === undefined ? {} : { projectIcon: command.projectIcon }), + ...(command.scripts === undefined ? {} : { scripts: command.scripts }), + updatedAt: occurredAt, + }, + }); + } + + case "project.delete": { + if (activeProject === undefined) return missingProject(); + // Thread children are deleted by ProjectService before this event commits. + return Result.succeed({ + ...base, + type: "project.deleted", + payload: { projectId: command.projectId, deletedAt: occurredAt }, + }); + } + } +} diff --git a/apps/server/src/orchestration-v2/ProjectSettingsUpgrade.integration.test.ts b/apps/server/src/orchestration-v2/ProjectSettingsUpgrade.integration.test.ts new file mode 100644 index 000000000000..7a9afc38399d --- /dev/null +++ b/apps/server/src/orchestration-v2/ProjectSettingsUpgrade.integration.test.ts @@ -0,0 +1,191 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, it } from "@effect/vitest"; +import { ProjectId } from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import * as CheckpointStore from "../checkpointing/CheckpointStore.ts"; +import { ServerConfig } from "../config.ts"; +import * as GitWorkflow from "../git/GitWorkflowService.ts"; +import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient"; + +import { layer as mcpSessionRegistryTestLayer } from "../mcp/McpSessionRegistry.testkit.ts"; +import { makeSqlitePersistenceLive } from "../persistence/Layers/Sqlite.ts"; +import { runMigrations } from "../persistence/Migrations.ts"; +import { ProjectEnrichmentService } from "../project/ProjectEnrichmentService.ts"; +import * as ProjectService from "../project/ProjectService.ts"; +import { ProviderInstanceRegistry } from "../provider/Services/ProviderInstanceRegistry.ts"; +import { ServerSettingsService } from "../serverSettings.ts"; +import { SourceControlProviderRegistry } from "../sourceControl/SourceControlProviderRegistry.ts"; +import * as VcsDriverRegistry from "../vcs/VcsDriverRegistry.ts"; +import * as VcsProcess from "../vcs/VcsProcess.ts"; +import { WorkspacePaths } from "../workspace/WorkspacePaths.ts"; +import { LegacyV1ThreadImporter } from "./LegacyV1ThreadImporter.ts"; +import { OrchestrationV2LayerLive, ProjectServiceLayerLive } from "./runtimeLayer.ts"; + +const projectId = ProjectId.make("project:upgrade"); +const icon = { kind: "emoji", emoji: "🦊" } as const; +const encodeJson = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown)); + +const expectedSettings = { + default_thread_env_mode: "worktree", + auto_pull: 1, + favicon_path: "brand/icon.svg", + project_icon_json: encodeJson(icon), +}; + +const readSettings = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const rows = yield* sql` + SELECT default_thread_env_mode, auto_pull, favicon_path, project_icon_json + FROM projection_projects + WHERE project_id = ${projectId} + `; + return rows[0]; +}); + +/** A released V1 database at migration 54 whose project carries all four settings. */ +const seedV1Database = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* runMigrations({ toMigrationInclusive: 54 }); + const events = [ + { + type: "project.created", + payload: { + projectId, + title: "Upgrade", + workspaceRoot: "/work/upgrade", + defaultModelSelection: null, + faviconPath: null, + projectIcon: null, + scripts: [], + createdAt: "2026-01-01T00:00:00.000Z", + updatedAt: "2026-01-01T00:00:00.000Z", + }, + }, + { + type: "project.meta-updated", + payload: { + projectId, + defaultThreadEnvMode: "worktree", + autoPull: true, + faviconPath: "brand/icon.svg", + projectIcon: icon, + updatedAt: "2026-01-02T00:00:00.000Z", + }, + }, + ] as const; + for (const [version, event] of events.entries()) { + yield* sql` + INSERT INTO orchestration_events ( + event_id, aggregate_kind, stream_id, stream_version, event_type, occurred_at, + command_id, causation_event_id, correlation_id, actor_kind, payload_json, metadata_json + ) VALUES ( + ${`v1:${event.type}`}, 'project', ${projectId}, ${version}, ${event.type}, + ${event.payload.updatedAt}, ${`command:${event.type}`}, NULL, ${`command:${event.type}`}, + 'client', ${encodeJson(event.payload)}, '{}' + ) + `; + } + // The V1 projector had applied both events before the upgrade. + yield* sql` + INSERT INTO projection_projects ( + project_id, title, workspace_root, default_model_selection_json, default_thread_env_mode, + auto_pull, favicon_path, project_icon_json, scripts_json, created_at, updated_at, deleted_at + ) VALUES ( + ${projectId}, 'Upgrade', '/work/upgrade', NULL, ${expectedSettings.default_thread_env_mode}, + 1, ${expectedSettings.favicon_path}, ${expectedSettings.project_icon_json}, '[]', + '2026-01-01T00:00:00.000Z', '2026-01-02T00:00:00.000Z', NULL + ) + `; + yield* sql` + INSERT INTO projection_state (projector, last_applied_sequence, updated_at) + VALUES ('projection.projects', 2, '2026-01-02T00:00:00.000Z') + `; +}); + +const unusedEnrichment = { + repositoryIdentity: null, + faviconPath: null, + repositoryIdentityResolved: false, +}; + +/** The production V2 runtime and project service against one file-backed database. */ +const makeRuntimeLayer = (dbPath: string) => { + const platform = Layer.merge( + NodeServices.layer, + Layer.mock(SourceControlProviderRegistry)({ resolveLink: () => Effect.die("unused") }), + ); + const serverConfig = ServerConfig.layerTest(process.cwd(), { prefix: "t3-project-upgrade-" }); + const checkpointStore = CheckpointStore.layer.pipe( + Layer.provide( + VcsDriverRegistry.layer.pipe( + Layer.provide(VcsProcess.layer), + Layer.provide(serverConfig), + Layer.provide(platform), + ), + ), + ); + return Layer.mergeAll( + OrchestrationV2LayerLive.pipe(Layer.provide(ProjectServiceLayerLive)), + ProjectServiceLayerLive, + ).pipe( + Layer.provide( + Layer.mock(ProjectEnrichmentService)({ + peek: () => Effect.succeed(unusedEnrichment), + getAvailable: () => Effect.succeed(unusedEnrichment), + invalidate: () => Effect.void, + }), + ), + Layer.provide( + Layer.mock(WorkspacePaths)({ + normalizeWorkspaceRoot: (workspaceRoot) => Effect.succeed(workspaceRoot), + }), + ), + Layer.provide(mcpSessionRegistryTestLayer), + Layer.provideMerge(makeSqlitePersistenceLive(dbPath)), + Layer.provide(checkpointStore), + Layer.provide(serverConfig), + Layer.provide(ServerSettingsService.layerTest()), + Layer.provide( + Layer.succeed(ProviderInstanceRegistry, { + getInstance: () => Effect.succeed(undefined), + listInstances: Effect.succeed([]), + listUnavailable: Effect.succeed([]), + streamChanges: Stream.empty, + subscribeChanges: Effect.never, + }), + ), + Layer.provide( + Layer.mock(GitWorkflow.GitWorkflowService)({ + pruneWorktrees: () => Effect.void, + }), + ), + Layer.provide(platform), + ); +}; + +it.live("keeps project settings through the V2 migrations and the first V2 boot", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const stateDir = yield* fs.makeTempDirectoryScoped({ prefix: "t3-project-upgrade-" }); + const dbPath = path.join(stateDir, "statev2.sqlite"); + + // Seed the released V1 schema, then boot the V2 runtime, which runs 055+, on the same file. + yield* seedV1Database.pipe(Effect.provide(NodeSqliteClient.layer({ filename: dbPath }))); + yield* Effect.gen(function* () { + yield* (yield* LegacyV1ThreadImporter).reconcileShells; + assert.deepEqual(yield* readSettings, expectedSettings); + const project = yield* (yield* ProjectService.ProjectService).getById(projectId); + assert.equal(project._tag, "Some"); + }).pipe(Effect.provide(makeRuntimeLayer(dbPath))); + }), + ).pipe(Effect.provide(NodeServices.layer)), +); diff --git a/apps/server/src/orchestration-v2/ProjectStore.ts b/apps/server/src/orchestration-v2/ProjectStore.ts new file mode 100644 index 000000000000..ce71d3ad12ee --- /dev/null +++ b/apps/server/src/orchestration-v2/ProjectStore.ts @@ -0,0 +1,274 @@ +import { + type ApplicationProjectEvent, + IsoDateTime, + ModelSelection, + type OrchestrationProjectShell, + ProjectIconOverride, + ProjectId, + ProjectScript, + ThreadEnvMode, +} from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; +import * as SqlSchema from "effect/unstable/sql/SqlSchema"; + +export class ProjectStoreV2Error extends Schema.TaggedError()( + "ProjectStoreV2Error", + { + operation: Schema.String, + cause: Schema.Defect(), + }, +) { + override get message(): string { + return `Project store operation '${this.operation}' failed.`; + } +} + +/** One row of `projection_projects`, the durable project read model. */ +export const ProjectRow = Schema.Struct({ + projectId: ProjectId, + title: Schema.String, + workspaceRoot: Schema.String, + defaultModelSelection: Schema.NullOr(ModelSelection), + defaultThreadEnvMode: Schema.NullOr(ThreadEnvMode), + autoPull: Schema.Boolean, + faviconPath: Schema.NullOr(Schema.String), + projectIcon: Schema.NullOr(ProjectIconOverride), + scripts: Schema.Array(ProjectScript), + createdAt: IsoDateTime, + updatedAt: IsoDateTime, + deletedAt: Schema.NullOr(IsoDateTime), +}); +export type ProjectRow = typeof ProjectRow.Type; + +const ProjectDbRow = Schema.Struct({ + ...ProjectRow.fields, + defaultModelSelection: Schema.NullOr(Schema.fromJsonString(ModelSelection)), + autoPull: Schema.BooleanFromBit, + projectIcon: Schema.NullOr(Schema.fromJsonString(ProjectIconOverride)), + scripts: Schema.fromJsonString(Schema.Array(ProjectScript)), +}); + +/** Shell fields without workspace-derived enrichment such as repository identity. */ +function toShell(row: ProjectRow): OrchestrationProjectShell { + return { + id: row.projectId, + title: row.title, + workspaceRoot: row.workspaceRoot, + repositoryIdentity: null, + defaultModelSelection: row.defaultModelSelection, + defaultThreadEnvMode: row.defaultThreadEnvMode, + autoPull: row.autoPull, + faviconPath: row.faviconPath, + projectIcon: row.projectIcon, + scripts: row.scripts, + createdAt: row.createdAt, + updatedAt: row.updatedAt, + }; +} + +export class ProjectStoreV2 extends Context.Service< + ProjectStoreV2, + { + /** Fold one committed project event into its row. Call inside the commit transaction. */ + readonly apply: (event: ApplicationProjectEvent) => Effect.Effect; + readonly get: ( + projectId: ProjectId, + options?: { readonly includeDeleted?: boolean }, + ) => Effect.Effect, ProjectStoreV2Error>; + readonly list: (options?: { + readonly projectIds?: ReadonlyArray; + readonly includeDeleted?: boolean; + }) => Effect.Effect, ProjectStoreV2Error>; + /** Workspace roots match by exact string; callers normalize before asking. */ + readonly findActiveByWorkspaceRoot: ( + workspaceRoot: string, + ) => Effect.Effect, ProjectStoreV2Error>; + readonly getShell: ( + projectId: ProjectId, + ) => Effect.Effect, ProjectStoreV2Error>; + readonly listShells: Effect.Effect< + ReadonlyArray, + ProjectStoreV2Error + >; + } +>()("t3/orchestration-v2/ProjectStore/ProjectStoreV2") {} + +/** @public Service construction is part of the canonical Effect module API. */ +export const make = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const encodeRow = Schema.encodeEffect(ProjectDbRow); + + const selectRows = SqlSchema.findAll({ + Request: Schema.Struct({ + projectId: Schema.optional(ProjectId), + projectIds: Schema.optional(Schema.Array(ProjectId)), + workspaceRoot: Schema.optional(Schema.String), + includeDeleted: Schema.Boolean, + }), + Result: ProjectDbRow, + execute: (request) => sql` + SELECT + project_id AS "projectId", + title, + workspace_root AS "workspaceRoot", + default_model_selection_json AS "defaultModelSelection", + default_thread_env_mode AS "defaultThreadEnvMode", + auto_pull AS "autoPull", + favicon_path AS "faviconPath", + project_icon_json AS "projectIcon", + scripts_json AS "scripts", + created_at AS "createdAt", + updated_at AS "updatedAt", + deleted_at AS "deletedAt" + FROM projection_projects + WHERE ${sql.and([ + ...(request.includeDeleted ? [] : [sql`deleted_at IS NULL`]), + ...(request.projectId === undefined ? [] : [sql`project_id = ${request.projectId}`]), + ...(request.projectIds === undefined ? [] : [sql.in("project_id", request.projectIds)]), + ...(request.workspaceRoot === undefined + ? [] + : [sql`workspace_root = ${request.workspaceRoot}`]), + ])} + ORDER BY created_at ASC, project_id ASC + `, + }); + + const upsertRow = (row: ProjectRow) => + encodeRow(row).pipe( + Effect.flatMap( + (encoded) => sql` + INSERT INTO projection_projects ( + project_id, + title, + workspace_root, + default_model_selection_json, + default_thread_env_mode, + auto_pull, + favicon_path, + project_icon_json, + scripts_json, + created_at, + updated_at, + deleted_at + ) + VALUES ( + ${encoded.projectId}, + ${encoded.title}, + ${encoded.workspaceRoot}, + ${encoded.defaultModelSelection}, + ${encoded.defaultThreadEnvMode}, + ${encoded.autoPull}, + ${encoded.faviconPath}, + ${encoded.projectIcon}, + ${encoded.scripts}, + ${encoded.createdAt}, + ${encoded.updatedAt}, + ${encoded.deletedAt} + ) + ON CONFLICT (project_id) + DO UPDATE SET + title = excluded.title, + workspace_root = excluded.workspace_root, + default_model_selection_json = excluded.default_model_selection_json, + default_thread_env_mode = excluded.default_thread_env_mode, + auto_pull = excluded.auto_pull, + favicon_path = excluded.favicon_path, + project_icon_json = excluded.project_icon_json, + scripts_json = excluded.scripts_json, + created_at = excluded.created_at, + updated_at = excluded.updated_at, + deleted_at = excluded.deleted_at + `, + ), + ); + + const mapError = + (operation: string) => + (effect: Effect.Effect) => + effect.pipe(Effect.mapError((cause) => new ProjectStoreV2Error({ operation, cause }))); + + const get: ProjectStoreV2["Service"]["get"] = (projectId, options) => + selectRows({ projectId, includeDeleted: options?.includeDeleted === true }).pipe( + Effect.map((rows) => Option.fromUndefinedOr(rows[0])), + mapError("get"), + ); + + const list: ProjectStoreV2["Service"]["list"] = (options) => + selectRows({ + ...(options?.projectIds === undefined ? {} : { projectIds: options.projectIds }), + includeDeleted: options?.includeDeleted === true, + }).pipe(mapError("list")); + + const findActiveByWorkspaceRoot: ProjectStoreV2["Service"]["findActiveByWorkspaceRoot"] = ( + workspaceRoot, + ) => + selectRows({ workspaceRoot, includeDeleted: false }).pipe( + Effect.map((rows) => Option.fromUndefinedOr(rows[0])), + mapError("findActiveByWorkspaceRoot"), + ); + + const apply: ProjectStoreV2["Service"]["apply"] = Effect.fn("ProjectStoreV2.apply")( + function* (event) { + if (event.type === "project.created") { + const payload = event.payload; + return yield* upsertRow({ + projectId: payload.projectId, + title: payload.title, + workspaceRoot: payload.workspaceRoot, + defaultModelSelection: payload.defaultModelSelection, + defaultThreadEnvMode: payload.defaultThreadEnvMode ?? null, + autoPull: false, + faviconPath: payload.faviconPath ?? null, + projectIcon: payload.projectIcon ?? null, + scripts: payload.scripts, + createdAt: payload.createdAt, + updatedAt: payload.updatedAt, + deletedAt: null, + }).pipe(mapError("apply")); + } + const existing = yield* get(event.payload.projectId, { includeDeleted: true }); + if (Option.isNone(existing)) return; + const row = existing.value; + if (event.type === "project.deleted") { + return yield* upsertRow({ + ...row, + deletedAt: event.payload.deletedAt, + updatedAt: event.payload.deletedAt, + }).pipe(mapError("apply")); + } + const payload = event.payload; + yield* upsertRow({ + ...row, + ...(payload.title === undefined ? {} : { title: payload.title }), + ...(payload.workspaceRoot === undefined ? {} : { workspaceRoot: payload.workspaceRoot }), + ...(payload.defaultModelSelection === undefined + ? {} + : { defaultModelSelection: payload.defaultModelSelection }), + ...(payload.defaultThreadEnvMode === undefined + ? {} + : { defaultThreadEnvMode: payload.defaultThreadEnvMode }), + ...(payload.autoPull === undefined ? {} : { autoPull: payload.autoPull }), + ...(payload.faviconPath === undefined ? {} : { faviconPath: payload.faviconPath }), + ...(payload.projectIcon === undefined ? {} : { projectIcon: payload.projectIcon }), + ...(payload.scripts === undefined ? {} : { scripts: payload.scripts }), + updatedAt: payload.updatedAt, + }).pipe(mapError("apply")); + }, + ); + + return ProjectStoreV2.of({ + apply, + get, + list, + findActiveByWorkspaceRoot, + getShell: (projectId) => get(projectId).pipe(Effect.map(Option.map(toShell))), + listShells: list().pipe(Effect.map((rows) => rows.map(toShell))), + }); +}); + +export const layer = Layer.effect(ProjectStoreV2, make); diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 40008f8a22b6..5ecb3e560ba8 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -15,6 +15,7 @@ import { RuntimeRequestId, TurnItemId, type ModelSelection, + type OrchestrationV2Run, ProjectId, ProviderDriverKind, ProviderInstanceId, @@ -37,9 +38,7 @@ import * as SqlClient from "effect/unstable/sql/SqlClient"; import * as CheckpointStore from "../checkpointing/CheckpointStore.ts"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import { ServerConfig } from "../config.ts"; -import { OrchestrationEngineService } from "../orchestration/Services/OrchestrationEngine.ts"; -import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts"; -import { OrchestrationLayerLive } from "../orchestration/runtimeLayer.ts"; +import { OrchestrationEventInfrastructureLayerLive } from "../orchestration/runtimeLayer.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import { ProjectionProjectRepositoryLive } from "../persistence/Layers/ProjectionProjects.ts"; import { OrchestrationEventStore } from "../persistence/Services/OrchestrationEventStore.ts"; @@ -330,10 +329,17 @@ it.layer(ProjectDeletionTestLayer)("project deletion during thread commands", (i ); }); -const SharedApplicationDataPlaneTestLayer = Layer.merge( - OrchestrationLayerLive, - OrchestrationV2LayerLive, +const SharedApplicationDataPlaneTestLayer = Layer.mergeAll( + OrchestrationV2LayerLive.pipe(Layer.provide(ProjectServiceLayerLive)), + ProjectServiceLayerLive, + OrchestrationV2EventSinkLayerLive, + OrchestrationEventInfrastructureLayerLive, ).pipe( + Layer.provide( + Layer.mock(WorkspacePaths)({ + normalizeWorkspaceRoot: (workspaceRoot) => Effect.succeed(workspaceRoot), + }), + ), Layer.provide( Layer.succeed(ProjectEnrichmentService, { peek: () => @@ -360,7 +366,6 @@ const SharedApplicationDataPlaneTestLayer = Layer.merge( Layer.provide(ServerSettingsService.layerTest()), Layer.provide(TestProviderInstanceRegistry), Layer.provide(GitWorkflowTestLayer), - Layer.provide(ProjectServiceTestLayer), Layer.provide(PlatformTestLayer), ); @@ -3060,22 +3065,18 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { it.layer(SharedApplicationDataPlaneTestLayer)("pending provider interruption", (it) => { it.effect("interrupts a pending provider start without launching provider work", () => Effect.gen(function* () { - const applicationEngine = yield* OrchestrationEngineService; + const projects = yield* ProjectService.ProjectService; const orchestrator = yield* OrchestratorV2; const threadManagement = yield* ThreadManagementService; const effectWorker = yield* OrchestrationEffectWorkerV2; const projectId = ProjectId.make("runtime-layer-pending-interrupt-project"); const threadId = ThreadId.make("runtime-layer-pending-interrupt-thread"); - yield* applicationEngine.dispatch({ - type: "project.create", + yield* projects.create({ commandId: CommandId.make("runtime-layer-pending-interrupt-project-create"), projectId, title: "Pending interrupt project", workspaceRoot: "/tmp/runtime-layer-pending-interrupt-project", - defaultModelSelection: modelSelection, - scripts: [], - createdAt: "2026-06-22T00:00:00.000Z", }); yield* orchestrator.dispatch({ type: "thread.create", @@ -3138,21 +3139,17 @@ it.layer(SharedApplicationDataPlaneTestLayer)("pending provider interruption", ( it.layer(SharedApplicationDataPlaneTestLayer)("snooze projection", (it) => { it.effect("carries snooze state through the V2 shell projection", () => Effect.gen(function* () { - const applicationEngine = yield* OrchestrationEngineService; + const projects = yield* ProjectService.ProjectService; const orchestrator = yield* OrchestratorV2; const projectId = ProjectId.make("runtime-layer-snoozed-project"); const threadId = ThreadId.make("runtime-layer-snoozed-thread"); const snoozedUntil = "2099-07-25T09:00:00.000Z"; - yield* applicationEngine.dispatch({ - type: "project.create", + yield* projects.create({ commandId: CommandId.make("runtime-layer-snoozed-project-create"), projectId, title: "Snoozed shell projection", workspaceRoot: "/tmp/runtime-layer-snoozed-project", - defaultModelSelection: modelSelection, - scripts: [], - createdAt: "2026-07-24T00:00:00.000Z", }); yield* orchestrator.dispatch({ type: "thread.create", @@ -3217,21 +3214,17 @@ it.layer(SharedApplicationDataPlaneTestLayer)("snooze projection", (it) => { it.layer(SharedApplicationDataPlaneTestLayer)("visited projection", (it) => { it.effect("carries the visited watermark through the V2 shell projection", () => Effect.gen(function* () { - const applicationEngine = yield* OrchestrationEngineService; + const projects = yield* ProjectService.ProjectService; const orchestrator = yield* OrchestratorV2; const projectId = ProjectId.make("runtime-layer-visited-project"); const threadId = ThreadId.make("runtime-layer-visited-thread"); const visitedAt = "2026-07-24T01:00:00.000Z"; - yield* applicationEngine.dispatch({ - type: "project.create", + yield* projects.create({ commandId: CommandId.make("runtime-layer-visited-project-create"), projectId, title: "Visited shell projection", workspaceRoot: "/tmp/runtime-layer-visited-project", - defaultModelSelection: modelSelection, - scripts: [], - createdAt: "2026-07-24T00:00:00.000Z", }); yield* orchestrator.dispatch({ type: "thread.create", @@ -3329,27 +3322,22 @@ it.layer(SharedApplicationDataPlaneTestLayer)("visited projection", (it) => { it.layer(SharedApplicationDataPlaneTestLayer)("shared application data plane", (it) => { it.effect("orders retained project transactions and V2 thread transactions in one source", () => Effect.gen(function* () { - const applicationEngine = yield* OrchestrationEngineService; + const projects = yield* ProjectService.ProjectService; const applicationEvents = yield* OrchestrationEventStore; const orchestrator = yield* OrchestratorV2; - const projectionSnapshot = yield* ProjectionSnapshotQuery; const sql = yield* SqlClient.SqlClient; const projectId = ProjectId.make("runtime-layer-shared-project"); const threadId = ThreadId.make("runtime-layer-shared-thread"); - const projectCommand = { - type: "project.create" as const, + const projectInput = { commandId: CommandId.make("runtime-layer-shared-project-create"), projectId, title: "Shared application source", workspaceRoot: "/tmp/runtime-layer-shared-project", - defaultModelSelection: modelSelection, - scripts: [], - createdAt: "2026-06-20T00:00:00.000Z", }; - const projectResult = yield* applicationEngine.dispatch(projectCommand); - const projectRetry = yield* applicationEngine.dispatch(projectCommand); - assert.equal(projectRetry.sequence, projectResult.sequence); + const created = yield* projects.create(projectInput); + const retried = yield* projects.create(projectInput); + assert.deepEqual(retried, created); const delivered = yield* Queue.unbounded(); yield* applicationEvents.streamApplicationEvents().pipe( @@ -3359,7 +3347,7 @@ it.layer(SharedApplicationDataPlaneTestLayer)("shared application data plane", ( ); const projectEvent = yield* Queue.take(delivered); - assert.equal(projectEvent.sequence, projectResult.sequence); + assert.isTrue("aggregateKind" in projectEvent && projectEvent.aggregateId === projectId); const threadResult = yield* orchestrator.dispatch({ type: "thread.create", @@ -3381,7 +3369,7 @@ it.layer(SharedApplicationDataPlaneTestLayer)("shared application data plane", ( assert.isAbove(threadEvent.sequence, projectEvent.sequence); assert.isTrue("aggregateKind" in projectEvent); assert.isTrue("event" in threadEvent); - assert.equal((yield* projectionSnapshot.getProjectShellById(projectId))._tag, "Some"); + assert.equal((yield* projects.getById(projectId))._tag, "Some"); const retainedReceipts = yield* sql<{ readonly aggregate_kind: string; @@ -3465,9 +3453,22 @@ it.layer(TestLayer)("usage-limit recovery", (it) => { createdBy: "user", creationSource: "web", }); - const source = (yield* orchestrator.getThreadProjection(threadId)).runs[0]!; + const [source, queued] = (yield* orchestrator.getThreadProjection(threadId)).runs as [ + OrchestrationV2Run, + OrchestrationV2Run, + ]; + // A user stop holds the queue as the run ends (thread.turn.interrupt with + // holdQueue). Without the hold, the terminal-run worker may start the + // queued run before the resume below, depending on fiber scheduling. yield* events.write({ events: [ + { + id: EventId.make(`manual-resume:hold:${reason}`), + type: "run.updated", + threadId, + occurredAt: now, + payload: { ...queued, queueHeld: true }, + }, { id: EventId.make(`manual-resume:stop:${reason}`), type: "run.updated", diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index b4de5bb5d633..9f1fab6d85a2 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -1,10 +1,7 @@ import * as UsageLimitRecoveryWorker from "./UsageLimitRecoveryWorker.ts"; import * as Scheduler from "../scheduling/Scheduler.ts"; import * as Layer from "effect/Layer"; -import { - OrchestrationEventInfrastructureLayerLive, - OrchestrationLayerLive, -} from "../orchestration/runtimeLayer.ts"; +import { OrchestrationEventInfrastructureLayerLive } from "../orchestration/runtimeLayer.ts"; import { ProjectionProjectRepositoryLive } from "../persistence/Layers/ProjectionProjects.ts"; import { layer as providerSessionRuntimeLayer } from "../persistence/ProviderSessionRuntime.ts"; import * as TextGeneration from "../textGeneration/TextGeneration.ts"; @@ -31,6 +28,7 @@ import { layer as legacyV1ThreadImporterLayer } from "./LegacyV1ThreadImporter.t import { layer as orchestratorLayer } from "./Orchestrator.ts"; import { layer as projectionStoreLayer } from "./ProjectionStore.ts"; import { layer as projectionMaintenanceLayer } from "./ProjectionMaintenance.ts"; +import * as ProjectStore from "./ProjectStore.ts"; import { layerFromProviderInstanceRegistry as providerAdapterRegistryLayerFromProviderInstances } from "./ProviderAdapterRegistry.ts"; import { layer as providerContinuationRequestsLayer } from "./ProviderContinuationRequests.ts"; import { workerLive as providerContinuationWorkerLive } from "./ProviderContinuationService.ts"; @@ -67,6 +65,7 @@ const storesLayer = Layer.mergeAll( OrchestrationEventInfrastructureLayerLive, eventStoreProvided, projectionStoreLayer, + ProjectStore.layer, commandReceiptStoreProvided, effectOutboxLayer, turnItemPositionStoreLayer, @@ -82,8 +81,7 @@ const legacyV1ThreadImporterProvided = legacyV1ThreadImporterLayer.pipe( export const ProjectServiceLayerLive = projectServiceLayer.pipe( Layer.provide( Layer.mergeAll( - ProjectionProjectRepositoryLive, - OrchestrationLayerLive, + ProjectStore.layer, projectionStoreLayer, eventSinkProvided, idAllocatorLayer, @@ -302,4 +300,7 @@ export const OrchestrationV2ProductionLayerLive = Layer.mergeAll( ), providerContinuationWorkerProvided, agentSessionImporterProvided, -).pipe(Layer.provide(Scheduler.layer), Layer.provideMerge(OrchestrationLayerLive)); +).pipe( + Layer.provide(Scheduler.layer), + Layer.provideMerge(OrchestrationEventInfrastructureLayerLive), +); diff --git a/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts b/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts index 54114648ebca..633a76ad6425 100644 --- a/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts +++ b/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts @@ -38,6 +38,7 @@ import { layer as eventStoreLayer } from "../EventStore.ts"; import { layer as idAllocatorLayer } from "../IdAllocator.ts"; import { layer as orchestratorLayer } from "../Orchestrator.ts"; import { layer as projectionStoreLayer } from "../ProjectionStore.ts"; +import { layer as projectStoreLayer } from "../ProjectStore.ts"; import { OrchestratorV2, type OrchestratorV2Error } from "../Orchestrator.ts"; import { ProviderAdapterRegistryV2 } from "../ProviderAdapterRegistry.ts"; import { ProviderAuthService } from "../../provider/Services/ProviderAuthService.ts"; @@ -285,6 +286,7 @@ export function makeOrchestratorV2ReplayLayerWithRegistry( const storesLayer = Layer.mergeAll( eventStoreLayer, projectionStoreLayer, + projectStoreLayer, commandReceiptStoreLayer, effectOutboxLayer, turnItemPositionStoreLayer, diff --git a/apps/server/src/orchestration/Errors.ts b/apps/server/src/orchestration/Errors.ts index 5e6376072ce1..144d22753613 100644 --- a/apps/server/src/orchestration/Errors.ts +++ b/apps/server/src/orchestration/Errors.ts @@ -33,7 +33,6 @@ export const OrchestrationCommandRejection = Schema.Union([ OrchestrationThreadSettleBlockedError, ]); export type OrchestrationCommandRejection = typeof OrchestrationCommandRejection.Type; -export const isOrchestrationCommandRejection = Schema.is(OrchestrationCommandRejection); export class OrchestrationCommandPreviouslyRejectedError extends Schema.TaggedError()( "OrchestrationCommandPreviouslyRejectedError", diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts deleted file mode 100644 index 2fac43d9acf4..000000000000 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts +++ /dev/null @@ -1,389 +0,0 @@ -import type { ProjectId } from "@t3tools/contracts"; -import type { - OrchestrationClientOrigin, - OrchestrationEvent, - OrchestrationReadModel, - ProjectOrchestrationCommand, -} from "@t3tools/contracts/legacy-orchestration"; -import * as Cause from "effect/Cause"; -import * as Clock from "effect/Clock"; -import * as Crypto from "effect/Crypto"; -import * as DateTime from "effect/DateTime"; -import * as Deferred from "effect/Deferred"; -import * as Duration from "effect/Duration"; -import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; -import * as Layer from "effect/Layer"; -import * as Metric from "effect/Metric"; -import * as Option from "effect/Option"; -import * as PubSub from "effect/PubSub"; -import * as Queue from "effect/Queue"; -import * as Schema from "effect/Schema"; -import * as Stream from "effect/Stream"; -import * as SqlClient from "effect/unstable/sql/SqlClient"; - -import { - metricAttributes, - orchestrationCommandAckDuration, - orchestrationCommandsTotal, - orchestrationCommandDuration, -} from "../../observability/Metrics.ts"; -import { toPersistenceSqlError } from "../../persistence/Errors.ts"; -import { OrchestrationEventStore } from "../../persistence/Services/OrchestrationEventStore.ts"; -import { OrchestrationCommandReceiptRepository } from "../../persistence/Services/OrchestrationCommandReceipts.ts"; -import { - isOrchestrationCommandRejection, - OrchestrationCommandIdConflictError, - OrchestrationCommandInvariantError, - OrchestrationCommandPreviouslyRejectedError, - type OrchestrationDispatchError, - type OrchestrationProjectorDecodeError, -} from "../Errors.ts"; -import { decideOrchestrationCommand } from "../decider.ts"; -import { createEmptyReadModel, projectEvent } from "../projector.ts"; -import { OrchestrationProjectionPipeline } from "../Services/ProjectionPipeline.ts"; -import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; -import { ThreadBackgroundLivenessService } from "../ThreadBackgroundLiveness.ts"; -import { - OrchestrationEngineService, - type OrchestrationEngineShape, -} from "../Services/OrchestrationEngine.ts"; -const isOrchestrationCommandPreviouslyRejectedError = Schema.is( - OrchestrationCommandPreviouslyRejectedError, -); -const isOrchestrationCommandIdConflictError = Schema.is(OrchestrationCommandIdConflictError); - -interface CommandEnvelope { - command: ProjectOrchestrationCommand; - origin: OrchestrationClientOrigin | undefined; - result: Deferred.Deferred<{ sequence: number }, OrchestrationDispatchError>; - startedAtMs: number; -} - -function commandToAggregateRef(command: ProjectOrchestrationCommand): { - readonly aggregateKind: "project"; - readonly aggregateId: ProjectId; -} { - return { aggregateKind: "project", aggregateId: command.projectId }; -} - -const makeOrchestrationEngine = Effect.gen(function* () { - const sql = yield* SqlClient.SqlClient; - const eventStore = yield* OrchestrationEventStore; - const commandReceiptRepository = yield* OrchestrationCommandReceiptRepository; - const projectionPipeline = yield* OrchestrationProjectionPipeline; - const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; - const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService; - const crypto = yield* Crypto.Crypto; - - const nowIso = Effect.map(DateTime.now, DateTime.formatIso); - let commandReadModel = createEmptyReadModel(yield* nowIso); - - const commandQueue = yield* Queue.unbounded(); - const eventPubSub = yield* PubSub.unbounded(); - - const projectEventsOntoReadModel = ( - baseReadModel: OrchestrationReadModel, - events: ReadonlyArray, - ): Effect.Effect => - Effect.gen(function* () { - let nextReadModel = baseReadModel; - for (const event of events) { - nextReadModel = yield* projectEvent(nextReadModel, event); - } - return nextReadModel; - }); - - const processEnvelope = (envelope: CommandEnvelope): Effect.Effect => { - const dispatchStartSequence = commandReadModel.snapshotSequence; - let processingStartedAtMs = 0; - const aggregateRef = commandToAggregateRef(envelope.command); - const baseMetricAttributes = { - commandType: envelope.command.type, - aggregateKind: aggregateRef.aggregateKind, - } as const; - const reconcileReadModelAfterDispatchFailure = Effect.gen(function* () { - const persistedEvents = yield* Stream.runCollect( - eventStore.readFromSequence(dispatchStartSequence), - ).pipe(Effect.map((chunk): OrchestrationEvent[] => Array.from(chunk))); - if (persistedEvents.length === 0) { - return; - } - - commandReadModel = yield* projectEventsOntoReadModel(commandReadModel, persistedEvents); - - yield* eventStore.publishCommitted( - persistedEvents.filter( - (event) => - event.type === "project.created" || - event.type === "project.meta-updated" || - event.type === "project.deleted", - ), - ); - - for (const persistedEvent of persistedEvents) { - yield* PubSub.publish(eventPubSub, persistedEvent); - } - }); - - return Effect.exit( - Effect.gen(function* () { - processingStartedAtMs = yield* Clock.currentTimeMillis; - yield* Effect.annotateCurrentSpan({ - "orchestration.command_id": envelope.command.commandId, - "orchestration.command_type": envelope.command.type, - "orchestration.aggregate_kind": aggregateRef.aggregateKind, - "orchestration.aggregate_id": aggregateRef.aggregateId, - }); - - const existingReceipt = yield* commandReceiptRepository.getByCommandId({ - commandId: envelope.command.commandId, - }); - if (Option.isSome(existingReceipt)) { - // A receipt only proves this exact command was handled. Replaying it - // for a command aimed at another aggregate would report success for - // work that never happened. - if ( - existingReceipt.value.aggregateKind !== aggregateRef.aggregateKind || - existingReceipt.value.aggregateId !== aggregateRef.aggregateId - ) { - return yield* new OrchestrationCommandIdConflictError({ - commandId: envelope.command.commandId, - receiptAggregateKind: existingReceipt.value.aggregateKind, - receiptAggregateId: existingReceipt.value.aggregateId, - commandAggregateKind: aggregateRef.aggregateKind, - commandAggregateId: aggregateRef.aggregateId, - }); - } - if (existingReceipt.value.status === "accepted") { - return { - sequence: existingReceipt.value.resultSequence, - }; - } - return yield* new OrchestrationCommandPreviouslyRejectedError({ - commandId: envelope.command.commandId, - detail: existingReceipt.value.error ?? "Previously rejected.", - }); - } - - const eventBase = yield* decideOrchestrationCommand({ - command: envelope.command, - readModel: commandReadModel, - }).pipe( - Effect.provideService(Crypto.Crypto, crypto), - Effect.mapError((cause) => - isOrchestrationCommandRejection(cause) - ? cause - : new OrchestrationCommandInvariantError({ - commandType: envelope.command.type, - detail: "Failed to generate an event identifier.", - cause, - }), - ), - ); - const plannedEvents = Array.isArray(eventBase) ? eventBase : [eventBase]; - // Stamp the dispatching client's origin onto every event the command - // produced. The decider stays pure; attribution is an engine concern. - const eventBases = - envelope.origin === undefined - ? plannedEvents - : plannedEvents.map((planned) => ({ - ...planned, - metadata: { ...planned.metadata, origin: envelope.origin }, - })); - const committedCommand = yield* sql - .withTransaction( - Effect.gen(function* () { - const committedEvents: OrchestrationEvent[] = []; - const attachmentCleanups: Effect.Effect[] = []; - let nextCommandReadModel = commandReadModel; - - for (const nextEvent of eventBases) { - const savedEvent = yield* eventStore.append(nextEvent); - nextCommandReadModel = yield* projectEvent(nextCommandReadModel, savedEvent); - const cleanup = yield* projectionPipeline.projectEventDeferred(savedEvent); - attachmentCleanups.push(cleanup); - committedEvents.push(savedEvent); - } - - const lastSavedEvent = committedEvents.at(-1) ?? null; - if (lastSavedEvent === null) { - return yield* new OrchestrationCommandInvariantError({ - commandType: envelope.command.type, - detail: "Command produced no events.", - }); - } - - yield* commandReceiptRepository.upsert({ - commandId: envelope.command.commandId, - aggregateKind: lastSavedEvent.aggregateKind, - aggregateId: lastSavedEvent.aggregateId, - commandType: envelope.command.type, - acceptedAt: lastSavedEvent.occurredAt, - resultSequence: lastSavedEvent.sequence, - status: "accepted", - error: null, - }); - - return { - committedEvents, - attachmentCleanups, - lastSequence: lastSavedEvent.sequence, - nextCommandReadModel, - } as const; - }), - ) - .pipe( - Effect.catchTag("SqlError", (sqlError) => - Effect.fail( - toPersistenceSqlError("OrchestrationEngine.processEnvelope:transaction")(sqlError), - ), - ), - ); - - commandReadModel = committedCommand.nextCommandReadModel; - for (const cleanup of committedCommand.attachmentCleanups) { - yield* cleanup; - } - yield* eventStore.publishCommitted( - committedCommand.committedEvents.filter( - (event) => - event.type === "project.created" || - event.type === "project.meta-updated" || - event.type === "project.deleted", - ), - ); - for (const [index, event] of committedCommand.committedEvents.entries()) { - yield* PubSub.publish(eventPubSub, event); - if (index === 0) { - yield* Metric.update( - Metric.withAttributes( - orchestrationCommandAckDuration, - metricAttributes({ - ...baseMetricAttributes, - ackEventType: event.type, - }), - ), - Duration.millis(Math.max(0, (yield* Clock.currentTimeMillis) - envelope.startedAtMs)), - ); - } - } - return { sequence: committedCommand.lastSequence }; - }).pipe(Effect.withSpan(`orchestration.command.${envelope.command.type}`)), - ).pipe( - Effect.flatMap((exit) => - Effect.gen(function* () { - const outcome = Exit.isSuccess(exit) - ? "success" - : Cause.hasInterruptsOnly(exit.cause) - ? "interrupt" - : "failure"; - yield* Metric.update( - Metric.withAttributes( - orchestrationCommandDuration, - metricAttributes(baseMetricAttributes), - ), - Duration.millis(Math.max(0, (yield* Clock.currentTimeMillis) - processingStartedAtMs)), - ); - yield* Metric.update( - Metric.withAttributes( - orchestrationCommandsTotal, - metricAttributes({ - ...baseMetricAttributes, - outcome, - }), - ), - 1, - ); - - if (Exit.isSuccess(exit)) { - yield* Deferred.succeed(envelope.result, exit.value); - return; - } - - const error = Cause.squash(exit.cause) as OrchestrationDispatchError; - if ( - !isOrchestrationCommandPreviouslyRejectedError(error) && - !isOrchestrationCommandIdConflictError(error) - ) { - yield* reconcileReadModelAfterDispatchFailure.pipe( - Effect.catch(() => - Effect.logWarning( - "failed to reconcile orchestration read model after dispatch failure", - ).pipe( - Effect.annotateLogs({ - commandId: envelope.command.commandId, - snapshotSequence: commandReadModel.snapshotSequence, - }), - ), - ), - ); - - if (isOrchestrationCommandRejection(error)) { - yield* commandReceiptRepository - .upsert({ - commandId: envelope.command.commandId, - aggregateKind: aggregateRef.aggregateKind, - aggregateId: aggregateRef.aggregateId, - commandType: envelope.command.type, - acceptedAt: yield* nowIso, - resultSequence: commandReadModel.snapshotSequence, - status: "rejected", - error: error.message, - }) - .pipe(Effect.ignore); - } - } - - yield* Deferred.fail(envelope.result, error); - }), - ), - ); - }; - - yield* projectionPipeline.bootstrap; - commandReadModel = yield* projectionSnapshotQuery.getCommandReadModel(); - - const worker = Effect.forever(Queue.take(commandQueue).pipe(Effect.flatMap(processEnvelope))); - yield* Effect.forkScoped(worker); - yield* Effect.logDebug("orchestration engine started").pipe( - Effect.annotateLogs({ sequence: commandReadModel.snapshotSequence }), - ); - - const readEvents: OrchestrationEngineShape["readEvents"] = (fromSequenceExclusive, limit) => - eventStore.readFromSequence(fromSequenceExclusive, limit); - - const dispatch: OrchestrationEngineShape["dispatch"] = (command, options) => - Effect.gen(function* () { - const result = yield* Deferred.make<{ sequence: number }, OrchestrationDispatchError>(); - yield* Queue.offer(commandQueue, { - command, - origin: options?.origin, - result, - startedAtMs: yield* Clock.currentTimeMillis, - }); - return yield* Deferred.await(result); - }); - - return { - readEvents, - dispatch, - subscribeDomainEvents: PubSub.subscribe(eventPubSub).pipe(Effect.map(Stream.fromSubscription)), - // Each access creates a fresh PubSub subscription so that multiple - // consumers (wsServer, ProviderRuntimeIngestion, CheckpointReactor, etc.) - // each independently receive all domain events. - get streamDomainEvents(): OrchestrationEngineShape["streamDomainEvents"] { - return Stream.fromPubSub(eventPubSub); - }, - // The command read model's snapshotSequence tracks the latest committed - // event sequence (updated on the worker fiber). A plain property read is a - // consistent, committed value — reassignment of `commandReadModel` is - // atomic on the single-threaded event loop. - latestSequence: Effect.sync(() => commandReadModel.snapshotSequence), - } satisfies OrchestrationEngineShape; -}); - -export const OrchestrationEngineLive = Layer.effect( - OrchestrationEngineService, - makeOrchestrationEngine, -); diff --git a/apps/server/src/orchestration/Services/OrchestrationEngine.ts b/apps/server/src/orchestration/Services/OrchestrationEngine.ts deleted file mode 100644 index 190768593dd7..000000000000 --- a/apps/server/src/orchestration/Services/OrchestrationEngine.ts +++ /dev/null @@ -1,100 +0,0 @@ -/** - * Historical name for the application event-sourcing engine. - * - * This is not the agent orchestrator. It retains serialized project-command - * validation, append, receipt, and projection transactions. Agent execution is - * owned by orchestration V2, whose thread events use the same event store. - * - * Uses Effect `Context.Service` for dependency injection. Command dispatch, - * replay, and unknown-input decoding all return typed domain errors. - * - * @module OrchestrationEngineService - */ -import type { - OrchestrationClientOrigin, - OrchestrationEvent, - ProjectOrchestrationCommand, -} from "@t3tools/contracts/legacy-orchestration"; -import * as Context from "effect/Context"; -import type * as Effect from "effect/Effect"; -import type * as Scope from "effect/Scope"; -import type * as Stream from "effect/Stream"; - -import type { OrchestrationDispatchError } from "../Errors.ts"; -import type { OrchestrationEventStoreError } from "../../persistence/Errors.ts"; - -/** - * OrchestrationEngineShape - Service API for orchestration command and event flow. - */ -export interface OrchestrationEngineShape { - /** - * Replay persisted orchestration events from an exclusive sequence cursor. - * - * @param fromSequenceExclusive - Sequence cursor (exclusive). - * @param limit - Maximum number of events to read. Defaults to the event - * store's page-bounded default; pass a higher value when the caller must - * read every event after the cursor (e.g. per-thread catch-up that filters - * a small subset out of a potentially larger global range). - * @returns Stream containing ordered events. - */ - readonly readEvents: ( - fromSequenceExclusive: number, - limit?: number, - ) => Stream.Stream; - - /** - * Dispatch a validated orchestration command. - * - * @param command - Valid orchestration command. - * @param options - Optional client origin (surface/app version) stamped into - * the metadata of every event the command produces. - * @returns Effect containing the sequence of the persisted event. - * - * Dispatch is serialized through an internal queue and deduplicated via - * command receipts. - */ - readonly dispatch: ( - command: ProjectOrchestrationCommand, - options?: { readonly origin?: OrchestrationClientOrigin }, - ) => Effect.Effect<{ sequence: number }, OrchestrationDispatchError, never>; - - /** - * Stream persisted domain events in dispatch order. - * - * This is a hot runtime stream (new events only), not a historical replay. - */ - readonly streamDomainEvents: Stream.Stream; - - /** - * Acquire a domain-event subscription before starting a consumer. - * The subscription is ready when this effect returns and closes with the scope. - */ - readonly subscribeDomainEvents: Effect.Effect< - Stream.Stream, - never, - Scope.Scope - >; - - /** - * The latest sequence reflected in the engine's authoritative command read - * model (0 if none). Used to gauge how far behind a resuming client is before - * choosing between an incremental replay and a fresh projected snapshot. - */ - readonly latestSequence: Effect.Effect; -} - -/** - * OrchestrationEngineService - Service tag for orchestration engine access. - * - * @example - * ```ts - * const program = Effect.gen(function* () { - * const engine = yield* OrchestrationEngineService - * return yield* engine.dispatch(command) - * }) - * ``` - */ -export class OrchestrationEngineService extends Context.Service< - OrchestrationEngineService, - OrchestrationEngineShape ->()("t3/orchestration/Services/OrchestrationEngine/OrchestrationEngineService") {} diff --git a/apps/server/src/orchestration/runtimeLayer.ts b/apps/server/src/orchestration/runtimeLayer.ts index 950585e8dd70..f009641132af 100644 --- a/apps/server/src/orchestration/runtimeLayer.ts +++ b/apps/server/src/orchestration/runtimeLayer.ts @@ -2,7 +2,6 @@ import * as Layer from "effect/Layer"; import { OrchestrationCommandReceiptRepositoryLive } from "../persistence/Layers/OrchestrationCommandReceipts.ts"; import { OrchestrationEventStoreLive } from "../persistence/Layers/OrchestrationEventStore.ts"; -import { OrchestrationEngineLive } from "./Layers/OrchestrationEngine.ts"; import { OrchestrationProjectionPipelineLive } from "./Layers/ProjectionPipeline.ts"; import { OrchestrationProjectionSnapshotQueryLive } from "./Layers/ProjectionSnapshotQuery.ts"; import * as ThreadBackgroundLiveness from "./ThreadBackgroundLiveness.ts"; @@ -29,8 +28,3 @@ export const OrchestrationInfrastructureLayerLive = Layer.mergeAll( Layer.provideMerge(ThreadBackgroundLiveness.layer), Layer.provideMerge(ThreadPlanProgress.layer), ); - -export const OrchestrationLayerLive = Layer.mergeAll( - OrchestrationInfrastructureLayerLive, - OrchestrationEngineLive.pipe(Layer.provide(OrchestrationInfrastructureLayerLive)), -); diff --git a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts index 0c427609be0d..d92a275e588c 100644 --- a/apps/server/src/persistence/Layers/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Layers/OrchestrationEventStore.ts @@ -323,6 +323,33 @@ const makeEventStore = Effect.gen(function* () { ), ); + const appendProjectEvent: OrchestrationEventStoreShape["appendProjectEvent"] = (event) => + appendEventRow({ + eventId: event.eventId, + aggregateKind: "project", + streamId: event.aggregateId, + type: event.type, + causationEventId: event.causationEventId, + correlationId: event.correlationId, + actorKind: inferActorKind(event), + occurredAt: event.occurredAt, + commandId: event.commandId, + payloadJson: + event.type === "project.deleted" || !event.payload.projectIcon + ? event.payload + : { ...event.payload, projectIcon: encodeProjectIcon(event.payload.projectIcon) }, + metadataJson: event.metadata, + applicationEventVersion: 2, + }).pipe( + Effect.flatMap((row) => decodeProjectEvent(row)), + Effect.mapError( + toPersistenceSqlOrDecodeError( + "OrchestrationEventStore.appendProjectEvent:insert", + "OrchestrationEventStore.appendProjectEvent:decode", + ), + ), + ); + const readFromSequence: OrchestrationEventStoreShape["readFromSequence"] = ( sequenceExclusive, limit = DEFAULT_READ_FROM_SEQUENCE_LIMIT, @@ -688,6 +715,7 @@ const makeEventStore = Effect.gen(function* () { append, readFromSequence, readAll: () => readFromSequence(0, Number.MAX_SAFE_INTEGER), + appendProjectEvent, appendAgentEvents, readAgentEvents, getAgentReplayStats, diff --git a/apps/server/src/persistence/Services/OrchestrationEventStore.ts b/apps/server/src/persistence/Services/OrchestrationEventStore.ts index f6c2297badbf..1941ef34f718 100644 --- a/apps/server/src/persistence/Services/OrchestrationEventStore.ts +++ b/apps/server/src/persistence/Services/OrchestrationEventStore.ts @@ -11,6 +11,7 @@ * @module OrchestrationEventStore */ import type { + ApplicationProjectEvent, ApplicationStoredEvent, CommandId, OrchestrationV2DomainEvent, @@ -24,6 +25,13 @@ import type * as Stream from "effect/Stream"; import type { OrchestrationEventStoreError } from "../Errors.ts"; +/** A project event before the store assigns its sequence. */ +export type UnsequencedProjectEvent = ApplicationProjectEvent extends infer Event + ? Event extends ApplicationProjectEvent + ? Omit + : never + : never; + /** * OrchestrationEventStoreShape - Service API for orchestration event persistence. */ @@ -61,6 +69,11 @@ export interface OrchestrationEventStoreShape { */ readonly readAll: () => Stream.Stream; + /** Append one project event to the shared application log. */ + readonly appendProjectEvent: ( + event: UnsequencedProjectEvent, + ) => Effect.Effect; + /** Append V2 agent events to the same globally ordered application log. */ readonly appendAgentEvents: (input: { readonly commandId?: CommandId; diff --git a/apps/server/src/project/ProjectService.deletion.test.ts b/apps/server/src/project/ProjectService.deletion.test.ts index 01825da3348c..dfa3312df276 100644 --- a/apps/server/src/project/ProjectService.deletion.test.ts +++ b/apps/server/src/project/ProjectService.deletion.test.ts @@ -17,7 +17,6 @@ import { TestClock } from "effect/testing"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { ServerConfig } from "../config.ts"; -import { OrchestrationLayerLive } from "../orchestration/runtimeLayer.ts"; import { OrchestrationEffectRequestV2 } from "../orchestration-v2/EffectOutbox.ts"; import { EventSinkV2, @@ -26,7 +25,7 @@ import { layer as eventSinkLayer, } from "../orchestration-v2/EventSink.ts"; import { layer as eventStoreLayer } from "../orchestration-v2/EventStore.ts"; -import { layer as idAllocatorLayer } from "../orchestration-v2/IdAllocator.ts"; +import { IdAllocatorV2, layer as idAllocatorLayer } from "../orchestration-v2/IdAllocator.ts"; import { LegacyV1ThreadImporter, layer as legacyImporterLayer, @@ -39,8 +38,9 @@ import { ProjectionStoreV2, layer as projectionStoreLayer, } from "../orchestration-v2/ProjectionStore.ts"; +import { layer as projectStoreLayer } from "../orchestration-v2/ProjectStore.ts"; import { layer as threadCommandExecutorLayer } from "../orchestration-v2/ThreadCommandExecutor.ts"; -import { ProjectionProjectRepositoryLive } from "../persistence/Layers/ProjectionProjects.ts"; +import { planThreadDeletion } from "../orchestration-v2/ThreadDeletion.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as WorkspacePaths from "../workspace/WorkspacePaths.ts"; import * as ProjectEnrichmentService from "./ProjectEnrichmentService.ts"; @@ -54,8 +54,7 @@ const eventPersistenceLayer = eventSinkLayer.pipe( const servicesLayer = Layer.mergeAll( legacyImporterLayer.pipe(Layer.provideMerge(eventPersistenceLayer)), projectionMaintenanceLayer.pipe(Layer.provide(eventPersistenceLayer)), - OrchestrationLayerLive, - ProjectionProjectRepositoryLive, + projectStoreLayer, idAllocatorLayer, threadCommandExecutorLayer, Layer.succeed(WorkspacePaths.WorkspacePaths, { @@ -333,22 +332,24 @@ it.effect( assert.lengthOf(projection.messages, 4); assert.equal(yield* importer.pendingThreadCount, 0); const rows = yield* sql<{ - readonly legacy_deleted_at: string | null; readonly v2_deleted_at: string | null; readonly project_deleted_at: string | null; }>` - SELECT legacy.deleted_at AS legacy_deleted_at, - v2.deleted_at AS v2_deleted_at, + SELECT v2.deleted_at AS v2_deleted_at, project.deleted_at AS project_deleted_at - FROM projection_threads AS legacy - JOIN orchestration_v2_projection_threads AS v2 ON v2.thread_id = legacy.thread_id - JOIN projection_projects AS project ON project.project_id = legacy.project_id - WHERE legacy.thread_id = ${threadId} + FROM orchestration_v2_projection_threads AS v2 + JOIN projection_projects AS project ON project.project_id = v2.project_id + WHERE v2.thread_id = ${threadId} `; assert.lengthOf(rows, 1); - assert.isNotNull(rows[0]?.legacy_deleted_at); assert.isNotNull(rows[0]?.v2_deleted_at); assert.isNotNull(rows[0]?.project_deleted_at); + // The legacy V1 row is only an import source; deletion writes V2 events only. + const legacyEvents = yield* sql` + SELECT sequence FROM orchestration_events + WHERE application_event_version = 1 AND stream_id = ${threadId} + `; + assert.deepEqual(legacyEvents, []); const cleanup = yield* sql<{ readonly command_id: string; @@ -420,3 +421,80 @@ it.effect("rejects a child deletion command ID already accepted for an unrelated }).pipe(Effect.provide(servicesLayer)); }).pipe(Effect.provide(databaseLayer)), ); + +it.effect("deletes a project without force once its imported threads were deleted in V2", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:imported-emptied"); + const threadId = ThreadId.make("thread:imported-emptied"); + yield* seedProject(projectId); + // V2 never writes the V1 thread table, so this row stays live there. + yield* sql` + INSERT INTO projection_threads ( + thread_id, project_id, title, model_selection_json, + runtime_mode, interaction_mode, branch, worktree_path, latest_turn_id, + created_at, updated_at, archived_at, deleted_at + ) VALUES ( + ${threadId}, ${projectId}, 'Imported thread', '{"instanceId":"codex","model":"gpt-5.4"}', + 'full-access', 'default', NULL, NULL, NULL, + '2026-01-01T00:00:00.000Z', '2026-01-01T00:00:00.000Z', NULL, NULL + ) + `; + yield* TestClock.setTime(Date.parse("2026-09-04T12:00:00.000Z")); + + yield* Effect.gen(function* () { + const importer = yield* LegacyV1ThreadImporter; + const projections = yield* ProjectionStoreV2; + const eventSink = yield* EventSinkV2; + const idAllocator = yield* IdAllocatorV2; + const service = yield* ProjectService.make; + yield* importer.reconcileShells; + const early = yield* service + .delete({ commandId: CommandId.make("command:emptied:early"), projectId }) + .pipe(Effect.flip); + assert.equal(early._tag, "ProjectNotEmptyError"); + + // Delete the imported thread the way thread.delete does. + yield* importer.ensureTranscript(threadId); + const command = { + type: "thread.delete" as const, + commandId: CommandId.make("command:emptied:thread-delete"), + threadId, + }; + const now = yield* DateTime.now; + const plan = yield* planThreadDeletion({ + command, + projection: yield* projections.getThreadRecords(threadId, [ + "runs", + "attempts", + "nodes", + "runtimeRequests", + "subagents", + "providerSessions", + ]), + attachmentIds: [], + now, + idAllocator, + }); + yield* eventSink.commitCommand({ + commandId: command.commandId, + commandType: command.type, + threadId, + acceptedAt: now, + events: plan.events, + effects: plan.effects, + }); + const legacy = yield* sql<{ readonly deleted_at: string | null }>` + SELECT deleted_at FROM projection_threads WHERE thread_id = ${threadId} + `; + assert.deepEqual(legacy, [{ deleted_at: null }]); + + const deleted = yield* service.delete({ + commandId: CommandId.make("command:emptied:project-delete"), + projectId, + }); + assert.isNotNull(deleted.deletedAt); + assert.isTrue(Option.isNone(yield* service.getById(projectId))); + }).pipe(Effect.provide(servicesLayer)); + }).pipe(Effect.provide(databaseLayer)), +); diff --git a/apps/server/src/project/ProjectService.test.ts b/apps/server/src/project/ProjectService.test.ts index cc3082b5b16e..8cae24f4d328 100644 --- a/apps/server/src/project/ProjectService.test.ts +++ b/apps/server/src/project/ProjectService.test.ts @@ -11,7 +11,16 @@ import { TestClock } from "effect/testing"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { ServerConfig } from "../config.ts"; -import { ProjectServiceLayerLive } from "../orchestration-v2/runtimeLayer.ts"; +import { EventSinkV2 } from "../orchestration-v2/EventSink.ts"; +import * as IdAllocator from "../orchestration-v2/IdAllocator.ts"; +import * as LegacyV1ThreadImporter from "../orchestration-v2/LegacyV1ThreadImporter.ts"; +import * as ProjectionStore from "../orchestration-v2/ProjectionStore.ts"; +import * as ProjectStore from "../orchestration-v2/ProjectStore.ts"; +import { + OrchestrationV2EventSinkLayerLive, + ProjectServiceLayerLive, +} from "../orchestration-v2/runtimeLayer.ts"; +import { layer as threadCommandExecutorLayer } from "../orchestration-v2/ThreadCommandExecutor.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import * as WorkspacePaths from "../workspace/WorkspacePaths.ts"; import * as ProjectEnrichmentService from "./ProjectEnrichmentService.ts"; @@ -60,6 +69,25 @@ const makeTestLayer = ( const TestLayer = makeTestLayer(metadataLayer); +/** Every dependency of ProjectService.make, so a test can swap one of them. */ +const ProjectServiceDependenciesLayer = Layer.mergeAll( + OrchestrationV2EventSinkLayerLive, + ProjectStore.layer, + ProjectionStore.layer, + IdAllocator.layer, + threadCommandExecutorLayer, +).pipe( + Layer.provideMerge( + LegacyV1ThreadImporter.layer.pipe(Layer.provide(OrchestrationV2EventSinkLayerLive)), + ), + Layer.provideMerge(ProjectEnrichmentService.layer), + Layer.provideMerge(workspacePathsLayer), + Layer.provideMerge(metadataLayer), + Layer.provideMerge(SqlitePersistenceMemory), + Layer.provide(ServerConfig.layerTest(process.cwd(), { prefix: "project-service-race-" })), + Layer.provide(NodeServices.layer), +); + const waitForProject = Effect.fn("ProjectServiceTest.waitForProject")(function* ( service: ProjectService.ProjectService["Service"], projectId: ProjectId, @@ -220,6 +248,247 @@ it.layer(TestLayer)("ProjectService", (it) => { assert.equal(second.project.id, first.project.id); }), ); + + it.effect("commits the event, the row and the receipt together, once per command id", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:retry"); + const input = { + commandId: CommandId.make("command:retry:create"), + projectId, + title: "Retry", + workspaceRoot: "/work/retry", + }; + const first = yield* service.create(input); + // A retry re-plans against the row it created; its receipt still answers. + const retried = yield* service.create(input); + assert.deepEqual(retried, first); + const updateInput = { + commandId: CommandId.make("command:retry:update"), + projectId, + autoPull: true, + }; + yield* service.update(updateInput); + yield* service.update(updateInput); + + const committed = yield* sql<{ + readonly sequence: number; + readonly event_type: string; + readonly receipt_sequence: number; + readonly status: string; + readonly aggregate_kind: string; + }>` + SELECT events.sequence, events.event_type, + receipts.result_sequence AS receipt_sequence, receipts.status, receipts.aggregate_kind + FROM orchestration_events AS events + JOIN orchestration_command_receipts AS receipts ON receipts.command_id = events.command_id + WHERE events.stream_id = ${projectId} + ORDER BY events.sequence + `; + assert.deepEqual( + committed.map((row) => [row.event_type, row.status, row.aggregate_kind]), + [ + ["project.created", "accepted", "project"], + ["project.meta-updated", "accepted", "project"], + ], + ); + for (const row of committed) assert.equal(row.receipt_sequence, row.sequence); + assert.isTrue(Option.getOrThrow(yield* service.getById(projectId)).autoPull); + }), + ); + + it.effect("replays a rejected receipt instead of re-planning the command", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:rejected"); + yield* service.create({ + commandId: CommandId.make("command:rejected:create"), + projectId, + title: "Rejected", + workspaceRoot: "/work/rejected", + }); + const invalid = { + commandId: CommandId.make("command:rejected:script"), + projectId, + scripts: [ + { + id: "Not.Valid", + name: "Invalid", + command: "true", + icon: "play" as const, + runOnWorktreeCreate: false, + }, + ], + }; + const first = yield* service.update(invalid).pipe(Effect.flip); + assert.equal(first._tag, "ProjectOperationError"); + // Same id, now a valid payload: the recorded rejection wins. + const replayed = yield* service + .update({ ...invalid, scripts: [{ ...invalid.scripts[0]!, id: "valid" }] }) + .pipe(Effect.flip); + assert.equal(replayed._tag, "ProjectOperationError"); + assert.include( + String(replayed._tag === "ProjectOperationError" && replayed.cause), + "Script ID", + ); + const receipts = yield* sql<{ readonly status: string; readonly aggregate_id: string }>` + SELECT status, aggregate_id FROM orchestration_command_receipts + WHERE command_id = ${invalid.commandId} + `; + assert.deepEqual(receipts, [{ status: "rejected", aggregate_id: projectId }]); + const events = yield* sql` + SELECT sequence FROM orchestration_events WHERE command_id = ${invalid.commandId} + `; + assert.deepEqual(events, []); + assert.deepEqual(Option.getOrThrow(yield* service.getById(projectId)).scripts, []); + }), + ); + + it.effect("rejects a workspace another active project holds", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const other = ProjectId.make("project:conflict:mover"); + yield* service.create({ + commandId: CommandId.make("command:conflict:holder"), + projectId: ProjectId.make("project:conflict:holder"), + title: "Holder", + workspaceRoot: "/work/conflict", + }); + yield* service.create({ + commandId: CommandId.make("command:conflict:mover"), + projectId: other, + title: "Mover", + workspaceRoot: "/work/conflict-mover", + }); + const move = yield* service + .update({ + commandId: CommandId.make("command:conflict:move"), + projectId: other, + workspaceRoot: "/work/conflict", + }) + .pipe(Effect.flip); + assert.equal(move._tag, "ProjectConflictError"); + assert.equal( + Option.getOrThrow(yield* service.getById(other)).workspaceRoot, + "/work/conflict-mover", + ); + }), + ); + + it.effect("replays a workspace conflict after the workspace frees up", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const holderId = ProjectId.make("project:freed:holder"); + yield* service.create({ + commandId: CommandId.make("command:freed:holder"), + projectId: holderId, + title: "Holder", + workspaceRoot: "/work/freed", + }); + const claim = { + commandId: CommandId.make("command:freed:claim"), + projectId: ProjectId.make("project:freed:claim"), + title: "Claim", + workspaceRoot: "/work/freed", + }; + const first = yield* service.create(claim).pipe(Effect.flip); + yield* service.delete({ + commandId: CommandId.make("command:freed:delete"), + projectId: holderId, + }); + // A fresh plan would now succeed; the recorded conflict still answers. + const replayed = yield* service.create(claim).pipe(Effect.flip); + assert.deepEqual(replayed, first); + assert.instanceOf(replayed, ProjectService.ProjectConflictError); + assert.isTrue( + Option.isNone(yield* service.getById(claim.projectId, { includeDeleted: true })), + ); + }), + ); + + it.effect("returns the deleted project when a completed delete is retried", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:delete-retry"); + yield* service.create({ + commandId: CommandId.make("command:delete-retry:create"), + projectId, + title: "Delete retry", + workspaceRoot: "/work/delete-retry", + }); + const input = { commandId: CommandId.make("command:delete-retry:delete"), projectId }; + const deleted = yield* service.delete(input); + assert.deepEqual(yield* service.delete(input), deleted); + const events = yield* sql<{ readonly command_id: string }>` + SELECT command_id FROM orchestration_events + WHERE stream_id = ${projectId} AND event_type = 'project.deleted' + `; + assert.deepEqual(events, [{ command_id: input.commandId }]); + }), + ); + + it.effect("treats a deleted project as missing for updates and new deletes", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:deleted-target"); + yield* service.create({ + commandId: CommandId.make("command:deleted-target:create"), + projectId, + title: "Deleted target", + workspaceRoot: "/work/deleted-target", + }); + yield* service.delete({ + commandId: CommandId.make("command:deleted-target:delete"), + projectId, + }); + const again = yield* service + .delete({ commandId: CommandId.make("command:deleted-target:delete-again"), projectId }) + .pipe(Effect.flip); + assert.instanceOf(again, ProjectService.ProjectNotFoundError); + const events = yield* sql<{ readonly event_type: string }>` + SELECT event_type FROM orchestration_events + WHERE stream_id = ${projectId} ORDER BY sequence ASC + `; + assert.deepEqual( + events.map((event) => event.event_type), + ["project.created", "project.deleted"], + ); + }), + ); + + it.effect("rolls back the event and receipt when the row write fails", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:atomic"); + yield* sql` + CREATE TRIGGER fail_project_row BEFORE INSERT ON projection_projects + WHEN NEW.project_id = 'project:atomic' + BEGIN SELECT RAISE(ABORT, 'injected row failure'); END + `; + const input = { + commandId: CommandId.make("command:atomic:create"), + projectId, + title: "Atomic", + workspaceRoot: "/work/atomic", + }; + const failure = yield* service.create(input).pipe(Effect.flip); + assert.equal(failure._tag, "ProjectOperationError"); + const leftovers = yield* sql<{ readonly count: number }>` + SELECT + (SELECT COUNT(*) FROM orchestration_events WHERE command_id = ${input.commandId}) + + (SELECT COUNT(*) FROM orchestration_command_receipts WHERE command_id = ${input.commandId}) + + (SELECT COUNT(*) FROM projection_projects WHERE project_id = ${projectId}) AS count + `; + assert.equal(leftovers[0]?.count, 0); + yield* sql`DROP TRIGGER fail_project_row`; + assert.equal((yield* service.create(input)).id, projectId); + }), + ); }); it.effect( @@ -434,3 +703,99 @@ it.effect("invalidates workspace-derived metadata when a project moves", () => }).pipe(Effect.provide(makeTestLayer(versionedMetadataLayer))); }), ); + +it.effect("serializes two projects claiming the same workspace root", () => + Effect.gen(function* () { + const eventSink = yield* EventSinkV2; + const firstReachedCommit = yield* Deferred.make(); + const releaseFirst = yield* Deferred.make(); + // Hold the first claim between its plan and its commit. + const gatedSink = EventSinkV2.of({ + ...eventSink, + commitProjectCommand: (input) => + input.projectId === "project:race:first" + ? Deferred.succeed(firstReachedCommit, undefined).pipe( + Effect.andThen(Deferred.await(releaseFirst)), + Effect.andThen(eventSink.commitProjectCommand(input)), + ) + : eventSink.commitProjectCommand(input), + }); + const service = yield* ProjectService.make.pipe(Effect.provideService(EventSinkV2, gatedSink)); + const claim = (name: string) => + service.create({ + commandId: CommandId.make(`command:race:${name}`), + projectId: ProjectId.make(`project:race:${name}`), + title: name, + workspaceRoot: "/work/race", + }); + const first = yield* claim("first").pipe(Effect.forkChild({ startImmediately: true })); + yield* Deferred.await(firstReachedCommit); + const second = yield* claim("second").pipe( + Effect.flip, + Effect.forkChild({ startImmediately: true }), + ); + yield* Effect.yieldNow; + assert.isUndefined(second.pollUnsafe()); + yield* Deferred.succeed(releaseFirst, undefined); + assert.equal((yield* Fiber.join(first)).id, "project:race:first"); + assert.equal((yield* Fiber.join(second))._tag, "ProjectConflictError"); + }).pipe(Effect.provide(ProjectServiceDependenciesLayer)), +); + +it.effect("rejects an update that waited on the lock while its project was deleted", () => + Effect.gen(function* () { + const eventSink = yield* EventSinkV2; + const deleteReachedCommit = yield* Deferred.make(); + const releaseDelete = yield* Deferred.make(); + // Hold the delete between its plan and its commit, inside the project lock. + const gatedSink = EventSinkV2.of({ + ...eventSink, + commitProjectCommand: (input) => + input.commandType === "project.delete" + ? Deferred.succeed(deleteReachedCommit, undefined).pipe( + Effect.andThen(Deferred.await(releaseDelete)), + Effect.andThen(eventSink.commitProjectCommand(input)), + ) + : eventSink.commitProjectCommand(input), + }); + const service = yield* ProjectService.make.pipe(Effect.provideService(EventSinkV2, gatedSink)); + const sql = yield* SqlClient.SqlClient; + const projectId = ProjectId.make("project:update-race"); + yield* service.create({ + commandId: CommandId.make("command:update-race:create"), + projectId, + title: "Update race", + workspaceRoot: "/work/update-race", + }); + const deletion = yield* service + .delete({ commandId: CommandId.make("command:update-race:delete"), projectId }) + .pipe(Effect.forkChild({ startImmediately: true })); + yield* Deferred.await(deleteReachedCommit); + const updateCommandId = CommandId.make("command:update-race:update"); + // The update sees the active row, then queues behind the delete's lock. + const update = yield* service + .update({ commandId: updateCommandId, projectId, title: "Too late" }) + .pipe(Effect.flip, Effect.forkChild({ startImmediately: true })); + yield* Effect.yieldNow; + yield* Deferred.succeed(releaseDelete, undefined); + assert.isNotNull((yield* Fiber.join(deletion)).deletedAt); + assert.instanceOf(yield* Fiber.join(update), ProjectService.ProjectNotFoundError); + // The planner rejected it under the lock, so the rejection has a receipt. + const receipts = yield* sql<{ readonly status: string }>` + SELECT status FROM orchestration_command_receipts WHERE command_id = ${updateCommandId} + `; + assert.deepEqual(receipts, [{ status: "rejected" }]); + const events = yield* sql<{ readonly event_type: string }>` + SELECT event_type FROM orchestration_events + WHERE stream_id = ${projectId} ORDER BY sequence ASC + `; + assert.deepEqual( + events.map((event) => event.event_type), + ["project.created", "project.deleted"], + ); + assert.equal( + Option.getOrThrow(yield* service.getById(projectId, { includeDeleted: true })).title, + "Update race", + ); + }).pipe(Effect.provide(ProjectServiceDependenciesLayer)), +); diff --git a/apps/server/src/project/ProjectService.ts b/apps/server/src/project/ProjectService.ts index d390e7379330..cd5d58cdb6d5 100644 --- a/apps/server/src/project/ProjectService.ts +++ b/apps/server/src/project/ProjectService.ts @@ -5,25 +5,33 @@ import { type ProjectCreatePayload, type ProjectUpdatePayload, type ProjectSnapshot, + type ThreadId, } from "@t3tools/contracts"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; -import { OrchestrationEngineService } from "../orchestration/Services/OrchestrationEngine.ts"; import { EventSinkV2 } from "../orchestration-v2/EventSink.ts"; -import { IdAllocatorV2 } from "../orchestration-v2/IdAllocator.ts"; +import * as IdAllocator from "../orchestration-v2/IdAllocator.ts"; +import { makeKeyedSerialExecutor } from "../orchestration-v2/KeyedSerialExecutor.ts"; import { LegacyV1ThreadImporter } from "../orchestration-v2/LegacyV1ThreadImporter.ts"; import { ProjectionStoreV2 } from "../orchestration-v2/ProjectionStore.ts"; +import { + decodeProjectCommandRejection, + encodeProjectCommandRejection, + planProjectCommand, + type ProjectCommand, +} from "../orchestration-v2/ProjectCommands.ts"; +import * as ProjectStore from "../orchestration-v2/ProjectStore.ts"; import { ThreadCommandExecutor, layer as threadCommandExecutorLayer, } from "../orchestration-v2/ThreadCommandExecutor.ts"; import { planThreadDeletion } from "../orchestration-v2/ThreadDeletion.ts"; -import * as ProjectionProjects from "../persistence/Services/ProjectionProjects.ts"; import { ProjectEnrichmentService, type ProjectEnrichment } from "./ProjectEnrichmentService.ts"; import * as WorkspacePaths from "../workspace/WorkspacePaths.ts"; @@ -128,18 +136,21 @@ export class ProjectService extends Context.Service< >()("t3/project/ProjectService") {} export const make = Effect.gen(function* () { - const engine = yield* OrchestrationEngineService; - const projects = yield* ProjectionProjects.ProjectionProjectRepository; + const projects = yield* ProjectStore.ProjectStoreV2; const projectEnrichment = yield* ProjectEnrichmentService; const workspacePaths = yield* WorkspacePaths.WorkspacePaths; const threadProjections = yield* ProjectionStoreV2; - const threadEvents = yield* EventSinkV2; - const idAllocator = yield* IdAllocatorV2; + const eventSink = yield* EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; const legacyImporter = yield* LegacyV1ThreadImporter; const threadCommands = yield* ThreadCommandExecutor; + // Commands for one project run in order. Commands that claim a workspace root + // also hold that root, so two projects cannot both claim it. + const projectLocks = yield* makeKeyedSerialExecutor(); + const workspaceLocks = yield* makeKeyedSerialExecutor(); const toProject = ( - row: ProjectionProjects.ProjectionProject, + row: ProjectStore.ProjectRow, enrichment: ProjectEnrichment | null, ): Project => ({ id: row.projectId, @@ -157,9 +168,7 @@ export const make = Effect.gen(function* () { deletedAt: row.deletedAt, }); - const hydrateAvailable = Effect.fn("ProjectService.hydrateAvailable")(function* ( - row: ProjectionProjects.ProjectionProject, - ) { + const hydrate = Effect.fn("ProjectService.hydrate")(function* (row: ProjectStore.ProjectRow) { const enrichment = row.deletedAt === null ? yield* projectEnrichment.getAvailable(row.workspaceRoot) @@ -167,200 +176,210 @@ export const make = Effect.gen(function* () { return toProject(row, enrichment); }); - const readRows = Effect.fn("ProjectService.readRows")(function* () { - return yield* projects - .listAll() + const readRow = (projectId: ProjectId, options?: { readonly includeDeleted?: boolean }) => + projects + .get(projectId, options) .pipe( Effect.mapError( - (cause) => new ProjectOperationError({ operation: "list-projects", cause }), + (cause) => new ProjectOperationError({ operation: "read-project", projectId, cause }), ), ); - }); - - const getById: ProjectService["Service"]["getById"] = Effect.fn("ProjectService.getById")( - function* (projectId, options) { - const row = yield* projects - .getById({ projectId }) - .pipe( - Effect.mapError( - (cause) => new ProjectOperationError({ operation: "read-project", projectId, cause }), - ), - ); - if (Option.isNone(row) || (row.value.deletedAt !== null && !options?.includeDeleted)) { - return Option.none(); - } - return Option.some(yield* hydrateAvailable(row.value)); - }, - ); - const getByWorkspaceRoot: ProjectService["Service"]["getByWorkspaceRoot"] = Effect.fn( - "ProjectService.getByWorkspaceRoot", - )(function* (workspaceRoot, options) { - const normalized = yield* workspacePaths.normalizeWorkspaceRoot(workspaceRoot).pipe( - Effect.mapError( - (cause) => - new ProjectOperationError({ - operation: "normalize-workspace", - workspaceRoot, - cause, - }), - ), - ); - const row = (yield* readRows()).find( - (candidate) => - candidate.workspaceRoot === normalized && - (options?.includeDeleted === true || candidate.deletedAt === null), - ); - return row === undefined ? Option.none() : Option.some(yield* hydrateAvailable(row)); - }); - - const readCommitted = Effect.fn("ProjectService.readCommitted")(function* (projectId: ProjectId) { - const row = yield* projects - .getById({ projectId }) + const normalizeWorkspaceRoot = (input: { + readonly projectId?: ProjectId; + readonly workspaceRoot: string; + readonly createIfMissing?: boolean; + }) => + workspacePaths + .normalizeWorkspaceRoot(input.workspaceRoot, { + createIfMissing: input.createIfMissing ?? false, + }) .pipe( Effect.mapError( - (cause) => new ProjectOperationError({ operation: "read-project", projectId, cause }), + (cause) => + new ProjectOperationError({ + operation: "normalize-workspace", + ...(input.projectId === undefined ? {} : { projectId: input.projectId }), + workspaceRoot: input.workspaceRoot, + cause, + }), ), ); + + /** + * Plan one command against rows read under its locks, then commit its event or + * its rejection. A reused command id resolves to the receipt it already has. + */ + const commit = Effect.fn("ProjectService.commit")(function* (command: ProjectCommand) { + const { projectId } = command; + const dispatchError = (cause: unknown) => + new ProjectOperationError({ operation: "dispatch-project-command", projectId, cause }); + const workspaceRoot = command.type === "project.delete" ? undefined : command.workspaceRoot; + const planAndCommit = Effect.gen(function* () { + const project = Option.getOrUndefined(yield* readRow(projectId, { includeDeleted: true })); + const workspaceOwner = + workspaceRoot === undefined + ? undefined + : Option.getOrUndefined( + yield* projects + .findActiveByWorkspaceRoot(workspaceRoot) + .pipe(Effect.mapError(dispatchError)), + ); + const now = yield* DateTime.now; + const eventId = yield* idAllocator.allocate + .event({ commandId: command.commandId }) + .pipe(Effect.mapError(dispatchError)); + const planned = planProjectCommand({ + command, + state: { project, workspaceOwner }, + eventId, + now, + }); + if (Result.isSuccess(planned)) { + const { receipt } = yield* eventSink.commitProjectCommand({ + commandId: command.commandId, + projectId, + commandType: command.type, + acceptedAt: now, + event: planned.success, + }); + return receipt; + } + return yield* eventSink.commitRejectedProjectCommand({ + commandId: command.commandId, + projectId, + commandType: command.type, + rejectedAt: now, + error: encodeProjectCommandRejection(planned.failure), + }); + }); + const receipt = yield* projectLocks + .withLock( + projectId, + workspaceRoot === undefined + ? planAndCommit + : workspaceLocks.withLock(workspaceRoot, planAndCommit), + ) + .pipe(Effect.mapError(dispatchError)); + if (receipt.projectId !== projectId || receipt.commandType !== command.type) { + return yield* dispatchError( + `Command ${command.commandId} was already used by ${receipt.commandType} for ${receipt.projectId}.`, + ); + } + // A retried command re-plans against the state it already produced, so its + // first receipt, not the new plan, decides the outcome. + if (receipt.status === "accepted") return; + const rejection = Option.getOrUndefined(decodeProjectCommandRejection(receipt.error)); + switch (rejection?._tag) { + case "ProjectWorkspaceConflictError": + return yield* new ProjectConflictError({ + projectId, + workspaceRoot: rejection.workspaceRoot, + conflictingProjectId: rejection.conflictingProjectId, + }); + case "ProjectCommandMissingProjectError": + return yield* new ProjectNotFoundError({ projectId }); + default: + return yield* dispatchError( + rejection ?? receipt.error ?? "The command was previously rejected.", + ); + } + }); + + const readCommitted = Effect.fn("ProjectService.readCommitted")(function* (projectId: ProjectId) { + const row = yield* readRow(projectId, { includeDeleted: true }); if (Option.isNone(row)) { return yield* new ProjectOperationError({ operation: "read-project", projectId, - cause: "The accepted project command did not produce a project projection.", + cause: "The accepted project command did not produce a project row.", }); } - return yield* hydrateAvailable(row.value); + return yield* hydrate(row.value); }); - const invalidateEnrichment = (...workspaceRoots: ReadonlyArray) => - projectEnrichment.invalidate(workspaceRoots); - - const dispatch = ( - projectId: ProjectId, - command: Parameters[0], - onCommitted: Effect.Effect, - ) => - engine.dispatch(command).pipe( - Effect.mapError( - (cause) => - new ProjectOperationError({ - operation: "dispatch-project-command", - projectId, - cause, - }), - ), - Effect.andThen(onCommitted), - ); + const getById: ProjectService["Service"]["getById"] = Effect.fn("ProjectService.getById")( + function* (projectId, options) { + const row = yield* readRow(projectId, options); + return Option.isNone(row) ? Option.none() : Option.some(yield* hydrate(row.value)); + }, + ); - const assertWorkspaceAvailable = Effect.fn("ProjectService.assertWorkspaceAvailable")(function* ( - projectId: ProjectId, - workspaceRoot: string, - ) { - const conflicting = (yield* readRows()).find( - (candidate) => - candidate.deletedAt === null && - candidate.projectId !== projectId && - candidate.workspaceRoot === workspaceRoot, + const getByWorkspaceRoot: ProjectService["Service"]["getByWorkspaceRoot"] = Effect.fn( + "ProjectService.getByWorkspaceRoot", + )(function* (workspaceRoot, options) { + const normalized = yield* normalizeWorkspaceRoot({ workspaceRoot }); + const row = yield* ( + options?.includeDeleted === true + ? projects + .list({ includeDeleted: true }) + .pipe( + Effect.map((rows) => + Option.fromUndefinedOr(rows.find((row) => row.workspaceRoot === normalized)), + ), + ) + : projects.findActiveByWorkspaceRoot(normalized) + ).pipe( + Effect.mapError((cause) => new ProjectOperationError({ operation: "list-projects", cause })), ); - if (conflicting !== undefined) { - return yield* new ProjectConflictError({ - projectId, - workspaceRoot, - conflictingProjectId: conflicting.projectId, - }); - } + return Option.isNone(row) ? Option.none() : Option.some(yield* hydrate(row.value)); }); const create: ProjectService["Service"]["create"] = Effect.fn("ProjectService.create")( function* (input) { - const workspaceRoot = yield* workspacePaths - .normalizeWorkspaceRoot(input.workspaceRoot, { - createIfMissing: input.createWorkspaceRootIfMissing ?? false, - }) - .pipe( - Effect.mapError( - (cause) => - new ProjectOperationError({ - operation: "normalize-workspace", - projectId: input.projectId, - workspaceRoot: input.workspaceRoot, - cause, - }), - ), - ); - yield* assertWorkspaceAvailable(input.projectId, workspaceRoot); - const now = DateTime.formatIso(yield* DateTime.now); - return yield* dispatch( - input.projectId, - { - type: "project.create", - commandId: input.commandId, - projectId: input.projectId, - title: input.title, - workspaceRoot, - defaultModelSelection: input.defaultModelSelection ?? null, - scripts: [...(input.scripts ?? [])], - createdAt: now, - }, - invalidateEnrichment(workspaceRoot).pipe(Effect.andThen(readCommitted(input.projectId))), - ); + const workspaceRoot = yield* normalizeWorkspaceRoot({ + projectId: input.projectId, + workspaceRoot: input.workspaceRoot, + createIfMissing: input.createWorkspaceRootIfMissing ?? false, + }); + yield* commit({ + type: "project.create", + commandId: input.commandId, + projectId: input.projectId, + title: input.title, + workspaceRoot, + ...(input.scripts === undefined ? {} : { scripts: input.scripts }), + }); + yield* projectEnrichment.invalidate([workspaceRoot]); + return yield* readCommitted(input.projectId); }, ); const update: ProjectService["Service"]["update"] = Effect.fn("ProjectService.update")( function* (input) { - const existing = yield* projects.getById({ projectId: input.projectId }).pipe( - Effect.mapError( - (cause) => - new ProjectOperationError({ - operation: "read-project", - projectId: input.projectId, - cause, - }), - ), - ); - if (Option.isNone(existing) || existing.value.deletedAt !== null) { + const existing = yield* readRow(input.projectId); + if (Option.isNone(existing)) { return yield* new ProjectNotFoundError({ projectId: input.projectId }); } + const previousRoot = existing.value.workspaceRoot; const workspaceRoot = input.workspaceRoot === undefined - ? existing.value.workspaceRoot - : yield* workspacePaths.normalizeWorkspaceRoot(input.workspaceRoot).pipe( - Effect.mapError( - (cause) => - new ProjectOperationError({ - operation: "normalize-workspace", - projectId: input.projectId, - workspaceRoot: input.workspaceRoot, - cause, - }), - ), - ); - yield* assertWorkspaceAvailable(input.projectId, workspaceRoot); - return yield* dispatch( - input.projectId, - { - type: "project.meta.update", - commandId: input.commandId, - projectId: input.projectId, - ...(input.title === undefined ? {} : { title: input.title }), - ...(workspaceRoot === existing.value.workspaceRoot ? {} : { workspaceRoot }), - ...(input.defaultModelSelection === undefined - ? {} - : { defaultModelSelection: input.defaultModelSelection }), - ...(input.autoPull === undefined ? {} : { autoPull: input.autoPull }), - ...(input.projectIcon === undefined ? {} : { projectIcon: input.projectIcon }), - ...(input.faviconPath === undefined ? {} : { faviconPath: input.faviconPath }), - ...(input.defaultThreadEnvMode === undefined - ? {} - : { defaultThreadEnvMode: input.defaultThreadEnvMode }), - ...(input.scripts === undefined ? {} : { scripts: [...input.scripts] }), - }, - (workspaceRoot === existing.value.workspaceRoot - ? Effect.void - : invalidateEnrichment(existing.value.workspaceRoot, workspaceRoot) - ).pipe(Effect.andThen(readCommitted(input.projectId))), - ); + ? previousRoot + : yield* normalizeWorkspaceRoot({ + projectId: input.projectId, + workspaceRoot: input.workspaceRoot, + }); + yield* commit({ + type: "project.meta.update", + commandId: input.commandId, + projectId: input.projectId, + ...(input.title === undefined ? {} : { title: input.title }), + ...(workspaceRoot === previousRoot ? {} : { workspaceRoot }), + ...(input.defaultModelSelection === undefined + ? {} + : { defaultModelSelection: input.defaultModelSelection }), + ...(input.autoPull === undefined ? {} : { autoPull: input.autoPull }), + ...(input.projectIcon === undefined ? {} : { projectIcon: input.projectIcon }), + ...(input.faviconPath === undefined ? {} : { faviconPath: input.faviconPath }), + ...(input.defaultThreadEnvMode === undefined + ? {} + : { defaultThreadEnvMode: input.defaultThreadEnvMode }), + ...(input.scripts === undefined ? {} : { scripts: input.scripts }), + }); + if (workspaceRoot !== previousRoot) { + yield* projectEnrichment.invalidate([previousRoot, workspaceRoot]); + } + return yield* readCommitted(input.projectId); }, ); @@ -372,120 +391,120 @@ export const make = Effect.gen(function* () { }, ); + /** Delete one child thread durably; a stable command id makes a retry resume the cascade. */ + const deleteChildThread = Effect.fn("ProjectService.deleteChildThread")(function* ( + input: ProjectDeleteInput, + threadId: ThreadId, + ) { + yield* legacyImporter.ensureTranscript(threadId); + const projection = yield* threadProjections.getThreadRecords(threadId, [ + "runs", + "attempts", + "nodes", + "runtimeRequests", + "subagents", + "providerSessions", + ]); + if (projection.thread.deletedAt !== null || projection.thread.projectId !== input.projectId) { + return; + } + const command = { + type: "thread.delete" as const, + commandId: CommandId.make(`${input.commandId}:delete-thread:${threadId}`), + threadId, + }; + const now = yield* DateTime.now; + const plan = yield* planThreadDeletion({ + command, + projection, + attachmentIds: yield* threadProjections.getThreadAttachmentIds(threadId), + now, + idAllocator, + }); + const committed = yield* eventSink.commitCommand({ + commandId: command.commandId, + commandType: command.type, + threadId, + acceptedAt: now, + events: plan.events, + effects: plan.effects, + }); + if ( + committed.receipt.threadId !== command.threadId || + committed.receipt.commandType !== command.type + ) { + return yield* Effect.fail("The thread deletion command ID belongs to a different command."); + } + if (committed.receipt.status === "rejected") { + return yield* Effect.fail( + committed.receipt.error ?? "Thread deletion was previously rejected.", + ); + } + }); + + /** Refuse a non-empty project without force, else delete its live threads first. */ + const deleteChildThreads = Effect.fn("ProjectService.deleteChildThreads")(function* ( + input: ProjectDeleteInput, + ) { + const { projectId } = input; + // The V2 shell is the only record of which threads are live. + const snapshot = yield* threadProjections + .getShellSnapshot() + .pipe( + Effect.mapError( + (cause) => new ProjectOperationError({ operation: "list-threads", projectId, cause }), + ), + ); + const projectThreads = [...snapshot.threads, ...snapshot.archivedThreads].filter( + (thread) => thread.projectId === projectId, + ); + if (projectThreads.length > 0 && input.force !== true) { + return yield* new ProjectNotEmptyError({ projectId }); + } + // Delete children durably before the project so a failed cascade can be retried. + yield* Effect.forEach( + projectThreads, + (thread) => + threadCommands + .withLock(thread.id, deleteChildThread(input, thread.id)) + .pipe( + Effect.mapError( + (cause) => + new ProjectOperationError({ operation: "delete-thread", projectId, cause }), + ), + ), + { concurrency: 1, discard: true }, + ); + }); + const deleteProject: ProjectService["Service"]["delete"] = Effect.fn("ProjectService.delete")( function* (input) { const { projectId } = input; - const existing = yield* projects - .getById({ projectId }) - .pipe( - Effect.mapError( - (cause) => new ProjectOperationError({ operation: "read-project", projectId, cause }), - ), - ); - if (Option.isNone(existing) || existing.value.deletedAt !== null) { + // A deleted row still reaches commit, so a retried command id replays its + // receipt and any other command id is rejected as not found. + const existing = yield* readRow(projectId, { includeDeleted: true }); + if (Option.isNone(existing)) { return yield* new ProjectNotFoundError({ projectId }); } - const snapshot = yield* threadProjections - .getShellSnapshot() - .pipe( - Effect.mapError( - (cause) => new ProjectOperationError({ operation: "list-threads", projectId, cause }), - ), - ); - const projectThreads = [...snapshot.threads, ...snapshot.archivedThreads].filter( - (thread) => thread.projectId === projectId, - ); - if (projectThreads.length > 0 && input.force !== true) { - return yield* new ProjectNotEmptyError({ projectId }); + if (existing.value.deletedAt === null) { + yield* deleteChildThreads(input); } - - // Delete children durably before the project. Stable command IDs let a retry - // finish a partially completed cascade without repeating cleanup effects. - yield* Effect.forEach( - projectThreads, - (thread) => - threadCommands - .withLock( - thread.id, - Effect.gen(function* () { - yield* legacyImporter.ensureTranscript(thread.id); - const projection = yield* threadProjections.getThreadRecords(thread.id, [ - "runs", - "attempts", - "nodes", - "runtimeRequests", - "subagents", - "providerSessions", - ]); - if ( - projection.thread.deletedAt !== null || - projection.thread.projectId !== projectId - ) { - return; - } - const command = { - type: "thread.delete" as const, - commandId: CommandId.make(`${input.commandId}:delete-thread:${thread.id}`), - threadId: thread.id, - }; - const now = yield* DateTime.now; - const plan = yield* planThreadDeletion({ - command, - projection, - attachmentIds: yield* threadProjections.getThreadAttachmentIds(thread.id), - now, - idAllocator, - }); - const committed = yield* threadEvents.commitCommand({ - commandId: command.commandId, - commandType: command.type, - threadId: command.threadId, - acceptedAt: now, - events: plan.events, - effects: plan.effects, - }); - if ( - committed.receipt.threadId !== command.threadId || - committed.receipt.commandType !== command.type - ) { - return yield* Effect.fail( - "The thread deletion command ID belongs to a different command.", - ); - } - if (committed.receipt.status === "rejected") { - return yield* Effect.fail( - committed.receipt.error ?? "Thread deletion was previously rejected.", - ); - } - }), - ) - .pipe( - Effect.mapError( - (cause) => - new ProjectOperationError({ operation: "delete-thread", projectId, cause }), - ), - ), - { concurrency: 1, discard: true }, - ); - return yield* dispatch( - projectId, - { - type: "project.delete", - commandId: input.commandId, - projectId, - ...(input.force === undefined ? {} : { force: input.force }), - }, - invalidateEnrichment(existing.value.workspaceRoot).pipe( - Effect.andThen(readCommitted(projectId)), - ), - ); + yield* commit({ type: "project.delete", commandId: input.commandId, projectId }); + yield* projectEnrichment.invalidate([existing.value.workspaceRoot]); + return yield* readCommitted(projectId); }, ); const snapshot = Effect.gen(function* () { - const rows = (yield* readRows()).filter((row) => row.deletedAt === null); - const hydrated = yield* Effect.forEach(rows, hydrateAvailable, { concurrency: 8 }); + const rows = yield* projects + .list() + .pipe( + Effect.mapError( + (cause) => new ProjectOperationError({ operation: "list-projects", cause }), + ), + ); + const hydrated = yield* Effect.forEach(rows, hydrate, { concurrency: 8 }); return { projects: hydrated, updatedAt: DateTime.formatIso(yield* DateTime.now), diff --git a/docs/operations/observability.md b/docs/operations/observability.md index c6537f6eac71..db41058a6d58 100644 --- a/docs/operations/observability.md +++ b/docs/operations/observability.md @@ -318,11 +318,11 @@ jq -r 'select(.traceId == "TRACE_ID_HERE") | [ Filter orchestration commands: ```bash -jq -c 'select(.attributes["orchestration.command_type"] != null) | { +jq -c 'select(.attributes["orchestration_v2.command_type"] != null) | { name, durationMs, - commandType: .attributes["orchestration.command_type"], - aggregateKind: .attributes["orchestration.aggregate_kind"] + commandType: .attributes["orchestration_v2.command_type"], + threadId: .attributes["orchestration_v2.thread_id"] }' "$TRACE_FILE" ``` @@ -365,7 +365,7 @@ Good first searches: `deployment.environment.name` - span names like `sendTurn` or a Git operation such as `GitVcsDriver.statusDetails.status` - Git spans whose `git.operation` attribute identifies the operation -- orchestration spans with attributes like `orchestration.command_type` +- orchestration spans with attributes like `orchestration_v2.command_type` Once you know traces are arriving, narrower TraceQL queries for names such as `sendTurn` or Git operation names become useful. @@ -377,15 +377,12 @@ Traces are best for one request. Metrics are best for trends. Good metric families to watch: - `t3_rpc_request_duration` -- `t3_orchestration_command_duration` -- `t3_orchestration_command_ack_duration` - `t3_provider_turn_duration` - `t3_git_command_duration` Counters tell you volume and failure rate: - `t3_rpc_requests_total` -- `t3_orchestration_commands_total` - `t3_provider_turns_total` - `t3_git_commands_total` @@ -401,21 +398,6 @@ Use traces when the question is: - "which child span caused this one slow interaction?" - "what logs were emitted inside the failing flow?" -### What The New Ack Metric Means - -`t3_orchestration_command_ack_duration` measures: - -- start: command dispatch enters the orchestration engine -- end: the first committed domain event for that command is published by the server - -That is a server-side acknowledgment metric. It does not measure: - -- websocket transit to the browser -- client receipt -- React render time - -If you need those later, add client-side instrumentation or a dedicated server fanout metric. - ## Common Workflows ### "Why did this request fail?" @@ -432,12 +414,6 @@ If you need those later, add client-side instrumentation or a dedicated server f 2. Check child spans for sqlite, git, provider, or terminal work. 3. Look at the matching duration metrics to see whether the slowness is systemic. -### "Did this command take too long to acknowledge?" - -1. Check `t3_orchestration_command_ack_duration` by `commandType`. -2. If it is high, inspect the corresponding orchestration trace. -3. Look at child spans for projection, sqlite, provider, or git work. - ### "Are git hooks causing latency?" 1. Filter `git.operation` spans. @@ -639,7 +615,6 @@ Current high-value span and metric boundaries include: - RPC request metrics in `apps/server/src/observability/RpcInstrumentation.ts` - startup phases - orchestration command processing -- orchestration command acknowledgment latency - provider session and turn operations - git command execution and git hook events - terminal session lifecycle