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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 25 additions & 12 deletions apps/server/src/observability/Metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
{
Expand All @@ -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.",
});

Expand Down Expand Up @@ -139,13 +158,7 @@ export const withMetrics: {
<A, E, R>(effect: Effect.Effect<A, E, R>, options: WithMetricsOptions): Effect.Effect<A, E, R>;
} = dual(2, withMetricsImpl);

const providerMetricAttributes = (provider: string, extra?: Readonly<Record<string, unknown>>) =>
compactMetricAttributes({
provider,
...extra,
});

const providerTurnMetricAttributes = (input: {
export const providerTurnMetricAttributes = (input: {
readonly provider: string;
readonly model: string | null | undefined;
readonly extra?: Readonly<Record<string, unknown>>;
Expand Down
54 changes: 54 additions & 0 deletions apps/server/src/orchestration-v2/EffectWorker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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;
Expand Down
5 changes: 5 additions & 0 deletions apps/server/src/orchestration-v2/EffectWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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)),
);
32 changes: 31 additions & 1 deletion apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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<OrchestratorV2DispatchResult, OrchestratorV2Error> {
yield* Effect.annotateCurrentSpan({
"orchestration_v2.command_id": command.commandId,
Expand Down Expand Up @@ -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));
}
Expand All @@ -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 },
}),
),
Comment thread
cursor[bot] marked this conversation as resolved.
);
});

const handleTerminalRun = (stored: OrchestrationV2StoredEvent) =>
Effect.gen(function* () {
Expand Down
Loading
Loading