diff --git a/apps/server/src/observability/Metrics.ts b/apps/server/src/observability/Metrics.ts index b75ace399ccc..edc03e87460a 100644 --- a/apps/server/src/observability/Metrics.ts +++ b/apps/server/src/observability/Metrics.ts @@ -19,13 +19,32 @@ export const rpcRequestDuration = Metric.timer("t3_rpc_request_duration", { description: "RPC request handling duration.", }); -const orchestrationEventsProcessedTotal = Metric.counter( +export const orchestrationEventsProcessedTotal = Metric.counter( "t3_orchestration_events_processed_total", { description: "Total orchestration intent events processed by runtime reactors.", }, ); +export const orchestrationCommandsTotal = Metric.counter("t3_orchestration_commands_total", { + description: "Total orchestration commands dispatched by result.", +}); + +export const orchestrationCommandDuration = Metric.timer("t3_orchestration_command_duration", { + description: + "Orchestration command dispatch duration while holding its dispatch lock, excluding lock wait.", +}); + +export const orchestrationCommandAckDuration = Metric.timer( + "t3_orchestration_command_ack_duration", + { + description: + "Time from before a command acquires its per-thread dispatch lock until its commit is " + + "durable, including lock contention. Distinct from t3_orchestration_command_duration, " + + "which excludes lock wait.", + }, +); + export const orchestrationEffectClaimsTotal = Metric.counter( "t3_orchestration_effect_claims_total", { @@ -38,19 +57,19 @@ export const orchestrationEffectQueueWait = Metric.timer("t3_orchestration_effec "Time from an orchestration effect's temporal availability until claim, including same-thread blocking.", }); -const providerSessionsTotal = Metric.counter("t3_provider_sessions_total", { +export const providerSessionsTotal = Metric.counter("t3_provider_sessions_total", { description: "Total provider session lifecycle operations.", }); -const providerTurnsTotal = Metric.counter("t3_provider_turns_total", { +export const providerTurnsTotal = Metric.counter("t3_provider_turns_total", { description: "Total provider turn lifecycle operations.", }); -const providerTurnDuration = Metric.timer("t3_provider_turn_duration", { +export const providerTurnDuration = Metric.timer("t3_provider_turn_duration", { description: "Provider turn request duration.", }); -const providerRuntimeEventsTotal = Metric.counter("t3_provider_runtime_events_total", { +export const providerRuntimeEventsTotal = Metric.counter("t3_provider_runtime_events_total", { description: "Total canonical provider runtime events processed.", }); @@ -139,13 +158,7 @@ export const withMetrics: { (effect: Effect.Effect, options: WithMetricsOptions): Effect.Effect; } = dual(2, withMetricsImpl); -const providerMetricAttributes = (provider: string, extra?: Readonly>) => - compactMetricAttributes({ - provider, - ...extra, - }); - -const providerTurnMetricAttributes = (input: { +export const providerTurnMetricAttributes = (input: { readonly provider: string; readonly model: string | null | undefined; readonly extra?: Readonly>; diff --git a/apps/server/src/orchestration-v2/EffectWorker.test.ts b/apps/server/src/orchestration-v2/EffectWorker.test.ts index af2667a8cf8e..e33b575a3464 100644 --- a/apps/server/src/orchestration-v2/EffectWorker.test.ts +++ b/apps/server/src/orchestration-v2/EffectWorker.test.ts @@ -14,6 +14,7 @@ import * as Deferred from "effect/Deferred"; 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 Queue from "effect/Queue"; import * as Ref from "effect/Ref"; @@ -721,6 +722,59 @@ it.effect("backs off briefly when a due deadline loses a claim race", () => }).pipe(Effect.provide(TestClock.layer())), ); +it.effect("records a processed-event metric for a successfully executed claim", () => + Effect.gen(function* () { + const now = DateTime.formatIso(yield* DateTime.now); + const effectId = "effect:worker-processed-metric"; + const workerId = "worker-processed-metric"; + const claimedEffect: EffectOutbox.OrchestrationEffectV2 = { + id: effectId, + commandId: CommandId.make("command:worker-processed-metric"), + threadId: ThreadId.make("thread:worker-processed-metric"), + request: { type: "terminal.cleanup" }, + status: "running", + attemptCount: 1, + availableAt: now, + leaseOwner: workerId, + leaseExpiresAt: now, + createdAt: now, + updatedAt: now, + completedAt: null, + lastError: null, + }; + const outboxLayer = Layer.mock(EffectOutbox.EffectOutboxV2)({ + claimNext: () => Effect.succeed(Option.some(claimedEffect)), + get: () => Effect.succeed(Option.some(claimedEffect)), + awaitCancellation: () => Effect.never, + clearCancellation: () => Effect.void, + succeed: () => Effect.succeed(true), + }); + const executorLayer = Layer.succeed( + EffectWorker.OrchestrationEffectExecutorV2, + EffectWorker.OrchestrationEffectExecutorV2.of({ execute: () => Effect.void }), + ); + const workerLayer = EffectWorker.layerWithOptions({ workerId }).pipe( + Layer.provide(Layer.merge(outboxLayer, executorLayer)), + ); + + const exit = yield* EffectWorker.OrchestrationEffectWorkerV2.pipe( + Effect.flatMap((worker) => worker.runOnce), + Effect.provide(workerLayer), + Effect.exit, + ); + + assert.isTrue(Exit.isSuccess(exit)); + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_orchestration_events_processed_total" && + snapshot.attributes?.eventType === "terminal.cleanup", + ), + ); + }), +); + it.effect("safely retries after replacement cleanup succeeds and start fails", () => Effect.gen(function* () { const now = yield* DateTime.now; diff --git a/apps/server/src/orchestration-v2/EffectWorker.ts b/apps/server/src/orchestration-v2/EffectWorker.ts index 1abf070b87ec..76a0e4533c40 100644 --- a/apps/server/src/orchestration-v2/EffectWorker.ts +++ b/apps/server/src/orchestration-v2/EffectWorker.ts @@ -15,6 +15,7 @@ import { metricAttributes, orchestrationEffectClaimsTotal, orchestrationEffectQueueWait, + orchestrationEventsProcessedTotal, } from "../observability/Metrics.ts"; import * as RunFinalizationService from "./RunFinalizationService.ts"; import * as ResourceCleanupService from "./ResourceCleanupService.ts"; @@ -646,6 +647,10 @@ export const layerWithOptions = ( }).pipe(Effect.onError((cause) => requeueClaim(effect, cause))); if (cancelledBeforeExecution) return true; + yield* increment(orchestrationEventsProcessedTotal, { + eventType: effect.request.type, + }); + const execution = executor .execute(effect, { willRetry: effect.attemptCount < maxAttempts }) .pipe(Effect.as("executed" as const)); diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index ad41ae6ffde6..14c24d169f79 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -20,6 +20,7 @@ import { import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Metric from "effect/Metric"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; @@ -534,3 +535,48 @@ it.effect("settles only the stopped run's background work, once", () => ]); }).pipe(Effect.provide(testLayer)), ); + +it.effect("records a command, duration and ack metric for a dispatched command", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const threadId = ThreadId.make("thread:dispatch-metrics"); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create-dispatch-metrics"), + threadId, + projectId: ProjectId.make("project:dispatch-metrics"), + title: "Dispatch metrics", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_orchestration_commands_total" && + snapshot.attributes?.commandType === "thread.create" && + snapshot.attributes?.outcome === "success", + ), + ); + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_orchestration_command_duration" && + snapshot.attributes?.commandType === "thread.create", + ), + ); + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_orchestration_command_ack_duration" && + snapshot.attributes?.ackEventType === "thread.created", + ), + ); + }).pipe(Effect.provide(testLayer)), +); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index c58dd46609ec..07dbef477627 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -50,10 +50,13 @@ import { derivePendingBackgroundWork, pendingBackgroundTurnItems, } from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; +import * as Clock from "effect/Clock"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; +import * as Metric from "effect/Metric"; import * as Path from "effect/Path"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -96,6 +99,13 @@ import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import { ProviderAdapterRegistryV2 } from "./ProviderAdapterRegistry.ts"; import { ProviderContinuationRequests } from "./ProviderContinuationRequests.ts"; import { makeProviderFailure } from "./ProviderFailure.ts"; +import { + metricAttributes, + orchestrationCommandAckDuration, + orchestrationCommandDuration, + orchestrationCommandsTotal, + withMetrics, +} from "../observability/Metrics.ts"; import { ProviderSessionManagerV2 } from "./ProviderSessionManager.ts"; import { ProviderSwitchServiceV2 } from "./ProviderSwitchService.ts"; import { isAutomaticCompletionRun, queuedRunsInDeliveryOrder } from "./QueuedRunOrder.ts"; @@ -9277,6 +9287,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const dispatchWithReceiptEffect = Effect.fn("orchestrationV2.dispatch.withReceipt")(function* ( command: OrchestrationV2ServerCommand, + dispatchStartedAtMs: number, ): Effect.fn.Return { yield* Effect.annotateCurrentSpan({ "orchestration_v2.command_id": command.commandId, @@ -9446,6 +9457,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio detail: committed.receipt.error ?? "Previously rejected.", }); } + const ackEventType = committed.storedEvents[0]?.event.type; + if (ackEventType !== undefined) { + yield* Metric.update( + Metric.withAttributes(orchestrationCommandAckDuration, metricAttributes({ ackEventType })), + Duration.millis(Math.max(0, (yield* Clock.currentTimeMillis) - dispatchStartedAtMs)), + ); + } if (command.type === "queue.resume") { yield* mapDispatchError(command)(startNextQueuedRun(command.threadId)); } @@ -9463,7 +9481,19 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); const dispatchWithReceipt = (command: OrchestrationV2ServerCommand) => - threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command)); + Effect.gen(function* () { + const dispatchStartedAtMs = yield* Clock.currentTimeMillis; + return yield* threadDispatch.withLock( + commandThreadId(command), + dispatchWithReceiptEffect(command, dispatchStartedAtMs).pipe( + withMetrics({ + counter: orchestrationCommandsTotal, + timer: orchestrationCommandDuration, + attributes: { commandType: command.type }, + }), + ), + ); + }); const handleTerminalRun = (stored: OrchestrationV2StoredEvent) => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts index f1470125a889..e296c27400ba 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts @@ -24,6 +24,7 @@ import * as Deferred from "effect/Deferred"; import * as Fiber from "effect/Fiber"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Metric from "effect/Metric"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; @@ -467,6 +468,120 @@ layer("ProviderEventIngestorV2", (it) => { }), ); + it.effect("records a provider runtime event metric for every normalized provider event", () => + Effect.gen(function* () { + const ingestor = yield* ProviderEventIngestor.ProviderEventIngestorV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const projectId = yield* idAllocator.allocate.project({ + fixtureName: "provider-event-runtime-metric", + }); + const threadId = yield* idAllocator.allocate.thread({ + fixtureName: "provider-event-runtime-metric", + projectId, + }); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + + yield* ingestor.normalize({ + providerSessionId, + providerInstanceId: modelSelection.instanceId, + threadId, + event: { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId: idAllocator.derive.providerThread({ + driver: CODEX_DRIVER, + nativeThreadId: "native-thread-runtime-metric", + }), + providerTurnId: idAllocator.derive.providerTurn({ + driver: CODEX_DRIVER, + nativeTurnId: "native-turn-runtime-metric", + }), + runOrdinal: 1, + status: "completed", + failure: null, + threadDisposition: "reusable", + }, + }); + + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_provider_runtime_events_total" && + snapshot.attributes?.provider === CODEX_DRIVER && + snapshot.attributes?.eventType === "turn.terminal", + ), + ); + }), + ); + + it.effect( + "records an additional failed-terminal metric when a provider turn ends in failure", + () => + Effect.gen(function* () { + const now = yield* DateTime.now; + const eventSink = yield* EventSink.EventSinkV2; + const ingestor = yield* ProviderEventIngestor.ProviderEventIngestorV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const threadEvent = yield* threadCreatedEvent(now); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId: threadEvent.threadId, + }); + const providerThreadId = idAllocator.derive.providerThread({ + driver: CODEX_DRIVER, + nativeThreadId: "native-thread-failed-metric", + }); + const providerTurnId = idAllocator.derive.providerTurn({ + driver: CODEX_DRIVER, + nativeTurnId: "native-turn-failed-metric", + }); + + yield* eventSink.write({ events: [threadEvent] }); + yield* ingestor.ingestNormalized({ + providerSessionId, + providerInstanceId: modelSelection.instanceId, + threadId: threadEvent.threadId, + event: { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId, + providerTurnId, + runOrdinal: 1, + failureItemOrdinal: 102, + status: "failed", + failure: makeProviderFailure({ + message: "Invalid reasoning effort.", + code: "invalid_request", + class: "validation_error", + }), + threadDisposition: "reusable", + }, + }); + + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_provider_runtime_events_total" && + snapshot.attributes?.provider === CODEX_DRIVER && + snapshot.attributes?.eventType === "turn.terminal", + ), + ); + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_provider_runtime_events_total" && + snapshot.attributes?.provider === CODEX_DRIVER && + snapshot.attributes?.eventType === "turn.terminal.failed", + ), + ); + }), + ); + it.effect("persists an interrupted run's inherited terminal through the live run router", () => Effect.gen(function* () { const now = yield* DateTime.now; diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts index 8082962b86de..412d4b8c6652 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts @@ -25,6 +25,7 @@ import * as Layer from "effect/Layer"; import * as Schema from "effect/Schema"; import { getModelSelectionStringOptionValue } from "@t3tools/shared/model"; +import { increment, providerRuntimeEventsTotal } from "../observability/Metrics.ts"; import * as AnalyticsService from "../telemetry/AnalyticsService.ts"; import * as EventSink from "./EventSink.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; @@ -340,6 +341,10 @@ export const layer: Layer.Layer< const normalize: ProviderEventIngestorV2Shape["normalize"] = (input) => Effect.gen(function* () { + yield* increment(providerRuntimeEventsTotal, { + provider: input.event.driver, + eventType: input.event.type, + }); switch (input.event.type) { case "app_thread.created": return [ @@ -462,6 +467,10 @@ export const layer: Layer.Layer< if (input.event.status !== "failed") { return dismissed; } + yield* increment(providerRuntimeEventsTotal, { + provider: input.event.driver, + eventType: "turn.terminal.failed", + }); const occurredAt = yield* DateTime.now; return [ ...dismissed, diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 63d4c57fb663..cbc0e4360417 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -23,6 +23,7 @@ import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; +import * as Metric from "effect/Metric"; import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; @@ -833,6 +834,57 @@ it.effect("ProviderSessionManagerV2 opens a duplicate session only once", () => }), ); +it.effect("ProviderSessionManagerV2 records provider session lifecycle metrics", () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread-provider-session-manager-metrics"); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + + yield* eventSink.write({ + events: [yield* makeThreadCreatedEvent({ idAllocator, threadId, now })], + }); + const sessionCount = (operation: string) => + Metric.snapshot.pipe( + Effect.map((snapshots) => { + const snapshot = snapshots.find( + (candidate) => + candidate.id === "t3_provider_sessions_total" && + candidate.attributes?.provider === CODEX_DRIVER && + candidate.attributes?.operation === operation && + candidate.attributes?.outcome === "success", + ); + return snapshot !== undefined && "count" in snapshot.state + ? Number(snapshot.state.count) + : 0; + }), + ); + const startsBefore = yield* sessionCount("start"); + const stopsBefore = yield* sessionCount("stop"); + + const open = manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + yield* open; + yield* open; + yield* manager.close(providerSessionId); + yield* manager.close(providerSessionId); + yield* open; + yield* manager.close(providerSessionId); + + assert.equal((yield* sessionCount("start")) - startsBefore, 2); + assert.equal((yield* sessionCount("stop")) - stopsBefore, 2); + }); + + yield* effect.pipe(Effect.provide(makeTestLayer({ state, idleTimeoutMs: 60_000 }))); + }), +); + it.effect("ProviderSessionManagerV2 releases live sessions when its layer shuts down", () => Effect.gen(function* () { const state = yield* Ref.make(emptyState); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 786eb5a6bf0a..d3a3b83f1d8a 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -32,6 +32,8 @@ import * as ProjectService from "../project/ProjectService.ts"; import * as McpProviderSession from "../mcp/McpProviderSession.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as McpSessionRegistry from "../mcp/McpSessionRegistry.ts"; +import { outcomeFromExit } from "../observability/Attributes.ts"; +import { increment, providerSessionsTotal } from "../observability/Metrics.ts"; import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "./IdAllocator.ts"; import { makeKeyedSerialExecutor } from "./KeyedSerialExecutor.ts"; @@ -806,7 +808,15 @@ export const layerWithOptions = ( if (Option.isSome(closeExit) && Exit.isFailure(closeExit.value)) { return yield* Effect.failCause(closeExit.value.cause); } - }), + }).pipe( + Effect.onExit((exit) => + increment(providerSessionsTotal, { + provider: entry.runtime.driver, + operation: "stop", + outcome: outcomeFromExit(exit), + }), + ), + ), }), (entry) => Option.match(entry, { @@ -1544,6 +1554,29 @@ export const layerWithOptions = ( }); yield* Effect.addFinalizer(() => shutdown); + // Reuse of a live session is not a start, so only fresh opens are counted. + const recordSessionOpen = + (input: Parameters[0]) => + (effect: Effect.Effect) => + Effect.gen(function* () { + if ((yield* Ref.get(sessions)).has(sessionKey(input.providerSessionId))) { + return yield* effect; + } + const exit = yield* Effect.exit(effect); + const provider = Exit.isSuccess(exit) + ? exit.value.driver + : yield* registry.get(input.modelSelection.instanceId).pipe( + Effect.map((adapter) => adapter.driver), + Effect.orElseSucceed(() => undefined), + ); + yield* increment(providerSessionsTotal, { + provider, + operation: input.resumeFromSession === undefined ? "start" : "recover", + outcome: outcomeFromExit(exit), + }); + return Exit.isSuccess(exit) ? exit.value : yield* Effect.failCause(exit.cause); + }); + return ProviderSessionManagerV2.of({ shutdown, open: (input) => @@ -1711,7 +1744,7 @@ export const layerWithOptions = ( yield* startEventPump(entry); yield* scheduleIdleRelease(input.providerSessionId); return exposedRuntime; - }), + }).pipe(recordSessionOpen(input)), ), get: (providerSessionId) => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts index eba830319aaf..b193caa40b42 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnControlService.test.ts @@ -18,6 +18,7 @@ import * as DateTime from "effect/DateTime"; 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 Ref from "effect/Ref"; import * as Stream from "effect/Stream"; @@ -322,3 +323,159 @@ it.effect( assert.equal(interrupted?.nativeThreadRef?.nativeId, "native-thread:restart-session"); }), ); + +it.effect("records a provider turn metric for a successful interrupt", () => + Effect.gen(function* () { + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:interrupt-metrics"); + const sessionId = ProviderSessionId.make("provider-session:interrupt-metrics"); + const providerThreadId = ProviderThreadId.make("provider-thread:interrupt-metrics"); + const providerTurnId = ProviderTurnId.make("provider-turn:interrupt-metrics"); + const attemptId = RunAttemptId.make("run-attempt:interrupt-metrics"); + const providerThread: OrchestrationV2ProviderThread = { + id: providerThreadId, + driver, + providerInstanceId, + providerSessionId: sessionId, + appThreadId: threadId, + ownerNodeId: null, + nativeThreadRef: { + driver, + nativeId: "native-thread:interrupt-metrics", + strength: "strong", + }, + nativeConversationHeadRef: null, + status: "active", + firstRunOrdinal: 1, + lastRunOrdinal: 1, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }; + const projection = yield* Ref.make( + makeProjection({ now, threadId, providerThread, providerTurnId, attemptId }), + ); + const providerSession = { + id: sessionId, + driver, + providerInstanceId, + status: "running" as const, + cwd: "/workspace", + model: modelSelection.model, + capabilities: CodexProviderCapabilitiesV2, + createdAt: now, + updatedAt: now, + lastError: null, + }; + const runtime: ProviderAdapterV2SessionRuntime = { + instanceId: providerInstanceId, + driver, + providerSessionId: sessionId, + providerSession, + events: Stream.empty, + ensureThread: () => Effect.die("unused ensureThread"), + resumeThread: () => Effect.die("unused resumeThread"), + startTurn: () => Effect.die("unused startTurn"), + steerTurn: () => Effect.die("unused steerTurn"), + interruptTurn: () => + Ref.update(projection, (current) => ({ + ...current, + providerTurns: current.providerTurns.map((turn) => + turn.id === providerTurnId + ? { ...turn, status: "interrupted" as const, completedAt: now } + : turn, + ), + })), + respondToRuntimeRequest: () => Effect.die("unused respondToRuntimeRequest"), + readThreadSnapshot: () => Effect.die("unused readThreadSnapshot"), + rollbackThread: () => Effect.die("unused rollbackThread"), + forkThread: () => Effect.die("unused forkThread"), + }; + const projectionLayer = Layer.succeed( + ProjectionStore.ProjectionStoreV2, + ProjectionStore.ProjectionStoreV2.of({ + apply: () => Effect.void, + getLimitRecoveryCandidates: () => Effect.die("unused getLimitRecoveryCandidates"), + getShellSnapshot: () => Effect.die("unused getShellSnapshot"), + getThreadShell: () => Effect.die("unused getThreadShell"), + getThread: () => Ref.get(projection).pipe(Effect.map((state) => state.thread)), + getSettlementCandidates: () => Effect.die("unused getSettlementCandidates"), + getThreadsWithPullRequests: () => Effect.die("unused getThreadsWithPullRequests"), + getThreadProjection: () => Effect.die("control effects must not load transcript"), + getTurnStartContext: () => Effect.die("unused"), + getTurnStartHistory: () => Effect.die("unused"), + getRuntimeRecoveryProjection: () => Effect.die("unused getRuntimeRecoveryProjection"), + getPlan: () => Effect.die("unused"), + hasUnpairedRunInterruptRequest: () => Effect.die("unused interrupt read"), + getThreadAttachmentIds: () => Effect.die("Unused attachment lookup"), + getTimelinePage: () => Effect.die("Unused timeline read"), + getMessageCount: () => Effect.die("unused message count"), + getNextTurnItemOrdinal: () => Effect.die("unused ordinal read"), + getThreadRecords: () => Effect.die("unused record read"), + getRuntimeRequest: () => Effect.die("unused getRuntimeRequest"), + getRunningTurnContext: () => Effect.die("unused getRunningTurnContext"), + getThreadProviderContext: () => Effect.die("unused getThreadProviderContext"), + getRuntimeResponseContext: () => Effect.die("unused getRuntimeResponseContext"), + getPendingNativeUserInputs: () => Effect.die("unused getPendingNativeUserInputs"), + getProviderControlContext: (_threadId, target) => + Ref.get(projection).pipe( + Effect.map((current) => ({ + providerThread: current.providerThreads.find( + (thread) => thread.id === target.providerThreadId, + ), + providerTurn: current.providerTurns.find((turn) => turn.id === target.providerTurnId), + attempt: current.attempts.find((attempt) => attempt.id === target.attemptId), + message: undefined, + run: undefined, + })), + ), + getCheckpointContext: () => Effect.die("not used"), + getCheckpointCaptureContext: () => Effect.die("not used"), + getRunMessage: () => Effect.die("not used"), + canStartQueuedRun: () => Effect.die("not used"), + getRecoveryThreadIds: () => Effect.die("unused getRecoveryThreadIds"), + getUnreadableThreadIds: () => Effect.die("unused getUnreadableThreadIds"), + getThreadSnapshot: () => Effect.die("unused getThreadSnapshot"), + getThreadSnapshotWindow: () => Effect.die("unused getThreadSnapshotWindow"), + }), + ); + const sessionManagerLayer = Layer.succeed( + ProviderSessionManager.ProviderSessionManagerV2, + ProviderSessionManager.ProviderSessionManagerV2.of({ + shutdown: Effect.void, + open: () => Effect.die("unused open"), + get: (providerSessionId) => + Effect.succeed(providerSessionId === sessionId ? Option.some(runtime) : Option.none()), + close: () => Effect.void, + closeInstance: () => Effect.void, + release: () => Effect.void, + detach: () => Effect.void, + }), + ); + const controlLayer = ProviderTurnControlService.layer.pipe( + Layer.provide(Layer.merge(projectionLayer, sessionManagerLayer)), + ); + + yield* Effect.gen(function* () { + const control = yield* ProviderTurnControlService.ProviderTurnControlServiceV2; + yield* control.interrupt({ + threadId, + providerSessionId: sessionId, + providerThreadId, + providerTurnId, + }); + }).pipe(Effect.provide(controlLayer)); + + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_provider_turns_total" && + snapshot.attributes?.provider === driver && + snapshot.attributes?.operation === "interrupt" && + snapshot.attributes?.outcome === "success", + ), + ); + }), +); diff --git a/apps/server/src/orchestration-v2/ProviderTurnControlService.ts b/apps/server/src/orchestration-v2/ProviderTurnControlService.ts index 47f7aa80bdc4..32c9ba8d9a37 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnControlService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnControlService.ts @@ -13,6 +13,7 @@ import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; +import { providerTurnsTotal, withMetrics } from "../observability/Metrics.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import * as ProviderSessionManager from "./ProviderSessionManager.ts"; @@ -169,9 +170,11 @@ export const layer: Layer.Layer< }); return ProviderTurnControlServiceV2.of({ - interrupt: (input) => - Effect.gen(function* () { + interrupt: (input) => { + let driver: string | undefined; + return Effect.gen(function* () { const loaded = yield* load({ ...input, operation: "interrupt" }); + driver = loaded.providerThread.driver; const session = Option.isSome(loaded.session) ? loaded.session : yield* sessions.get(input.providerSessionId); @@ -186,6 +189,13 @@ export const layer: Layer.Layer< requestRuntimeRestart: true, }); }).pipe( + withMetrics({ + counter: providerTurnsTotal, + attributes: () => ({ + ...(driver === undefined ? {} : { provider: driver }), + operation: "interrupt", + }), + }), Effect.mapError((cause) => isProviderTurnControlError(cause) ? cause @@ -196,10 +206,13 @@ export const layer: Layer.Layer< cause, }), ), - ), - interruptAndAwaitTerminal: (input) => - Effect.gen(function* () { + ); + }, + interruptAndAwaitTerminal: (input) => { + let driver: string | undefined; + return Effect.gen(function* () { const loaded = yield* load({ ...input, operation: "restart" }); + driver = loaded.providerThread.driver; if (Option.isNone(loaded.session)) { // No live adapter: nothing can emit a terminal provider-turn update // from interrupt. Do not poll for projection terminalization or the @@ -252,6 +265,13 @@ export const layer: Layer.Layer< cause: `Provider turn ${input.providerTurnId} did not terminalize before restart.`, }); }).pipe( + withMetrics({ + counter: providerTurnsTotal, + attributes: () => ({ + ...(driver === undefined ? {} : { provider: driver }), + operation: "restart", + }), + }), Effect.mapError((cause) => isProviderTurnControlError(cause) ? cause @@ -262,9 +282,11 @@ export const layer: Layer.Layer< cause, }), ), - ), - steer: (input) => - Effect.gen(function* () { + ); + }, + steer: (input) => { + let driver: string | undefined; + return Effect.gen(function* () { const context = yield* projections.getProviderControlContext(input.threadId, input); const ownership = context.message?.delegatedCompletion; if (ownership !== undefined) { @@ -283,6 +305,7 @@ export const layer: Layer.Layer< return; } const loaded = yield* load({ ...input, operation: "steer" }); + driver = loaded.providerThread.driver; if (Option.isNone(loaded.session)) return; const { message, run } = loaded.context; if (message === undefined || run === undefined) { @@ -334,6 +357,13 @@ export const layer: Layer.Layer< ), ); }).pipe( + withMetrics({ + counter: providerTurnsTotal, + attributes: () => ({ + ...(driver === undefined ? {} : { provider: driver }), + operation: "steer", + }), + }), Effect.mapError((cause) => isProviderTurnControlError(cause) ? cause @@ -344,7 +374,8 @@ export const layer: Layer.Layer< cause, }), ), - ), + ); + }, }); }), ); diff --git a/apps/server/src/orchestration-v2/RunExecutionService.test.ts b/apps/server/src/orchestration-v2/RunExecutionService.test.ts index ec292a02b86f..4a70aa5e84b1 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.test.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.test.ts @@ -33,6 +33,7 @@ import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Metric from "effect/Metric"; import * as Option from "effect/Option"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; @@ -590,6 +591,86 @@ it.effect("rechecks run ownership immediately before calling the provider", () = }).pipe(Effect.provide(RunExecutionTestLayer)), ); +it.effect("records a provider turn metric for a successful send", () => + Effect.gen(function* () { + const runExecution = yield* RunExecutionService.RunExecutionServiceV2; + const threadId = ThreadId.make("thread:run-execution-send-metrics"); + const runId = RunId.make("run:run-execution-send-metrics"); + const attemptId = RunAttemptId.make("attempt:run-execution-send-metrics"); + const providerThreadId = ProviderThreadId.make("provider-thread:run-execution-send-metrics"); + const providerInstanceId = ProviderInstanceId.make("codex"); + const providerSessionId = ProviderSessionId.make("session:run-execution-send-metrics"); + const rootNodeId = NodeId.make("node:run-execution-send-metrics"); + + yield* runExecution.startRootRun({ + commandId: CommandId.make("command:run-execution-send-metrics"), + appThread: { id: threadId } as OrchestrationV2AppThread, + providerSessionId, + session: { + driver, + events: Stream.never, + startTurn: () => Effect.void, + } as unknown as ProviderAdapterV2SessionRuntime, + run: { + id: runId, + threadId, + ordinal: 1, + providerInstanceId, + } as OrchestrationV2Run, + rootNode: { id: rootNodeId } as OrchestrationV2ExecutionNode, + checkpointScope: { + id: CheckpointScopeId.make("checkpoint-scope:run-execution-send-metrics"), + } as OrchestrationV2CheckpointScope, + providerThread: { + id: providerThreadId, + driver, + } as OrchestrationV2ProviderThread, + attempt: { + id: attemptId, + providerTurnId: null, + } as OrchestrationV2RunAttempt, + attemptId, + providerTurnOrdinal: 1, + message: { + messageId: MessageId.make("message:run-execution-send-metrics"), + text: "Record a send metric.", + attachments: [], + createdBy: "user", + creationSource: "web", + }, + modelSelection: { instanceId: providerInstanceId, model: "gpt-5.4" }, + runtimePolicy: { + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + approvalPolicy: "never", + sandboxPolicy: { + type: "readOnly", + access: { type: "fullAccess" }, + networkAccess: false, + }, + }, + }); + + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_provider_turns_total" && + snapshot.attributes?.provider === driver && + snapshot.attributes?.operation === "send" && + snapshot.attributes?.outcome === "success", + ), + ); + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_provider_turn_duration" && snapshot.attributes?.provider === driver, + ), + ); + }).pipe(Effect.provide(RunExecutionTestLayer)), +); + it.effect( "dispatches only attachment-free compact commands through the native compaction path", () => diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 71211c126545..f93c556738a5 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -35,6 +35,12 @@ import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import * as McpSessionRegistry from "../mcp/McpSessionRegistry.ts"; +import { + providerTurnDuration, + providerTurnMetricAttributes, + providerTurnsTotal, + withMetrics, +} from "../observability/Metrics.ts"; import * as ServerSettings from "../serverSettings.ts"; import * as CheckpointService from "./CheckpointService.ts"; import * as EventSink from "./EventSink.ts"; @@ -1374,6 +1380,15 @@ export const layer: Layer.Layer< )) : input.session.startTurn(turnInput); yield* startTurn.pipe( + withMetrics({ + counter: providerTurnsTotal, + timer: providerTurnDuration, + attributes: providerTurnMetricAttributes({ + provider: input.session.driver, + model: input.modelSelection.model, + extra: { operation: "send" }, + }), + }), Effect.catchCause((cause) => Effect.logError("orchestration V2 provider turn start failed", { runId: input.run.id, diff --git a/apps/server/src/project/ProjectService.test.ts b/apps/server/src/project/ProjectService.test.ts index 98e576a6fe9b..c1d2ac4e99a5 100644 --- a/apps/server/src/project/ProjectService.test.ts +++ b/apps/server/src/project/ProjectService.test.ts @@ -5,6 +5,7 @@ import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Metric from "effect/Metric"; import * as Option from "effect/Option"; import * as Ref from "effect/Ref"; import { TestClock } from "effect/testing"; @@ -227,6 +228,36 @@ it.layer(TestLayer)("ProjectService", (it) => { }), ); + it.effect("records a command metric for a committed project command", () => + Effect.gen(function* () { + const service = yield* ProjectService.ProjectService; + yield* TestClock.setTime(Date.parse("2026-06-20T10:00:00.000Z")); + yield* service.create({ + commandId: CommandId.make("command:metrics:create"), + projectId: ProjectId.make("project:metrics"), + title: "Metrics", + workspaceRoot: "/work/metrics", + }); + + const snapshots = yield* Metric.snapshot; + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_orchestration_commands_total" && + snapshot.attributes?.commandType === "project.create" && + snapshot.attributes?.outcome === "success", + ), + ); + assert.isTrue( + snapshots.some( + (snapshot) => + snapshot.id === "t3_orchestration_command_duration" && + snapshot.attributes?.commandType === "project.create", + ), + ); + }), + ); + it.effect("auto-bootstraps a workspace exactly once", () => Effect.gen(function* () { const service = yield* ProjectService.ProjectService; diff --git a/apps/server/src/project/ProjectService.ts b/apps/server/src/project/ProjectService.ts index 5d0e6b68e626..f31de0ce1d05 100644 --- a/apps/server/src/project/ProjectService.ts +++ b/apps/server/src/project/ProjectService.ts @@ -16,6 +16,11 @@ import * as Option from "effect/Option"; import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; +import { + orchestrationCommandDuration, + orchestrationCommandsTotal, + withMetrics, +} from "../observability/Metrics.ts"; import * as EventSink from "../orchestration-v2/EventSink.ts"; import * as IdAllocator from "../orchestration-v2/IdAllocator.ts"; import { makeKeyedSerialExecutor } from "../orchestration-v2/KeyedSerialExecutor.ts"; @@ -219,7 +224,7 @@ export const make = Effect.gen(function* () { * 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 commitEffect = Effect.fn("ProjectService.commit")(function* (command: ProjectCommand) { const { projectId } = command; const dispatchError = (cause: unknown) => new ProjectOperationError({ operation: "dispatch-project-command", projectId, cause }); @@ -262,12 +267,18 @@ export const make = Effect.gen(function* () { error: encodeProjectCommandRejection(planned.failure), }); }); + const timedPlanAndCommit = planAndCommit.pipe( + withMetrics({ + timer: orchestrationCommandDuration, + attributes: { commandType: command.type }, + }), + ); const receipt = yield* projectLocks .withLock( projectId, workspaceRoot === undefined - ? planAndCommit - : workspaceLocks.withLock(workspaceRoot, planAndCommit), + ? timedPlanAndCommit + : workspaceLocks.withLock(workspaceRoot, timedPlanAndCommit), ) .pipe(Effect.mapError(dispatchError)); if (receipt.projectId !== projectId || receipt.commandType !== command.type) { @@ -295,6 +306,14 @@ export const make = Effect.gen(function* () { } }); + const commit = (command: ProjectCommand) => + commitEffect(command).pipe( + withMetrics({ + counter: orchestrationCommandsTotal, + attributes: { commandType: command.type }, + }), + ); + const readCommitted = Effect.fn("ProjectService.readCommitted")(function* (projectId: ProjectId) { const row = yield* readRow(projectId, { includeDeleted: true }); if (Option.isNone(row)) {