Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 0 additions & 16 deletions apps/server/src/observability/Metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
{
Expand Down
57 changes: 51 additions & 6 deletions apps/server/src/orchestration-v2/CommandReceiptStore.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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<boolean, CommandReceiptStoreV2Error>;
readonly upsert: (receipt: CommandReceiptV2) => Effect.Effect<void, CommandReceiptStoreV2Error>;
readonly upsert: (
receipt: AnyCommandReceiptV2,
) => Effect.Effect<void, CommandReceiptStoreV2Error>;
/** A thread command's receipt; none when the command id belongs to a project command. */
readonly getByCommandId: (
commandId: CommandId,
) => Effect.Effect<Option.Option<CommandReceiptV2>, CommandReceiptStoreV2Error>;
/** A project command's receipt; none when the command id belongs to a thread command. */
readonly getProjectByCommandId: (
commandId: CommandId,
) => Effect.Effect<Option.Option<ProjectCommandReceiptV2>, CommandReceiptStoreV2Error>;
}

export class CommandReceiptStoreV2 extends Context.Service<
Expand All @@ -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({
Expand All @@ -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,
Expand Down Expand Up @@ -158,6 +186,23 @@ const baseLayer: Layer.Layer<CommandReceiptStoreV2, never, OrchestrationCommandR
}),
),
),
getProjectByCommandId: (commandId) =>
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);
}),
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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, {
Expand Down Expand Up @@ -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({
Expand Down
111 changes: 110 additions & 1 deletion apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
RunId,
RuntimeRequestId,
NodeId,
type ProjectId,
ThreadId,
} from "@t3tools/contracts";
import * as Context from "effect/Context";
Expand All @@ -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 {
Expand All @@ -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,
Expand Down Expand Up @@ -156,6 +160,28 @@ export interface EventSinkV2Shape {
readonly rejectedAt: DateTime.Utc;
readonly error: string;
}) => Effect.Effect<CommandReceiptV2, EventSinkV2Error>;
/**
* 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<ProjectCommandReceiptV2, EventSinkV2Error>;
readonly stream: (input?: {
readonly threadId?: ThreadId;
readonly afterSequence?: number;
Expand Down Expand Up @@ -186,6 +212,7 @@ const baseLayer: Layer.Layer<
| EffectOutboxV2
| EventStoreV2
| ProjectionStoreV2
| ProjectStore.ProjectStoreV2
| SqlClient.SqlClient
| TurnItemPositionStoreV2
> = Layer.effect(
Expand All @@ -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<OrchestrationV2StoredEvent>();
const liveEventsByType = new Map<
Expand Down Expand Up @@ -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<EventSinkV2Shape["commitProjectCommand"]>[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<EventSinkV2Shape["commitRejectedProjectCommand"]>[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;
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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,
),
),
);
Loading
Loading