From 5b23fcb22dbd19ef4c535c47f26fe908073551fb Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 06:12:29 -0400 Subject: [PATCH 1/5] fix(server): orchestrator records turn, session, and command metrics again Deleting the old orchestration stack dropped every t3_provider_* and t3_orchestration_* metric with it, leaving the Grafana dashboard empty with no signal that anything had regressed. Signed-off-by: Yordis Prieto --- apps/server/src/observability/Metrics.ts | 35 +- .../src/orchestration-v2/EffectWorker.test.ts | 54 +++ .../src/orchestration-v2/EffectWorker.ts | 5 + .../Orchestrator.control-reads.test.ts | 46 +++ .../src/orchestration-v2/Orchestrator.ts | 32 +- .../ProviderEventIngestor.test.ts | 115 ++++++ .../orchestration-v2/ProviderEventIngestor.ts | 9 + .../ProviderSessionManager.test.ts | 49 +++ .../ProviderSessionManager.ts | 339 ++++++++++-------- .../ProviderTurnControlService.test.ts | 157 ++++++++ .../ProviderTurnControlService.ts | 49 ++- .../RunExecutionService.test.ts | 81 +++++ .../orchestration-v2/RunExecutionService.ts | 15 + .../server/src/project/ProjectService.test.ts | 31 ++ apps/server/src/project/ProjectService.ts | 16 +- 15 files changed, 861 insertions(+), 172 deletions(-) diff --git a/apps/server/src/observability/Metrics.ts b/apps/server/src/observability/Metrics.ts index b75ace399ccc..c1c0aba06b94 100644 --- a/apps/server/src/observability/Metrics.ts +++ b/apps/server/src/observability/Metrics.ts @@ -19,13 +19,31 @@ 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, from receipt lookup through commit.", +}); + +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 +56,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 +157,16 @@ export const withMetrics: { (effect: Effect.Effect, options: WithMetricsOptions): Effect.Effect; } = dual(2, withMetricsImpl); -const providerMetricAttributes = (provider: string, extra?: Readonly>) => +export 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..261557f0672b 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,54 @@ 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"); + + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + yield* manager.close(providerSessionId); + yield* manager.close(providerSessionId); + + assert.equal((yield* sessionCount("start")) - startsBefore, 1); + assert.equal((yield* sessionCount("stop")) - stopsBefore, 1); + }); + + 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..5cdbe4622c6b 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"; @@ -715,8 +717,9 @@ export const layerWithOptions = ( readonly cancelIdleFiber?: boolean; readonly onlyIfIdleGeneration?: number; readonly gracefulSubscribers?: boolean; - }) => - Effect.acquireUseRelease( + }) => { + let releasedDriver: string | undefined; + return Effect.acquireUseRelease( Ref.modify(sessions, (current) => { const key = sessionKey(input.providerSessionId); const existing = current.get(key); @@ -731,6 +734,7 @@ export const layerWithOptions = ( } const updated = new Map(current); updated.delete(key); + releasedDriver = existing.runtime.driver; return [Option.some(existing), updated] as const; }), (entry) => @@ -857,7 +861,17 @@ export const layerWithOptions = ( }), ), ), + Effect.onExit((exit) => + releasedDriver === undefined + ? Effect.void + : increment(providerSessionsTotal, { + provider: releasedDriver, + operation: "stop", + outcome: outcomeFromExit(exit), + }), + ), ); + }; // Annotated to break the releaseIfStillIdle <-> scheduleIdleReleaseInternal // inference cycle introduced by the pin re-arm below. @@ -1546,108 +1560,50 @@ export const layerWithOptions = ( return ProviderSessionManagerV2.of({ shutdown, - open: (input) => - sessionOpen.withLock( - input.providerSessionId, - Effect.gen(function* () { - const cwd = input.runtimePolicy.cwd; - if (cwd !== null) { - const workspaceIsDirectory = yield* fileSystem.stat(cwd).pipe( - Effect.map((stat) => stat.type === "Directory"), - Effect.catch((error) => Effect.succeed(error.reason._tag !== "NotFound")), - ); - if (!workspaceIsDirectory) { - return yield* new ProviderWorkspaceMissingError({ - threadId: input.threadId, - cwd, - }); + open: (input) => { + let openedDriver: string | undefined; + let reusedExisting = false; + return sessionOpen + .withLock( + input.providerSessionId, + Effect.gen(function* () { + const cwd = input.runtimePolicy.cwd; + if (cwd !== null) { + const workspaceIsDirectory = yield* fileSystem.stat(cwd).pipe( + Effect.map((stat) => stat.type === "Directory"), + Effect.catch((error) => Effect.succeed(error.reason._tag !== "NotFound")), + ); + if (!workspaceIsDirectory) { + return yield* new ProviderWorkspaceMissingError({ + threadId: input.threadId, + cwd, + }); + } } - } - const key = sessionKey(input.providerSessionId); - const existing = (yield* Ref.get(sessions)).get(key); - if (existing !== undefined) { - if ( - !existing.attachedThreadIds.has(input.threadId) && - !existing.supportsMultipleProviderThreads - ) { - return yield* new ProviderSessionOpenError({ - instanceId: input.modelSelection.instanceId, + const key = sessionKey(input.providerSessionId); + const existing = (yield* Ref.get(sessions)).get(key); + if (existing !== undefined) { + if ( + !existing.attachedThreadIds.has(input.threadId) && + !existing.supportsMultipleProviderThreads + ) { + return yield* new ProviderSessionOpenError({ + instanceId: input.modelSelection.instanceId, + providerSessionId: input.providerSessionId, + cause: `Provider ${existing.runtime.driver} does not support attaching multiple app threads to one session.`, + }); + } + yield* ensureThreadAttached({ providerSessionId: input.providerSessionId, - cause: `Provider ${existing.runtime.driver} does not support attaching multiple app threads to one session.`, + threadId: input.threadId, + providerInstanceId: existing.runtime.instanceId, }); + yield* touchActivity(input.providerSessionId); + reusedExisting = true; + return existing.exposedRuntime; } - yield* ensureThreadAttached({ - providerSessionId: input.providerSessionId, - threadId: input.threadId, - providerInstanceId: existing.runtime.instanceId, - }); - yield* touchActivity(input.providerSessionId); - return existing.exposedRuntime; - } - const adapter = yield* registry.get(input.modelSelection.instanceId).pipe( - Effect.mapError( - (cause) => - new ProviderSessionOpenError({ - instanceId: input.modelSelection.instanceId, - providerSessionId: input.providerSessionId, - cause, - }), - ), - ); - const prepared = yield* prepareMcpSession( - input.threadId, - input.modelSelection.instanceId, - ); - const mcpCredentialId = prepared.mcpCredentialId; - // The reservation from prepare protects the credential (which - // eager adapters bake into the provider process during - // openSession) from racing releases until this session's entry - // is recorded below. Dropped exactly once on every path. - let reservationDropped = mcpCredentialId === undefined; - const dropReservation = Effect.sync(() => { - if (!reservationDropped && mcpCredentialId !== undefined) { - reservationDropped = true; - dropMcpCredentialReservation(input.threadId, mcpCredentialId); - } - }); - const sessionScope = yield* Scope.make(); - const runtime = yield* adapter - .openSession({ - threadId: input.threadId, - providerSessionId: input.providerSessionId, - modelSelection: input.modelSelection, - runtimePolicy: input.runtimePolicy, - ...(input.resumeFromSession === undefined - ? {} - : { resumeFromSession: input.resumeFromSession }), - ...(input.initialNativeThreadId === undefined - ? {} - : { initialNativeThreadId: input.initialNativeThreadId }), - ...(input.initialProviderItemIdentityVersion === undefined - ? {} - : { - initialProviderItemIdentityVersion: - input.initialProviderItemIdentityVersion, - }), - }) - .pipe( - Effect.provideService(Scope.Scope, sessionScope), - Effect.tapError(() => - Scope.close(sessionScope, Exit.void).pipe( - Effect.ignore, - Effect.andThen(dropReservation), - // Revoke only a credential this open freshly minted: a - // reused credential is held by another live provider - // process and must survive this open's failure. - Effect.andThen( - prepared.issued - ? clearMcpSession(input.threadId, mcpCredentialId) - : Effect.void, - ), - ), - ), - Effect.onInterrupt(() => dropReservation), + const adapter = yield* registry.get(input.modelSelection.instanceId).pipe( Effect.mapError( (cause) => new ProviderSessionOpenError({ @@ -1657,62 +1613,137 @@ export const layerWithOptions = ( }), ), ); - const eventSubscribers = yield* Ref.make< - ReadonlyMap> - >(new Map()); - const exposedRuntime = decorateRuntime(runtime, eventSubscribers); - const now = yield* Clock.currentTimeMillis; - const entry: LiveSessionEntry = { - attachedThreadIds: new Set([input.threadId]), - loadedProviderThreadKeyByThread: new Map(), - mcpCredentialIdByThread: - mcpCredentialId === undefined - ? new Map() - : new Map([[input.threadId, mcpCredentialId]]), - supportsMultipleProviderThreads: - runtime.providerSession.capabilities.sessions - .supportsMultipleProviderThreadsPerSession, - runtime, - exposedRuntime, - eventSubscribers, - requestEventPermit: yield* Semaphore.make(1), - scope: sessionScope, - idleGeneration: 0, - busyCount: 0, - lastActivityAtMs: now, - idleFiber: null, - pinnedSinceMs: null, - }; - yield* Ref.update(sessions, (current) => { - const updated = new Map(current); - updated.set(key, entry); - return updated; - }); - // The entry now guards the credential via its recorded id, so - // the pre-open reservation can be dropped. - yield* dropReservation; - yield* withActivityError( - input.providerSessionId, - writeProviderSessionEvents({ - runtime, - threadIds: [input.threadId], - type: "provider-session.attached", - payload: runtime.providerSession, - }), - ).pipe( - Effect.tapError(() => - releaseEntry({ + const prepared = yield* prepareMcpSession( + input.threadId, + input.modelSelection.instanceId, + ); + const mcpCredentialId = prepared.mcpCredentialId; + // The reservation from prepare protects the credential (which + // eager adapters bake into the provider process during + // openSession) from racing releases until this session's entry + // is recorded below. Dropped exactly once on every path. + let reservationDropped = mcpCredentialId === undefined; + const dropReservation = Effect.sync(() => { + if (!reservationDropped && mcpCredentialId !== undefined) { + reservationDropped = true; + dropMcpCredentialReservation(input.threadId, mcpCredentialId); + } + }); + const sessionScope = yield* Scope.make(); + const runtime = yield* adapter + .openSession({ + threadId: input.threadId, providerSessionId: input.providerSessionId, - reason: "runtime_error", - detail: "Failed to persist the provider-session attachment.", - }).pipe(Effect.ignore), - ), - ); - yield* startEventPump(entry); - yield* scheduleIdleRelease(input.providerSessionId); - return exposedRuntime; - }), - ), + modelSelection: input.modelSelection, + runtimePolicy: input.runtimePolicy, + ...(input.resumeFromSession === undefined + ? {} + : { resumeFromSession: input.resumeFromSession }), + ...(input.initialNativeThreadId === undefined + ? {} + : { initialNativeThreadId: input.initialNativeThreadId }), + ...(input.initialProviderItemIdentityVersion === undefined + ? {} + : { + initialProviderItemIdentityVersion: + input.initialProviderItemIdentityVersion, + }), + }) + .pipe( + Effect.provideService(Scope.Scope, sessionScope), + Effect.tapError(() => + Scope.close(sessionScope, Exit.void).pipe( + Effect.ignore, + Effect.andThen(dropReservation), + // Revoke only a credential this open freshly minted: a + // reused credential is held by another live provider + // process and must survive this open's failure. + Effect.andThen( + prepared.issued + ? clearMcpSession(input.threadId, mcpCredentialId) + : Effect.void, + ), + ), + ), + Effect.onInterrupt(() => dropReservation), + Effect.mapError( + (cause) => + new ProviderSessionOpenError({ + instanceId: input.modelSelection.instanceId, + providerSessionId: input.providerSessionId, + cause, + }), + ), + ); + openedDriver = runtime.driver; + const eventSubscribers = yield* Ref.make< + ReadonlyMap> + >(new Map()); + const exposedRuntime = decorateRuntime(runtime, eventSubscribers); + const now = yield* Clock.currentTimeMillis; + const entry: LiveSessionEntry = { + attachedThreadIds: new Set([input.threadId]), + loadedProviderThreadKeyByThread: new Map(), + mcpCredentialIdByThread: + mcpCredentialId === undefined + ? new Map() + : new Map([[input.threadId, mcpCredentialId]]), + supportsMultipleProviderThreads: + runtime.providerSession.capabilities.sessions + .supportsMultipleProviderThreadsPerSession, + runtime, + exposedRuntime, + eventSubscribers, + requestEventPermit: yield* Semaphore.make(1), + scope: sessionScope, + idleGeneration: 0, + busyCount: 0, + lastActivityAtMs: now, + idleFiber: null, + pinnedSinceMs: null, + }; + yield* Ref.update(sessions, (current) => { + const updated = new Map(current); + updated.set(key, entry); + return updated; + }); + // The entry now guards the credential via its recorded id, so + // the pre-open reservation can be dropped. + yield* dropReservation; + yield* withActivityError( + input.providerSessionId, + writeProviderSessionEvents({ + runtime, + threadIds: [input.threadId], + type: "provider-session.attached", + payload: runtime.providerSession, + }), + ).pipe( + Effect.tapError(() => + releaseEntry({ + providerSessionId: input.providerSessionId, + reason: "runtime_error", + detail: "Failed to persist the provider-session attachment.", + }).pipe(Effect.ignore), + ), + ); + yield* startEventPump(entry); + yield* scheduleIdleRelease(input.providerSessionId); + return exposedRuntime; + }), + ) + .pipe( + Effect.onExit((exit) => + reusedExisting + ? Effect.void + : increment(providerSessionsTotal, { + provider: openedDriver, + operation: input.resumeFromSession === undefined ? "start" : "recover", + outcome: outcomeFromExit(exit), + }), + ), + ); + }, get: (providerSessionId) => Effect.gen(function* () { const entry = (yield* Ref.get(sessions)).get(sessionKey(providerSessionId)); 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..0c60f2497c50 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 }); @@ -295,6 +300,15 @@ export const make = Effect.gen(function* () { } }); + const commit = (command: ProjectCommand) => + commitEffect(command).pipe( + withMetrics({ + counter: orchestrationCommandsTotal, + timer: orchestrationCommandDuration, + attributes: { commandType: command.type }, + }), + ); + const readCommitted = Effect.fn("ProjectService.readCommitted")(function* (projectId: ProjectId) { const row = yield* readRow(projectId, { includeDeleted: true }); if (Option.isNone(row)) { From 555cf5e0514840de07cb27bc14e3668eea89e056 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 06:25:38 -0400 Subject: [PATCH 2/5] fix(observability): drop unused provider metric attributes helper Signed-off-by: Yordis Prieto --- apps/server/src/observability/Metrics.ts | 9 --------- 1 file changed, 9 deletions(-) diff --git a/apps/server/src/observability/Metrics.ts b/apps/server/src/observability/Metrics.ts index c1c0aba06b94..55cdc1519582 100644 --- a/apps/server/src/observability/Metrics.ts +++ b/apps/server/src/observability/Metrics.ts @@ -157,15 +157,6 @@ export const withMetrics: { (effect: Effect.Effect, options: WithMetricsOptions): Effect.Effect; } = dual(2, withMetricsImpl); -export const providerMetricAttributes = ( - provider: string, - extra?: Readonly>, -) => - compactMetricAttributes({ - provider, - ...extra, - }); - export const providerTurnMetricAttributes = (input: { readonly provider: string; readonly model: string | null | undefined; From 7cb5404defeb2d0b1a468f9ce4793b6a86731ddb Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 14:19:48 -0400 Subject: [PATCH 3/5] fix(observability): session and command metrics stay comparable across paths Signed-off-by: Yordis Prieto --- apps/server/src/observability/Metrics.ts | 3 ++- .../src/orchestration-v2/ProviderSessionManager.ts | 4 ++-- apps/server/src/project/ProjectService.ts | 11 ++++++++--- 3 files changed, 12 insertions(+), 6 deletions(-) diff --git a/apps/server/src/observability/Metrics.ts b/apps/server/src/observability/Metrics.ts index 55cdc1519582..edc03e87460a 100644 --- a/apps/server/src/observability/Metrics.ts +++ b/apps/server/src/observability/Metrics.ts @@ -31,7 +31,8 @@ export const orchestrationCommandsTotal = Metric.counter("t3_orchestration_comma }); export const orchestrationCommandDuration = Metric.timer("t3_orchestration_command_duration", { - description: "Orchestration command dispatch duration, from receipt lookup through commit.", + description: + "Orchestration command dispatch duration while holding its dispatch lock, excluding lock wait.", }); export const orchestrationCommandAckDuration = Metric.timer( diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 5cdbe4622c6b..8e9c88fbfc10 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -1583,6 +1583,7 @@ export const layerWithOptions = ( const key = sessionKey(input.providerSessionId); const existing = (yield* Ref.get(sessions)).get(key); if (existing !== undefined) { + reusedExisting = true; if ( !existing.attachedThreadIds.has(input.threadId) && !existing.supportsMultipleProviderThreads @@ -1599,7 +1600,6 @@ export const layerWithOptions = ( providerInstanceId: existing.runtime.instanceId, }); yield* touchActivity(input.providerSessionId); - reusedExisting = true; return existing.exposedRuntime; } @@ -1613,6 +1613,7 @@ export const layerWithOptions = ( }), ), ); + openedDriver = adapter.driver; const prepared = yield* prepareMcpSession( input.threadId, input.modelSelection.instanceId, @@ -1675,7 +1676,6 @@ export const layerWithOptions = ( }), ), ); - openedDriver = runtime.driver; const eventSubscribers = yield* Ref.make< ReadonlyMap> >(new Map()); diff --git a/apps/server/src/project/ProjectService.ts b/apps/server/src/project/ProjectService.ts index 0c60f2497c50..f31de0ce1d05 100644 --- a/apps/server/src/project/ProjectService.ts +++ b/apps/server/src/project/ProjectService.ts @@ -267,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) { @@ -304,7 +310,6 @@ export const make = Effect.gen(function* () { commitEffect(command).pipe( withMetrics({ counter: orchestrationCommandsTotal, - timer: orchestrationCommandDuration, attributes: { commandType: command.type }, }), ); From aa3fa628e4defa75047e66f643820c59c3ead35a Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 14:34:06 -0400 Subject: [PATCH 4/5] fix(observability): session metrics stay accurate when an open is re-run Signed-off-by: Yordis Prieto --- .../ProviderSessionManager.ts | 370 +++++++++--------- 1 file changed, 186 insertions(+), 184 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 8e9c88fbfc10..d3a3b83f1d8a 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -717,9 +717,8 @@ export const layerWithOptions = ( readonly cancelIdleFiber?: boolean; readonly onlyIfIdleGeneration?: number; readonly gracefulSubscribers?: boolean; - }) => { - let releasedDriver: string | undefined; - return Effect.acquireUseRelease( + }) => + Effect.acquireUseRelease( Ref.modify(sessions, (current) => { const key = sessionKey(input.providerSessionId); const existing = current.get(key); @@ -734,7 +733,6 @@ export const layerWithOptions = ( } const updated = new Map(current); updated.delete(key); - releasedDriver = existing.runtime.driver; return [Option.some(existing), updated] as const; }), (entry) => @@ -810,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, { @@ -861,17 +867,7 @@ export const layerWithOptions = ( }), ), ), - Effect.onExit((exit) => - releasedDriver === undefined - ? Effect.void - : increment(providerSessionsTotal, { - provider: releasedDriver, - operation: "stop", - outcome: outcomeFromExit(exit), - }), - ), ); - }; // Annotated to break the releaseIfStillIdle <-> scheduleIdleReleaseInternal // inference cycle introduced by the pin re-arm below. @@ -1558,52 +1554,133 @@ 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) => { - let openedDriver: string | undefined; - let reusedExisting = false; - return sessionOpen - .withLock( - input.providerSessionId, - Effect.gen(function* () { - const cwd = input.runtimePolicy.cwd; - if (cwd !== null) { - const workspaceIsDirectory = yield* fileSystem.stat(cwd).pipe( - Effect.map((stat) => stat.type === "Directory"), - Effect.catch((error) => Effect.succeed(error.reason._tag !== "NotFound")), - ); - if (!workspaceIsDirectory) { - return yield* new ProviderWorkspaceMissingError({ - threadId: input.threadId, - cwd, - }); - } + open: (input) => + sessionOpen.withLock( + input.providerSessionId, + Effect.gen(function* () { + const cwd = input.runtimePolicy.cwd; + if (cwd !== null) { + const workspaceIsDirectory = yield* fileSystem.stat(cwd).pipe( + Effect.map((stat) => stat.type === "Directory"), + Effect.catch((error) => Effect.succeed(error.reason._tag !== "NotFound")), + ); + if (!workspaceIsDirectory) { + return yield* new ProviderWorkspaceMissingError({ + threadId: input.threadId, + cwd, + }); } - const key = sessionKey(input.providerSessionId); - const existing = (yield* Ref.get(sessions)).get(key); - if (existing !== undefined) { - reusedExisting = true; - if ( - !existing.attachedThreadIds.has(input.threadId) && - !existing.supportsMultipleProviderThreads - ) { - return yield* new ProviderSessionOpenError({ - instanceId: input.modelSelection.instanceId, - providerSessionId: input.providerSessionId, - cause: `Provider ${existing.runtime.driver} does not support attaching multiple app threads to one session.`, - }); - } - yield* ensureThreadAttached({ + } + const key = sessionKey(input.providerSessionId); + const existing = (yield* Ref.get(sessions)).get(key); + if (existing !== undefined) { + if ( + !existing.attachedThreadIds.has(input.threadId) && + !existing.supportsMultipleProviderThreads + ) { + return yield* new ProviderSessionOpenError({ + instanceId: input.modelSelection.instanceId, providerSessionId: input.providerSessionId, - threadId: input.threadId, - providerInstanceId: existing.runtime.instanceId, + cause: `Provider ${existing.runtime.driver} does not support attaching multiple app threads to one session.`, }); - yield* touchActivity(input.providerSessionId); - return existing.exposedRuntime; } + yield* ensureThreadAttached({ + providerSessionId: input.providerSessionId, + threadId: input.threadId, + providerInstanceId: existing.runtime.instanceId, + }); + yield* touchActivity(input.providerSessionId); + return existing.exposedRuntime; + } - const adapter = yield* registry.get(input.modelSelection.instanceId).pipe( + const adapter = yield* registry.get(input.modelSelection.instanceId).pipe( + Effect.mapError( + (cause) => + new ProviderSessionOpenError({ + instanceId: input.modelSelection.instanceId, + providerSessionId: input.providerSessionId, + cause, + }), + ), + ); + const prepared = yield* prepareMcpSession( + input.threadId, + input.modelSelection.instanceId, + ); + const mcpCredentialId = prepared.mcpCredentialId; + // The reservation from prepare protects the credential (which + // eager adapters bake into the provider process during + // openSession) from racing releases until this session's entry + // is recorded below. Dropped exactly once on every path. + let reservationDropped = mcpCredentialId === undefined; + const dropReservation = Effect.sync(() => { + if (!reservationDropped && mcpCredentialId !== undefined) { + reservationDropped = true; + dropMcpCredentialReservation(input.threadId, mcpCredentialId); + } + }); + const sessionScope = yield* Scope.make(); + const runtime = yield* adapter + .openSession({ + threadId: input.threadId, + providerSessionId: input.providerSessionId, + modelSelection: input.modelSelection, + runtimePolicy: input.runtimePolicy, + ...(input.resumeFromSession === undefined + ? {} + : { resumeFromSession: input.resumeFromSession }), + ...(input.initialNativeThreadId === undefined + ? {} + : { initialNativeThreadId: input.initialNativeThreadId }), + ...(input.initialProviderItemIdentityVersion === undefined + ? {} + : { + initialProviderItemIdentityVersion: + input.initialProviderItemIdentityVersion, + }), + }) + .pipe( + Effect.provideService(Scope.Scope, sessionScope), + Effect.tapError(() => + Scope.close(sessionScope, Exit.void).pipe( + Effect.ignore, + Effect.andThen(dropReservation), + // Revoke only a credential this open freshly minted: a + // reused credential is held by another live provider + // process and must survive this open's failure. + Effect.andThen( + prepared.issued + ? clearMcpSession(input.threadId, mcpCredentialId) + : Effect.void, + ), + ), + ), + Effect.onInterrupt(() => dropReservation), Effect.mapError( (cause) => new ProviderSessionOpenError({ @@ -1613,137 +1690,62 @@ export const layerWithOptions = ( }), ), ); - openedDriver = adapter.driver; - const prepared = yield* prepareMcpSession( - input.threadId, - input.modelSelection.instanceId, - ); - const mcpCredentialId = prepared.mcpCredentialId; - // The reservation from prepare protects the credential (which - // eager adapters bake into the provider process during - // openSession) from racing releases until this session's entry - // is recorded below. Dropped exactly once on every path. - let reservationDropped = mcpCredentialId === undefined; - const dropReservation = Effect.sync(() => { - if (!reservationDropped && mcpCredentialId !== undefined) { - reservationDropped = true; - dropMcpCredentialReservation(input.threadId, mcpCredentialId); - } - }); - const sessionScope = yield* Scope.make(); - const runtime = yield* adapter - .openSession({ - threadId: input.threadId, - providerSessionId: input.providerSessionId, - modelSelection: input.modelSelection, - runtimePolicy: input.runtimePolicy, - ...(input.resumeFromSession === undefined - ? {} - : { resumeFromSession: input.resumeFromSession }), - ...(input.initialNativeThreadId === undefined - ? {} - : { initialNativeThreadId: input.initialNativeThreadId }), - ...(input.initialProviderItemIdentityVersion === undefined - ? {} - : { - initialProviderItemIdentityVersion: - input.initialProviderItemIdentityVersion, - }), - }) - .pipe( - Effect.provideService(Scope.Scope, sessionScope), - Effect.tapError(() => - Scope.close(sessionScope, Exit.void).pipe( - Effect.ignore, - Effect.andThen(dropReservation), - // Revoke only a credential this open freshly minted: a - // reused credential is held by another live provider - // process and must survive this open's failure. - Effect.andThen( - prepared.issued - ? clearMcpSession(input.threadId, mcpCredentialId) - : Effect.void, - ), - ), - ), - Effect.onInterrupt(() => dropReservation), - Effect.mapError( - (cause) => - new ProviderSessionOpenError({ - instanceId: input.modelSelection.instanceId, - providerSessionId: input.providerSessionId, - cause, - }), - ), - ); - const eventSubscribers = yield* Ref.make< - ReadonlyMap> - >(new Map()); - const exposedRuntime = decorateRuntime(runtime, eventSubscribers); - const now = yield* Clock.currentTimeMillis; - const entry: LiveSessionEntry = { - attachedThreadIds: new Set([input.threadId]), - loadedProviderThreadKeyByThread: new Map(), - mcpCredentialIdByThread: - mcpCredentialId === undefined - ? new Map() - : new Map([[input.threadId, mcpCredentialId]]), - supportsMultipleProviderThreads: - runtime.providerSession.capabilities.sessions - .supportsMultipleProviderThreadsPerSession, + const eventSubscribers = yield* Ref.make< + ReadonlyMap> + >(new Map()); + const exposedRuntime = decorateRuntime(runtime, eventSubscribers); + const now = yield* Clock.currentTimeMillis; + const entry: LiveSessionEntry = { + attachedThreadIds: new Set([input.threadId]), + loadedProviderThreadKeyByThread: new Map(), + mcpCredentialIdByThread: + mcpCredentialId === undefined + ? new Map() + : new Map([[input.threadId, mcpCredentialId]]), + supportsMultipleProviderThreads: + runtime.providerSession.capabilities.sessions + .supportsMultipleProviderThreadsPerSession, + runtime, + exposedRuntime, + eventSubscribers, + requestEventPermit: yield* Semaphore.make(1), + scope: sessionScope, + idleGeneration: 0, + busyCount: 0, + lastActivityAtMs: now, + idleFiber: null, + pinnedSinceMs: null, + }; + yield* Ref.update(sessions, (current) => { + const updated = new Map(current); + updated.set(key, entry); + return updated; + }); + // The entry now guards the credential via its recorded id, so + // the pre-open reservation can be dropped. + yield* dropReservation; + yield* withActivityError( + input.providerSessionId, + writeProviderSessionEvents({ runtime, - exposedRuntime, - eventSubscribers, - requestEventPermit: yield* Semaphore.make(1), - scope: sessionScope, - idleGeneration: 0, - busyCount: 0, - lastActivityAtMs: now, - idleFiber: null, - pinnedSinceMs: null, - }; - yield* Ref.update(sessions, (current) => { - const updated = new Map(current); - updated.set(key, entry); - return updated; - }); - // The entry now guards the credential via its recorded id, so - // the pre-open reservation can be dropped. - yield* dropReservation; - yield* withActivityError( - input.providerSessionId, - writeProviderSessionEvents({ - runtime, - threadIds: [input.threadId], - type: "provider-session.attached", - payload: runtime.providerSession, - }), - ).pipe( - Effect.tapError(() => - releaseEntry({ - providerSessionId: input.providerSessionId, - reason: "runtime_error", - detail: "Failed to persist the provider-session attachment.", - }).pipe(Effect.ignore), - ), - ); - yield* startEventPump(entry); - yield* scheduleIdleRelease(input.providerSessionId); - return exposedRuntime; - }), - ) - .pipe( - Effect.onExit((exit) => - reusedExisting - ? Effect.void - : increment(providerSessionsTotal, { - provider: openedDriver, - operation: input.resumeFromSession === undefined ? "start" : "recover", - outcome: outcomeFromExit(exit), - }), - ), - ); - }, + threadIds: [input.threadId], + type: "provider-session.attached", + payload: runtime.providerSession, + }), + ).pipe( + Effect.tapError(() => + releaseEntry({ + providerSessionId: input.providerSessionId, + reason: "runtime_error", + detail: "Failed to persist the provider-session attachment.", + }).pipe(Effect.ignore), + ), + ); + yield* startEventPump(entry); + yield* scheduleIdleRelease(input.providerSessionId); + return exposedRuntime; + }).pipe(recordSessionOpen(input)), + ), get: (providerSessionId) => Effect.gen(function* () { const entry = (yield* Ref.get(sessions)).get(sessionKey(providerSessionId)); From 8f3a62b3494dc9ca00a553bab257b33d7df11201 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 14:49:34 -0400 Subject: [PATCH 5/5] test(observability): session metrics survive a re-run open Signed-off-by: Yordis Prieto --- .../orchestration-v2/ProviderSessionManager.test.ts | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 261557f0672b..cbc0e4360417 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -869,13 +869,16 @@ it.effect("ProviderSessionManagerV2 records provider session lifecycle metrics", const startsBefore = yield* sessionCount("start"); const stopsBefore = yield* sessionCount("stop"); - yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); - yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + 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, 1); - assert.equal((yield* sessionCount("stop")) - stopsBefore, 1); + assert.equal((yield* sessionCount("start")) - startsBefore, 2); + assert.equal((yield* sessionCount("stop")) - stopsBefore, 2); }); yield* effect.pipe(Effect.provide(makeTestLayer({ state, idleTimeoutMs: 60_000 })));