diff --git a/src/services/AsyncIngestBun.ts b/src/services/AsyncIngestBun.ts index 5e4a5dd..a842020 100644 --- a/src/services/AsyncIngestBun.ts +++ b/src/services/AsyncIngestBun.ts @@ -1,9 +1,9 @@ import * as BunWorker from "@effect/platform-bun/BunWorker" -import { Context, Duration, Effect, Exit, Layer, Scope } from "effect" +import { Duration, Effect, Exit, Layer, Scope } from "effect" import * as RpcClient from "effect/unstable/rpc/RpcClient" -import type { RpcClientError } from "effect/unstable/rpc/RpcClientError" +import { RpcClientError } from "effect/unstable/rpc/RpcClientError" import * as RpcSerialization from "effect/unstable/rpc/RpcSerialization" -import type { WorkerError } from "effect/unstable/workers/WorkerError" +import { WorkerError } from "effect/unstable/workers/WorkerError" import { IngestRpcs } from "./ingestRpc.ts" import { AsyncIngest } from "./AsyncIngest.js" @@ -39,14 +39,32 @@ export const AsyncIngestLive = Layer.effect( // HTTP health wait for SQLite bootstrap. Managed readiness still verifies // the worker through explicit ingest probes. yield* Effect.forkScoped(getClient.pipe(Effect.ignore)) + // Cancellation and payload errors belong to one request. Only a broken + // transport invalidates the worker shared by all ingest callers. return { ingestTraces: (input) => Effect.flatMap(getClient, ({ client, clientScope }) => - client.ingestTraces(input).pipe(Effect.onError(() => Effect.andThen(Scope.close(clientScope, Exit.void), invalidateClient))), + client + .ingestTraces(input) + .pipe( + Effect.tapError((error) => + error instanceof RpcClientError || error instanceof WorkerError + ? Effect.andThen(Scope.close(clientScope, Exit.void), invalidateClient) + : Effect.void, + ), + ), ), ingestLogs: (input) => Effect.flatMap(getClient, ({ client, clientScope }) => - client.ingestLogs(input).pipe(Effect.onError(() => Effect.andThen(Scope.close(clientScope, Exit.void), invalidateClient))), + client + .ingestLogs(input) + .pipe( + Effect.tapError((error) => + error instanceof RpcClientError || error instanceof WorkerError + ? Effect.andThen(Scope.close(clientScope, Exit.void), invalidateClient) + : Effect.void, + ), + ), ), } }), diff --git a/src/services/TelemetryQuery.ts b/src/services/TelemetryQuery.ts index 41d8ba0..b556173 100644 --- a/src/services/TelemetryQuery.ts +++ b/src/services/TelemetryQuery.ts @@ -2,8 +2,8 @@ import * as BunWorker from "@effect/platform-bun/BunWorker" import { Duration, Effect, Exit, Layer, Scope } from "effect" import * as RpcClient from "effect/unstable/rpc/RpcClient" import * as RpcSerialization from "effect/unstable/rpc/RpcSerialization" -import type { WorkerError } from "effect/unstable/workers/WorkerError" -import type { RpcClientError } from "effect/unstable/rpc/RpcClientError" +import { WorkerError } from "effect/unstable/workers/WorkerError" +import { RpcClientError } from "effect/unstable/rpc/RpcClientError" import { TelemetryStoreReadonly, type TelemetryStoreReader } from "./TelemetryStore.js" import { QueryRpcs } from "./queryRpc.js" @@ -16,28 +16,45 @@ const WorkerProtocol = RpcClient.layerProtocolWorker({ size: 1 }).pipe( Layer.provide(BunWorker.layer(() => new Worker(new URL("./telemetryQueryWorker.ts", import.meta.url)))), ) -const query = (getClient: Effect.Effect, invalidateClient: Effect.Effect, method: QueryMethod, args: readonly unknown[] = []) => - Effect.flatMap(getClient, ({ client, clientScope }) => client.query({ method, args }).pipe( - Effect.onError(() => Effect.andThen(Scope.close(clientScope, Exit.void), invalidateClient)), - )).pipe( +const query = ( + getClient: Effect.Effect, + invalidateClient: Effect.Effect, + method: QueryMethod, + args: readonly unknown[] = [], +) => + Effect.flatMap(getClient, ({ client, clientScope }) => + client + .query({ method, args }) + .pipe( + Effect.tapError((error) => + error instanceof RpcClientError || error instanceof WorkerError + ? Effect.andThen(Scope.close(clientScope, Exit.void), invalidateClient) + : Effect.void, + ), + ), + ).pipe( Effect.map((result) => result as A), - Effect.mapError((error) => error instanceof Error ? error : new Error(String(error))), + Effect.mapError((error) => (error instanceof Error ? error : new Error(String(error)))), ) export const TelemetryQueryLive = Layer.effect( TelemetryStoreReadonly, - Effect.gen(function*() { + Effect.gen(function* () { const scope = yield* Scope.Scope - const [getClient, invalidateClient] = yield* Effect.cachedInvalidateWithTTL(Effect.gen(function*() { - const clientScope = yield* Scope.fork(scope, "sequential") - const protocolContext = yield* Layer.buildWithScope(WorkerProtocol, clientScope) - const client = yield* RpcClient.make(QueryRpcs).pipe( - Effect.provide(protocolContext), - Effect.provideService(Scope.Scope, clientScope), - ) - return { client, clientScope } - }), Duration.infinity) - const run = (method: QueryMethod, args: readonly unknown[] = []) => query(getClient, invalidateClient, method, args) + const [getClient, invalidateClient] = yield* Effect.cachedInvalidateWithTTL( + Effect.gen(function* () { + const clientScope = yield* Scope.fork(scope, "sequential") + const protocolContext = yield* Layer.buildWithScope(WorkerProtocol, clientScope) + const client = yield* RpcClient.make(QueryRpcs).pipe( + Effect.provide(protocolContext), + Effect.provideService(Scope.Scope, clientScope), + ) + return { client, clientScope } + }), + Duration.infinity, + ) + const run = (method: QueryMethod, args: readonly unknown[] = []) => + query(getClient, invalidateClient, method, args) return TelemetryStoreReadonly.of({ listServices: run("listServices"), listRecentTraces: (serviceName, options) => run("listRecentTraces", [serviceName, options]),