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
22 changes: 22 additions & 0 deletions apps/server/src/keybindings.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -604,6 +604,28 @@ it.layer(NodeServices.layer)("keybindings", (it) => {
}).pipe(Effect.provide(makeKeybindingsLayer())),
);

it.effect("retries a failed config read instead of keeping the failure", () =>
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const { keybindingsConfigPath } = yield* ServerConfig.ServerConfig;
// A directory where the file should be makes the read itself fail.
yield* fileSystem.makeDirectory(keybindingsConfigPath, { recursive: true });

const keybindings = yield* Keybindings.Keybindings;
const failed = yield* toDetailResult(keybindings.loadConfigState);
assertFailure(failed, "failed to read keybindings config");

yield* fileSystem.remove(keybindingsConfigPath, { recursive: true });
yield* writeKeybindingsConfig(keybindingsConfigPath, [
{ key: "mod+j", command: "terminal.toggle" },
]);

const configState = yield* keybindings.loadConfigState;
assert.deepEqual(configState.issues, []);
assert.isTrue(configState.keybindings.some((entry) => entry.command === "terminal.toggle"));
}).pipe(Effect.provide(makeKeybindingsLayer())),
);

it.effect("updates cached resolved config after upsert", () =>
Effect.gen(function* () {
const { keybindingsConfigPath } = yield* ServerConfig.ServerConfig;
Expand Down
7 changes: 4 additions & 3 deletions apps/server/src/keybindings.ts
Original file line number Diff line number Diff line change
Expand Up @@ -480,13 +480,14 @@ const make = Effect.gen(function* () {
})),
);

const resolvedConfigCache = yield* Cache.make<
// A failed read is not kept: the next read retries instead of replaying the failure.
const resolvedConfigCache = yield* Cache.makeWith<
typeof resolvedConfigCacheKey,
KeybindingsConfigState,
KeybindingsConfigError
>({
>(() => loadConfigStateFromDisk, {
capacity: 1,
lookup: () => loadConfigStateFromDisk,
timeToLive: (exit) => (Exit.isSuccess(exit) ? Duration.infinity : Duration.zero),
});

const loadConfigStateFromCacheOrDisk = Cache.get(resolvedConfigCache, resolvedConfigCacheKey);
Expand Down
59 changes: 59 additions & 0 deletions apps/server/src/observability/Metrics.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { assert, describe, it } from "@effect/vitest";
import { ProviderDriverKind } from "@t3tools/contracts";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
Expand Down Expand Up @@ -127,6 +128,64 @@ describe("withMetrics", () => {
}),
);

it.effect("counts interrupted work with an interrupt outcome and its duration", () =>
Effect.gen(function* () {
const counter = Metric.counter("with_metrics_interrupt_total");
const timer = Metric.timer("with_metrics_interrupt_duration");
const started = yield* Deferred.make<void>();

const fiber = yield* Deferred.succeed(started, undefined).pipe(
Effect.andThen(Effect.never),
withMetrics({ counter, timer, attributes: { operation: "interrupt" } }),
Effect.forkChild,
);
yield* Deferred.await(started);
yield* TestClock.adjust(Duration.millis(5));
yield* Fiber.interrupt(fiber);

const snapshots = yield* Metric.snapshot;
assert.equal(
hasMetricSnapshot(snapshots, "with_metrics_interrupt_total", {
operation: "interrupt",
outcome: "interrupt",
}),
true,
);
const duration = findHistogramSnapshot(snapshots, "with_metrics_interrupt_duration", {
operation: "interrupt",
});
assert.equal(duration?.state.count, 1);
assert.equal(duration?.state.sum, 5);
}),
);

it.effect("measures durations on the monotonic clock, not the wall clock", () =>
Effect.gen(function* () {
const timer = Metric.timer("with_metrics_monotonic_duration");
const started = yield* Deferred.make<void>();
const finish = yield* Deferred.make<void>();

const fiber = yield* Deferred.succeed(started, undefined).pipe(
Effect.andThen(Deferred.await(finish)),
withMetrics({ timer, attributes: { operation: "monotonic" } }),
Effect.forkChild,
);
yield* Deferred.await(started);
yield* TestClock.adjust(Duration.millis(10));
// A backward wall-clock correction must not shorten the measured duration.
yield* TestClock.setTime(0);
yield* Deferred.succeed(finish, undefined);
yield* Fiber.join(fiber);

const snapshots = yield* Metric.snapshot;
const duration = findHistogramSnapshot(snapshots, "with_metrics_monotonic_duration", {
operation: "monotonic",
});
assert.equal(duration?.state.count, 1);
assert.equal(duration?.state.sum, 10);
}),
);

it.effect("records timer durations from nanosecond clock readings", () =>
Effect.gen(function* () {
const duration = Duration.nanos(1_500_000n);
Expand Down
76 changes: 22 additions & 54 deletions apps/server/src/observability/Metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,7 @@ import * as Exit from "effect/Exit";
import * as Metric from "effect/Metric";
import { dual } from "effect/Function";

import {
compactMetricAttributes,
normalizeModelMetricLabel,
outcomeFromExit,
} from "./Attributes.ts";
import { compactMetricAttributes, outcomeFromExit } from "./Attributes.ts";

export const rpcRequestsTotal = Metric.counter("t3_rpc_requests_total", {
description: "Total RPC requests handled by the websocket RPC server.",
Expand All @@ -19,13 +15,6 @@ export const rpcRequestDuration = Metric.timer("t3_rpc_request_duration", {
description: "RPC request handling duration.",
});

const orchestrationEventsProcessedTotal = Metric.counter(
"t3_orchestration_events_processed_total",
{
description: "Total orchestration intent events processed by runtime reactors.",
},
);

export const orchestrationEffectClaimsTotal = Metric.counter(
"t3_orchestration_effect_claims_total",
{
Expand All @@ -38,20 +27,16 @@ 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", {
description: "Provider turn request duration.",
});

const providerRuntimeEventsTotal = Metric.counter("t3_provider_runtime_events_total", {
description: "Total canonical provider runtime events processed.",
export const providerTurnDuration = Metric.timer("t3_provider_turn_duration", {
description: "Time for the provider adapter to start a turn, not how long the turn runs.",
});

export const gitCommandsTotal = Metric.counter("t3_git_commands_total", {
Expand Down Expand Up @@ -91,16 +76,13 @@ export interface WithMetricsOptions {
) => Readonly<Record<string, unknown>>;
}

const withMetricsImpl = <A, E, R>(
effect: Effect.Effect<A, E, R>,
const recordMetrics = (
options: WithMetricsOptions,
): Effect.Effect<A, E, R> =>
startedAt: bigint,
exit: Exit.Exit<unknown, unknown>,
) =>
Effect.gen(function* () {
const startedAt = yield* Clock.currentTimeNanos;
const exit = yield* Effect.exit(effect);
const endedAt = yield* Clock.currentTimeNanos;
const elapsedNanos = endedAt > startedAt ? endedAt - startedAt : 0n;
const duration = Duration.nanos(elapsedNanos);
const duration = Duration.nanos((yield* Clock.monotonicTimeNanos) - startedAt);
const baseAttributes =
typeof options.attributes === "function" ? options.attributes() : (options.attributes ?? {});

Expand All @@ -125,35 +107,21 @@ const withMetricsImpl = <A, E, R>(
1,
);
}

if (Exit.isSuccess(exit)) {
return exit.value;
}
return yield* Effect.failCause(exit.cause);
});

// Durations come from the monotonic clock, so wall-clock corrections cannot skew them, and
// metrics are recorded in an exit finalizer, so interrupted work is counted as "interrupt".
const withMetricsImpl = <A, E, R>(
effect: Effect.Effect<A, E, R>,
options: WithMetricsOptions,
): Effect.Effect<A, E, R> =>
Effect.flatMap(Clock.monotonicTimeNanos, (startedAt) =>
Effect.onExit(effect, (exit) => recordMetrics(options, startedAt, exit)),
);

export const withMetrics: {
<A, E, R>(
(
options: WithMetricsOptions,
): (effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E, R>;
): <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E, R>;
<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: {
readonly provider: string;
readonly model: string | null | undefined;
readonly extra?: Readonly<Record<string, unknown>>;
}) => {
const modelFamily = normalizeModelMetricLabel(input.model);
return compactMetricAttributes({
provider: input.provider,
...(modelFamily ? { modelFamily } : {}),
...input.extra,
});
};
22 changes: 8 additions & 14 deletions apps/server/src/observability/RpcInstrumentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,12 +67,9 @@ const recordRpcStreamMetrics = <E>(
exit: Exit.Exit<unknown, E>,
): Effect.Effect<void, never, never> =>
Effect.gen(function* () {
const endedAt = yield* Clock.currentTimeNanos;
const elapsedNanos = endedAt > startedAt ? endedAt - startedAt : 0n;

yield* Metric.update(
Metric.withAttributes(rpcRequestDuration, metricAttributes({ method })),
Duration.nanos(elapsedNanos),
Duration.nanos((yield* Clock.monotonicTimeNanos) - startedAt),
);
yield* Metric.update(
Metric.withAttributes(
Expand Down Expand Up @@ -111,7 +108,7 @@ export const observeRpcStream = <A, E, R>(
): Stream.Stream<A, E, R> => {
const instrumented = Stream.unwrap(
Effect.gen(function* () {
const startedAt = yield* Clock.currentTimeNanos;
const startedAt = yield* Clock.monotonicTimeNanos;
return stream.pipe(Stream.onExit((exit) => recordRpcStreamMetrics(method, startedAt, exit)));
}),
);
Expand All @@ -126,15 +123,12 @@ export const observeRpcStreamEffect = <A, StreamError, StreamContext, EffectErro
): Stream.Stream<A, StreamError | EffectError, StreamContext | EffectContext> => {
const instrumented = Stream.unwrap(
Effect.gen(function* () {
const startedAt = yield* Clock.currentTimeNanos;
const exit = yield* Effect.exit(effect);

if (Exit.isFailure(exit)) {
yield* recordRpcStreamMetrics(method, startedAt, exit);
return yield* Effect.failCause(exit.cause);
}

return exit.value.pipe(
const startedAt = yield* Clock.monotonicTimeNanos;
// onError also runs when the stream is interrupted before it is produced.
const stream = yield* effect.pipe(
Effect.onError((cause) => recordRpcStreamMetrics(method, startedAt, Exit.failCause(cause))),
);
return stream.pipe(
Stream.onExit((streamExit) => recordRpcStreamMetrics(method, startedAt, streamExit)),
);
}),
Expand Down
Loading
Loading