From 1b9274f9430c99596086e04aa4cd5ace708edf18 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Mon, 5 Oct 2026 15:06:35 +0200 Subject: [PATCH] refactor(alerts): drop fetch from AlertRuntime, use the runtime HttpClient Alert delivery already runs on the HttpClient from the app runtime, but AlertRuntime still carried a raw `fetch` that dispatchDelivery used only to override FetchHttpClient.Fetch per send, and that the save-time checks (Telegram chat discovery, Telegram credential check, PagerDuty routing key check) called directly. - Remove `fetch` from AlertRuntimeApi and `fetchFn` from dispatchDelivery / TransportRuntime (NotificationDispatcher passed globalThis.fetch too). - fetchTelegramChats, verifyTelegramCredentials and verifyPagerDutyRoutingKey require HttpClient; AlertDestinationsService provides the client it acquires in its constructor. Telegram calls disable the client span because the bot token is in the URL path. - Tests provide their fake wire as FetchHttpClient.Fetch. --- .../v2/alchemy-provider.integration.test.ts | 1 - apps/api/src/routes/v2/alerts.http.test.ts | 1 - .../AlertDeliveryDispatch.providers.test.ts | 16 +- .../alerts/AlertDeliveryDispatch.test.ts | 22 ++- .../alerts/AlertDestinationDelivery.ts | 1 - .../alerts/AlertDestinationsService.ts | 10 +- .../services/alerts/AlertRulesService.test.ts | 1 - .../src/services/alerts/AlertRuntime.ts | 2 - .../src/services/alerts/AlertsService.test.ts | 13 +- .../services/alerts/NotificationDispatcher.ts | 1 - .../alerts/delivery/delivery-spans.test.ts | 16 +- .../src/services/alerts/delivery/dispatch.ts | 3 +- .../services/alerts/delivery/runTransport.ts | 4 +- .../alerts/delivery/transports/chat.test.ts | 8 +- .../alerts/delivery/transports/pagerduty.ts | 31 ++-- .../delivery/transports/telegram.test.ts | 29 ++-- .../alerts/delivery/transports/telegram.ts | 160 +++++++++--------- 17 files changed, 181 insertions(+), 138 deletions(-) diff --git a/apps/api/src/routes/v2/alchemy-provider.integration.test.ts b/apps/api/src/routes/v2/alchemy-provider.integration.test.ts index 9fb873e715..8fc567e82b 100644 --- a/apps/api/src/routes/v2/alchemy-provider.integration.test.ts +++ b/apps/api/src/routes/v2/alchemy-provider.integration.test.ts @@ -113,7 +113,6 @@ const makeHarness = () => { const runtimeLive = Layer.succeed(AlertRuntime, { now: Effect.sync(() => Date.now()), makeUuid: () => crypto.randomUUID(), - fetch: globalThis.fetch, deliveryTimeoutMs: () => 15_000, }) const hazelOAuthLive = HazelOAuthService.layer.pipe(Layer.provide(Layer.mergeAll(envLive, testDb.layer))) diff --git a/apps/api/src/routes/v2/alerts.http.test.ts b/apps/api/src/routes/v2/alerts.http.test.ts index 0d6ea0612e..33614adb65 100644 --- a/apps/api/src/routes/v2/alerts.http.test.ts +++ b/apps/api/src/routes/v2/alerts.http.test.ts @@ -97,7 +97,6 @@ const makeHarness = ( const runtimeLive = Layer.succeed(AlertRuntime, { now: Effect.sync(() => Date.now()), makeUuid: () => crypto.randomUUID(), - fetch: globalThis.fetch, deliveryTimeoutMs: () => 15_000, }) const hazelOAuthLive = diff --git a/packages/backend/src/services/alerts/AlertDeliveryDispatch.providers.test.ts b/packages/backend/src/services/alerts/AlertDeliveryDispatch.providers.test.ts index 7ea26d4e11..42551a3924 100644 --- a/packages/backend/src/services/alerts/AlertDeliveryDispatch.providers.test.ts +++ b/packages/backend/src/services/alerts/AlertDeliveryDispatch.providers.test.ts @@ -8,9 +8,19 @@ import { dispatchDelivery as dispatchDeliveryRaw } from "./delivery/dispatch" import { FetchHttpClient } from "effect/http" import type { EffectTransportDeps } from "./delivery/Transport" -/** The runtime provides the HTTP client; each test passes its own `fetch`. */ -const dispatchDelivery = (...args: Parameters) => - dispatchDeliveryRaw(...args).pipe(Effect.provide(FetchHttpClient.layer)) +type DispatchArgs = Parameters + +/** The runtime provides the HTTP client; each test fakes the wire as `FetchHttpClient.Fetch`. */ +const dispatchDelivery = ( + context: DispatchArgs[0], + payloadJson: DispatchArgs[1], + fetchFn: typeof fetch, + ...rest: [DispatchArgs[2], DispatchArgs[3], DispatchArgs[4], DispatchArgs[5]] +) => + dispatchDeliveryRaw(context, payloadJson, ...rest).pipe( + Effect.provide(FetchHttpClient.layer), + Effect.provideService(FetchHttpClient.Fetch, fetchFn), + ) /** * Characterization tests: they pin what each provider ACTUALLY sends today, diff --git a/packages/backend/src/services/alerts/AlertDeliveryDispatch.test.ts b/packages/backend/src/services/alerts/AlertDeliveryDispatch.test.ts index 7870abc2b3..1ac2331d63 100644 --- a/packages/backend/src/services/alerts/AlertDeliveryDispatch.test.ts +++ b/packages/backend/src/services/alerts/AlertDeliveryDispatch.test.ts @@ -23,9 +23,19 @@ import { resolveSignalDisplay } from "./alert-signal-display" import { renderTemplate } from "./alert-templating/renderer" import { DEFAULT_BODY_TEMPLATE, DEFAULT_TITLE_TEMPLATE } from "./alert-templating/defaultTemplates" -/** The runtime provides the HTTP client; each test passes its own `fetch`. */ -const dispatchDelivery = (...args: Parameters) => - dispatchDeliveryRaw(...args).pipe(Effect.provide(FetchHttpClient.layer)) +type DispatchArgs = Parameters + +/** The runtime provides the HTTP client; each test fakes the wire as `FetchHttpClient.Fetch`. */ +const dispatchDelivery = ( + context: DispatchArgs[0], + payloadJson: DispatchArgs[1], + fetchFn: typeof fetch, + ...rest: [DispatchArgs[2], DispatchArgs[3], DispatchArgs[4], DispatchArgs[5]] +) => + dispatchDeliveryRaw(context, payloadJson, ...rest).pipe( + Effect.provide(FetchHttpClient.layer), + Effect.provideService(FetchHttpClient.Fetch, fetchFn), + ) /** Chat posts must not happen for these destinations. */ const failingChatPost = () => @@ -278,11 +288,11 @@ describe("dispatchDelivery", () => { }), ) - it.effect("calls an unguarded transport's fetch detached from the runtime object", () => + it.effect("calls an unguarded transport's fetch as a bare function", () => Effect.gen(function* () { - // Regression: the unguarded transports used to call `runtime.fetchFn(...)`, + // Regression: the unguarded transports once called `runtime.fetchFn(...)`, // a method call that hands workerd's global `fetch` a `this` of the - // runtime object — "Illegal invocation", every such delivery dead. + // runtime object: "Illegal invocation", every such delivery dead. // A `function` (not an arrow) is what makes `this` observable here. let called = false let receiver: typeof globalThis | undefined diff --git a/packages/backend/src/services/alerts/AlertDestinationDelivery.ts b/packages/backend/src/services/alerts/AlertDestinationDelivery.ts index 7ce979c58d..01d07d3f8a 100644 --- a/packages/backend/src/services/alerts/AlertDestinationDelivery.ts +++ b/packages/backend/src/services/alerts/AlertDestinationDelivery.ts @@ -125,7 +125,6 @@ export const makeAlertDestinationDelivery = (options: { dispatchDeliveryImpl( context, payloadJson, - options.runtime.fetch, options.runtime.deliveryTimeoutMs(), context.linkUrl, composeChatUrl(context), diff --git a/packages/backend/src/services/alerts/AlertDestinationsService.ts b/packages/backend/src/services/alerts/AlertDestinationsService.ts index 8017af19c1..c24c27544e 100644 --- a/packages/backend/src/services/alerts/AlertDestinationsService.ts +++ b/packages/backend/src/services/alerts/AlertDestinationsService.ts @@ -442,10 +442,9 @@ export class AlertDestinationsService extends Context.Service< } const result = yield* verifyPagerDutyRoutingKey( integrationKey, - runtime.fetch, runtime.deliveryTimeoutMs(), `maple-keycheck-${runtime.makeUuid()}`, - ) + ).pipe(Effect.provideService(HttpClient.HttpClient, httpClient)) if (result.status === "invalid") { return yield* Effect.fail( makeValidationError(`PagerDuty rejected this routing key: ${result.reason}`), @@ -463,9 +462,8 @@ export class AlertDestinationsService extends Context.Service< const result = yield* verifyTelegramCredentials( botToken, chatId, - runtime.fetch, runtime.deliveryTimeoutMs(), - ) + ).pipe(Effect.provideService(HttpClient.HttpClient, httpClient)) if (result.status === "invalid") { return yield* Effect.fail(makeValidationError(result.reason)) } @@ -482,7 +480,9 @@ export class AlertDestinationsService extends Context.Service< if (!TELEGRAM_BOT_TOKEN_PATTERN.test(trimmed)) { return yield* Effect.fail(makeValidationError(TELEGRAM_MALFORMED_TOKEN_MESSAGE)) } - const result = yield* fetchTelegramChats(trimmed, runtime.fetch, runtime.deliveryTimeoutMs()) + const result = yield* fetchTelegramChats(trimmed, runtime.deliveryTimeoutMs()).pipe( + Effect.provideService(HttpClient.HttpClient, httpClient), + ) if (result.status === "invalid") return yield* Effect.fail(makeValidationError(result.reason)) return result.chats }) diff --git a/packages/backend/src/services/alerts/AlertRulesService.test.ts b/packages/backend/src/services/alerts/AlertRulesService.test.ts index e125575e74..a57e59b8c6 100644 --- a/packages/backend/src/services/alerts/AlertRulesService.test.ts +++ b/packages/backend/src/services/alerts/AlertRulesService.test.ts @@ -41,7 +41,6 @@ afterEach(() => cleanupTestDbs(createdDbs)) const runtime: AlertRuntimeApi = { now: Effect.succeed(NOW), makeUuid: () => RULE, - fetch: globalThis.fetch, deliveryTimeoutMs: () => 15_000, } diff --git a/packages/backend/src/services/alerts/AlertRuntime.ts b/packages/backend/src/services/alerts/AlertRuntime.ts index 4ba27a3277..e0aac999d0 100644 --- a/packages/backend/src/services/alerts/AlertRuntime.ts +++ b/packages/backend/src/services/alerts/AlertRuntime.ts @@ -7,7 +7,6 @@ export interface AlertRuntimeApi { /** Current wall-clock time in epoch ms, sourced from Effect's `Clock` so tests drive it via `TestClock`. */ readonly now: Effect.Effect readonly makeUuid: () => string - readonly fetch: typeof fetch readonly deliveryTimeoutMs: () => number } @@ -15,7 +14,6 @@ export class AlertRuntime extends Context.Reference("@maple/api defaultValue: (): AlertRuntimeApi => ({ now: Clock.currentTimeMillis, makeUuid: () => randomUUID(), - fetch: globalThis.fetch, deliveryTimeoutMs: () => DELIVERY_TIMEOUT_MS_DEFAULT, }), }) {} diff --git a/packages/backend/src/services/alerts/AlertsService.test.ts b/packages/backend/src/services/alerts/AlertsService.test.ts index 17f036841d..66be43ff10 100644 --- a/packages/backend/src/services/alerts/AlertsService.test.ts +++ b/packages/backend/src/services/alerts/AlertsService.test.ts @@ -194,10 +194,14 @@ const defaultTestRuntime: AlertRuntimeApi = { // TestClock.adjust. Real `fetch`/`Effect.timeout` settle on the live event loop. now: Clock.currentTimeMillis, makeUuid: () => crypto.randomUUID(), - fetch: globalThis.fetch, deliveryTimeoutMs: () => 15_000, } +interface TestOverrides extends Partial { + /** Provided as `FetchHttpClient.Fetch`, so every outbound call goes through it. */ + readonly fetch?: typeof fetch +} + // The fixed epoch scheduler tests start TestClock at, mirroring the previous // manual clock's default start time. const DEFAULT_CLOCK_EPOCH_MS = 1_700_000_000_000 @@ -245,10 +249,11 @@ const stubOrgMembersService = ( const makeLayer = ( testDb: TestDb, warehouseStub: WarehouseQueryServiceApi, - runtimeOverrides?: Partial, + overrides: TestOverrides = {}, emailStub?: (typeof EmailService)["Service"], chatAlertPoster: Layer.Layer = ChatAlertPoster.layer, ) => { + const { fetch: fetchImpl = globalThis.fetch, ...runtimeOverrides } = overrides const configLive = makeConfig() const envLive = Env.layer.pipe(Layer.provide(configLive)) const databaseLive = testDb.layer @@ -312,7 +317,9 @@ const makeLayer = ( Layer.provide(alertReadModelsLive), Layer.provide(alertRulesLive), ) - return Layer.mergeAll(alertDestinationsLive, alertReadModelsLive, alertRulesLive, alertsLive) + return Layer.mergeAll(alertDestinationsLive, alertReadModelsLive, alertRulesLive, alertsLive).pipe( + Layer.provideMerge(Layer.succeed(FetchHttpClient.Fetch, fetchImpl)), + ) } const asOrgId = Schema.decodeUnknownSync(OrgId) diff --git a/packages/backend/src/services/alerts/NotificationDispatcher.ts b/packages/backend/src/services/alerts/NotificationDispatcher.ts index dc7e4ebc2e..9165b54f6c 100644 --- a/packages/backend/src/services/alerts/NotificationDispatcher.ts +++ b/packages/backend/src/services/alerts/NotificationDispatcher.ts @@ -211,7 +211,6 @@ const make: Effect.Effect< const result = yield* dispatchDeliveryImpl( context, payloadJson, - globalThis.fetch, DELIVERY_TIMEOUT_MS, request.linkUrl, chatUrl, diff --git a/packages/backend/src/services/alerts/delivery/delivery-spans.test.ts b/packages/backend/src/services/alerts/delivery/delivery-spans.test.ts index 16e898afff..8a37b727d3 100644 --- a/packages/backend/src/services/alerts/delivery/delivery-spans.test.ts +++ b/packages/backend/src/services/alerts/delivery/delivery-spans.test.ts @@ -8,9 +8,19 @@ import { FetchHttpClient } from "effect/http" import type { EffectTransportDeps } from "./Transport" import type { DispatchContext } from "./context" -/** The runtime provides the HTTP client; each test passes its own `fetch`. */ -const dispatchDelivery = (...args: Parameters) => - dispatchDeliveryRaw(...args).pipe(Effect.provide(FetchHttpClient.layer)) +type DispatchArgs = Parameters + +/** The runtime provides the HTTP client; each test fakes the wire as `FetchHttpClient.Fetch`. */ +const dispatchDelivery = ( + context: DispatchArgs[0], + payloadJson: DispatchArgs[1], + fetchFn: typeof fetch, + ...rest: [DispatchArgs[2], DispatchArgs[3], DispatchArgs[4], DispatchArgs[5]] +) => + dispatchDeliveryRaw(context, payloadJson, ...rest).pipe( + Effect.provide(FetchHttpClient.layer), + Effect.provideService(FetchHttpClient.Fetch, fetchFn), + ) /** * The outbound provider call must be a **Client-kind span carrying diff --git a/packages/backend/src/services/alerts/delivery/dispatch.ts b/packages/backend/src/services/alerts/delivery/dispatch.ts index 1853887cf5..1b7b385bb5 100644 --- a/packages/backend/src/services/alerts/delivery/dispatch.ts +++ b/packages/backend/src/services/alerts/delivery/dispatch.ts @@ -26,14 +26,13 @@ import type { DispatchContext, DispatchResult } from "./context" export const dispatchDelivery = ( context: DispatchContext, payloadJson: string, - fetchFn: typeof fetch, timeoutMs: number, linkUrl: string, chatUrl: string, /** `sendEmail` (the platform email channel) and `postChatAlert` (a chat connector). */ deps: EffectTransportDeps, ): Effect.Effect => { - const runtime: TransportRuntime = { fetchFn, timeoutMs } + const runtime: TransportRuntime = { timeoutMs } /** * Resolved once here rather than by each provider: the template is a * property of the rule and the destination TYPE, not of the transport's diff --git a/packages/backend/src/services/alerts/delivery/runTransport.ts b/packages/backend/src/services/alerts/delivery/runTransport.ts index e49a0b8ca6..a47a0ec8e7 100644 --- a/packages/backend/src/services/alerts/delivery/runTransport.ts +++ b/packages/backend/src/services/alerts/delivery/runTransport.ts @@ -12,7 +12,7 @@ import { } from "@maple/domain/http" import { Duration, Effect, Option, Result } from "effect" import { constTrue } from "effect/Function" -import { FetchHttpClient, HttpBody, HttpClient, HttpClientRequest } from "effect/http" +import { HttpBody, HttpClient, HttpClientRequest } from "effect/http" import type { HttpClientResponse } from "effect/http" import { describeHttpClientError, guard } from "@maple/safe-fetch" import type { @@ -25,7 +25,6 @@ import type { import type { DispatchResult } from "./context" export interface TransportRuntime { - readonly fetchFn: typeof fetch readonly timeoutMs: number } @@ -134,7 +133,6 @@ const sendHttp = Effect.fn("AlertDelivery.http", { kind: "client" })(function* ( // cancelled when `runHttpTransport`'s scope closes. const scoped = HttpClient.withScope(client) const response = yield* (spec.guarded ? guard(scoped) : scoped).execute(request).pipe( - Effect.provideService(FetchHttpClient.Fetch, runtime.fetchFn), // The client's own span records `url.full`, and Discord, Hazel and // Telegram carry their delivery token in the URL path. This span is the // client span. diff --git a/packages/backend/src/services/alerts/delivery/transports/chat.test.ts b/packages/backend/src/services/alerts/delivery/transports/chat.test.ts index 41dfd59d7c..41334cdb67 100644 --- a/packages/backend/src/services/alerts/delivery/transports/chat.test.ts +++ b/packages/backend/src/services/alerts/delivery/transports/chat.test.ts @@ -3,6 +3,7 @@ import { ChatOutboundError, chatConnectorId, type ChatAlertBlock, type ChatBlock import { AlertDestinationId, ChatWorkspaceId, OrgId } from "@maple/domain/http" import { assert, describe, it } from "@effect/vitest" import { Effect, Schema } from "effect" +import { FetchHttpClient } from "effect/http" import { chatDeliveryFailure } from "../../ChatAlertPoster" import type { DispatchContext } from "../context" import { dispatchDelivery } from "../dispatch" @@ -180,14 +181,17 @@ describe("chat through dispatchDelivery", () => { throw new Error("a chat destination made an HTTP request of its own") } return Effect.gen(function* () { - const result = yield* dispatchDelivery(context, "{}", fetchFn, 5_000, LINK, CHAT, { + const result = yield* dispatchDelivery(context, "{}", 5_000, LINK, CHAT, { sendEmail: () => Effect.die("the chat transport sent an email"), postChatAlert: (post) => Effect.sync(() => { posts.push(post) return { connectorName: "Test Chat", messageId: "message-1" } }), - }) + }).pipe( + Effect.provide(FetchHttpClient.layer), + Effect.provideService(FetchHttpClient.Fetch, fetchFn), + ) assert.strictEqual(result.providerMessage, "Delivered to Test Chat #incidents") assert.strictEqual(posts[0]?.workspaceId, WORKSPACE) assert.include(cardOf(posts[0]?.blocks ?? []).title, "Checkout error rate") diff --git a/packages/backend/src/services/alerts/delivery/transports/pagerduty.ts b/packages/backend/src/services/alerts/delivery/transports/pagerduty.ts index 1866a67918..990d59b19d 100644 --- a/packages/backend/src/services/alerts/delivery/transports/pagerduty.ts +++ b/packages/backend/src/services/alerts/delivery/transports/pagerduty.ts @@ -1,4 +1,5 @@ import { Duration, Effect } from "effect" +import { HttpClient, HttpClientRequest } from "effect/http" import { displayGroupKey, truncate } from "../../alert-formatting" import type { HttpTransport, RenderInput, SecretConfigOf } from "../Transport" @@ -84,25 +85,27 @@ export type PagerDutyKeyVerification = */ export const verifyPagerDutyRoutingKey = ( integrationKey: string, - fetchFn: typeof fetch, timeoutMs: number, dedupKey: string, -): Effect.Effect => - Effect.tryPromise(() => - fetchFn("https://events.pagerduty.com/v2/enqueue", { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ - routing_key: integrationKey, - event_action: "resolve", - dedup_key: dedupKey, - }), - }), +): Effect.Effect => + HttpClient.execute( + HttpClientRequest.post("https://events.pagerduty.com/v2/enqueue").pipe( + HttpClientRequest.bodyText( + JSON.stringify({ + routing_key: integrationKey, + event_action: "resolve", + dedup_key: dedupKey, + }), + "application/json", + ), + ), ).pipe( Effect.flatMap((response) => { - if (response.ok) return Effect.succeed({ status: "valid" }) + if (response.status >= 200 && response.status < 300) { + return Effect.succeed({ status: "valid" }) + } if (response.status === 400) { - return Effect.promise(() => response.text().catch(() => "")).pipe( + return Effect.orElseSucceed(response.text, () => "").pipe( Effect.map((body): PagerDutyKeyVerification => { const reason = truncate(body.trim().replace(/\s+/g, " "), 500) return { status: "invalid", reason: reason || "Invalid routing key" } diff --git a/packages/backend/src/services/alerts/delivery/transports/telegram.test.ts b/packages/backend/src/services/alerts/delivery/transports/telegram.test.ts index bb4c4500e4..867f87bb3a 100644 --- a/packages/backend/src/services/alerts/delivery/transports/telegram.test.ts +++ b/packages/backend/src/services/alerts/delivery/transports/telegram.test.ts @@ -2,6 +2,7 @@ import type { AlertDestinationRow } from "@maple/db" import { AlertDestinationId } from "@maple/domain/http" import { assert, describe, it } from "@effect/vitest" import { Effect, Result, Schema } from "effect" +import { FetchHttpClient, type HttpClient } from "effect/http" import type { DispatchContext } from "../context" import type { RenderInput, SecretConfigOf } from "../Transport" import { @@ -18,6 +19,10 @@ const CHAT_ID = "-1001234567890" const LINK = "https://web.localhost/alerts" const CHAT = "https://web.localhost/chat" +/** Runs a Bot API call on the fetch client with `fetchFn` as the wire. */ +const withFetch = (effect: Effect.Effect, fetchFn: typeof fetch) => + effect.pipe(Effect.provide(FetchHttpClient.layer), Effect.provideService(FetchHttpClient.Fetch, fetchFn)) + const destinationRow: AlertDestinationRow = { id: DESTINATION_ID, orgId: "org_1" as AlertDestinationRow["orgId"], @@ -235,7 +240,7 @@ describe("verifyTelegramCredentials", () => { it.effect("is valid when both getMe and getChat succeed", () => Effect.gen(function* () { const { fetchFn, calls } = stub([ok, ok]) - const result = yield* verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, fetchFn, 1000) + const result = yield* withFetch(verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, 1000), fetchFn) assert.deepStrictEqual(result, { status: "valid" }) assert.strictEqual(calls.length, 2) assert.include(calls[1]!, `chat_id=${encodeURIComponent(CHAT_ID)}`) @@ -245,7 +250,7 @@ describe("verifyTelegramCredentials", () => { it.effect("rejects a bad token without asking about the chat", () => Effect.gen(function* () { const { fetchFn, calls } = stub([{ status: 401, body: { ok: false, error_code: 401 } }]) - const result = yield* verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, fetchFn, 1000) + const result = yield* withFetch(verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, 1000), fetchFn) assert.strictEqual(result.status, "invalid") assert.strictEqual(calls.length, 1) }), @@ -261,7 +266,7 @@ describe("verifyTelegramCredentials", () => { body: { ok: false, error_code: 400, description: "Bad Request: chat not found" }, }, ]) - const result = yield* verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, fetchFn, 1000) + const result = yield* withFetch(verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, 1000), fetchFn) assert.strictEqual(result.status, "invalid") if (result.status === "invalid") assert.include(result.reason, "chat not found") }), @@ -270,7 +275,7 @@ describe("verifyTelegramCredentials", () => { it.effect("fails open on a 5xx so a Telegram outage cannot block a save", () => Effect.gen(function* () { const { fetchFn } = stub([{ status: 503, body: {} }]) - const result = yield* verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, fetchFn, 1000) + const result = yield* withFetch(verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, 1000), fetchFn) assert.deepStrictEqual(result, { status: "unknown" }) }), ) @@ -278,7 +283,7 @@ describe("verifyTelegramCredentials", () => { it.effect("fails open when the request throws", () => Effect.gen(function* () { const fetchFn: typeof fetch = () => Promise.reject(new Error("network down")) - const result = yield* verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, fetchFn, 1000) + const result = yield* withFetch(verifyTelegramCredentials(BOT_TOKEN, CHAT_ID, 1000), fetchFn) assert.deepStrictEqual(result, { status: "unknown" }) }), ) @@ -312,7 +317,7 @@ describe("fetchTelegramChats", () => { ok: true, result: [{ my_chat_member: { chat: chatOf(-1001234567890, "Acme On-call", "supergroup") } }], }) - const result = yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + const result = yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.deepStrictEqual(result, { status: "ok", chats: [{ id: "-1001234567890", title: "Acme On-call", type: "supergroup" }], @@ -328,7 +333,7 @@ describe("fetchTelegramChats", () => { it.effect("does not confirm updates — no offset is ever sent", () => Effect.gen(function* () { const { fetchFn, calls } = respondWith(200, { ok: true, result: [] }) - yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.lengthOf(calls, 1) assert.notInclude(calls[0]!, "offset") }), @@ -344,7 +349,7 @@ describe("fetchTelegramChats", () => { { channel_post: { chat: chatOf(-100222, "Newer", "channel") } }, ], }) - const result = yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + const result = yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.strictEqual(result.status, "ok") if (result.status === "ok") { assert.deepStrictEqual( @@ -361,7 +366,7 @@ describe("fetchTelegramChats", () => { ok: true, result: [{ message: { chat: { id: 42, type: "private", username: "ada" } } }], }) - const result = yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + const result = yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.strictEqual(result.status, "ok") if (result.status === "ok") { assert.deepStrictEqual(result.chats, [{ id: "42", title: "ada", type: "private" }]) @@ -377,7 +382,7 @@ describe("fetchTelegramChats", () => { error_code: 409, description: "Conflict: can't use getUpdates method while webhook is active", }) - const result = yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + const result = yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.strictEqual(result.status, "invalid") if (result.status === "invalid") assert.include(result.reason, "webhook") }), @@ -386,7 +391,7 @@ describe("fetchTelegramChats", () => { it.effect("reports a rejected token as such", () => Effect.gen(function* () { const { fetchFn } = respondWith(401, { ok: false, error_code: 401, description: "Unauthorized" }) - const result = yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + const result = yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.deepStrictEqual(result, { status: "invalid", reason: "Telegram rejected the bot token" }) }), ) @@ -394,7 +399,7 @@ describe("fetchTelegramChats", () => { it.effect("never throws when the request fails", () => Effect.gen(function* () { const fetchFn: typeof fetch = () => Promise.reject(new Error("network down")) - const result = yield* fetchTelegramChats(BOT_TOKEN, fetchFn, 1000) + const result = yield* withFetch(fetchTelegramChats(BOT_TOKEN, 1000), fetchFn) assert.strictEqual(result.status, "invalid") }), ) diff --git a/packages/backend/src/services/alerts/delivery/transports/telegram.ts b/packages/backend/src/services/alerts/delivery/transports/telegram.ts index fc6d4349dc..1dee65d8bb 100644 --- a/packages/backend/src/services/alerts/delivery/transports/telegram.ts +++ b/packages/backend/src/services/alerts/delivery/transports/telegram.ts @@ -7,6 +7,8 @@ import { type AlertDeliveryFailure, } from "@maple/domain/http" import { Duration, Effect, Result, Schema } from "effect" +import { constTrue } from "effect/Function" +import { HttpClient } from "effect/http" import { buildTelegramText, buildTelegramTextFromTemplate } from "../../AlertDeliveryDispatch" import { truncate } from "../../alert-formatting" import type { HttpTransport, ProviderAck, RenderInput, SecretConfigOf } from "../Transport" @@ -141,27 +143,38 @@ export type TelegramCredentialVerification = /** Network error / timeout / 429 / 5xx — can't conclude; caller should fail open. */ | { status: "unknown" } -const telegramApiCall = ( - url: string, - fetchFn: typeof fetch, -): Effect.Effect<{ ok: boolean; status: number; description: string }> => - Effect.tryPromise(() => fetchFn(url, { method: "GET" })).pipe( +/** + * A Bot API GET with its body read as text (empty when unreadable). The bot + * token is in the URL path, so the client's own span, which records + * `url.full`, stays off. + */ +const telegramGet = (url: string) => + HttpClient.get(url).pipe( Effect.flatMap((response) => - Effect.promise(() => response.text().catch(() => "")).pipe( - Effect.map((body) => { - const decoded = decodeTelegramResponseJson(body) - const description = Result.getOrElse( - Result.map(decoded, (payload) => payload.description ?? ""), - () => "", - ) - return { - ok: response.ok, - status: response.status, - description: truncate(description.replace(/\s+/g, " ").trim(), 300), - } + Effect.map( + Effect.orElseSucceed(response.text, () => ""), + (body) => ({ + ok: response.status >= 200 && response.status < 300, + status: response.status, + body, }), ), ), + Effect.provideService(HttpClient.TracerDisabledWhen, constTrue), + ) + +const telegramApiCall = ( + url: string, +): Effect.Effect<{ ok: boolean; status: number; description: string }, never, HttpClient.HttpClient> => + telegramGet(url).pipe( + Effect.map(({ ok, status, body }) => { + const decoded = decodeTelegramResponseJson(body) + const description = Result.getOrElse( + Result.map(decoded, (payload) => payload.description ?? ""), + () => "", + ) + return { ok, status, description: truncate(description.replace(/\s+/g, " ").trim(), 300) } + }), Effect.orElseSucceed(() => ({ ok: false, status: 0, description: "" })), ) @@ -179,18 +192,17 @@ const telegramApiCall = ( export const verifyTelegramCredentials = ( botToken: string, chatId: string, - fetchFn: typeof fetch, timeoutMs: number, -): Effect.Effect => +): Effect.Effect => Effect.gen(function* () { const base = `${TELEGRAM_API_ORIGIN}/bot${botToken}` - const me = yield* telegramApiCall(`${base}/getMe`, fetchFn) + const me = yield* telegramApiCall(`${base}/getMe`) if (me.status === 401 || me.status === 404) { return { status: "invalid", reason: "Telegram rejected the bot token" } as const } if (!me.ok) return { status: "unknown" } as const - const chat = yield* telegramApiCall(`${base}/getChat?chat_id=${encodeURIComponent(chatId)}`, fetchFn) + const chat = yield* telegramApiCall(`${base}/getChat?chat_id=${encodeURIComponent(chatId)}`) if (chat.status === 400 || chat.status === 403 || chat.status === 404) { return { status: "invalid", @@ -292,69 +304,61 @@ const narrowChatType = (raw: string): TelegramChat["type"] | null => */ export const fetchTelegramChats = ( botToken: string, - fetchFn: typeof fetch, timeoutMs: number, -): Effect.Effect => - Effect.tryPromise(() => - fetchFn( - `${TELEGRAM_API_ORIGIN}/bot${botToken}/getUpdates?limit=100&timeout=0&allowed_updates=${encodeURIComponent( - JSON.stringify(["message", "edited_message", "channel_post", "my_chat_member"]), - )}`, - { method: "GET" }, - ), +): Effect.Effect => + telegramGet( + `${TELEGRAM_API_ORIGIN}/bot${botToken}/getUpdates?limit=100&timeout=0&allowed_updates=${encodeURIComponent( + JSON.stringify(["message", "edited_message", "channel_post", "my_chat_member"]), + )}`, ).pipe( - Effect.flatMap((response) => - Effect.promise(() => response.text().catch(() => "")).pipe( - Effect.map((raw): TelegramChatDiscovery => { - const decoded = decodeGetUpdates(raw) - if (Result.isFailure(decoded)) { - return response.ok - ? { status: "invalid", reason: "Telegram returned an unexpected response" } - : { - status: "invalid", - reason: `Telegram rejected the request (${response.status})`, - } - } - const payload = decoded.success - if (!payload.ok) { - const description = payload.description ?? "" - if (payload.error_code === 409 || description.includes("webhook is active")) { - return { - status: "invalid", - reason: "This bot has a webhook registered, so Maple cannot read its recent chats. Delete the webhook (or enter the chat ID by hand).", - } - } - if (payload.error_code === 401) { - return { status: "invalid", reason: "Telegram rejected the bot token" } - } - return { + Effect.map((response): TelegramChatDiscovery => { + const decoded = decodeGetUpdates(response.body) + if (Result.isFailure(decoded)) { + return response.ok + ? { status: "invalid", reason: "Telegram returned an unexpected response" } + : { status: "invalid", - reason: - truncate(description.replace(/\s+/g, " ").trim(), 300) || - "Telegram rejected the request", + reason: `Telegram rejected the request (${response.status})`, } + } + const payload = decoded.success + if (!payload.ok) { + const description = payload.description ?? "" + if (payload.error_code === 409 || description.includes("webhook is active")) { + return { + status: "invalid", + reason: "This bot has a webhook registered, so Maple cannot read its recent chats. Delete the webhook (or enter the chat ID by hand).", } + } + if (payload.error_code === 401) { + return { status: "invalid", reason: "Telegram rejected the bot token" } + } + return { + status: "invalid", + reason: + truncate(description.replace(/\s+/g, " ").trim(), 300) || + "Telegram rejected the request", + } + } - const byId = new Map() - // Newest first: `getUpdates` returns ascending `update_id`, and the - // chat someone just added the bot to is the one they are looking for. - for (const update of [...(payload.result ?? [])].reverse()) { - const chat = ( - update.my_chat_member ?? - update.message ?? - update.channel_post ?? - update.edited_message - )?.chat - if (chat === undefined) continue - const type = narrowChatType(chat.type) - if (type === null) continue - const id = String(chat.id) - if (!byId.has(id)) byId.set(id, { id, title: chatLabel(chat), type }) - } - return { status: "ok", chats: [...byId.values()] } - }), - ), - ), + const byId = new Map() + // Newest first: `getUpdates` returns ascending `update_id`, and the + // chat someone just added the bot to is the one they are looking for. + for (const update of [...(payload.result ?? [])].reverse()) { + const chat = ( + update.my_chat_member ?? + update.message ?? + update.channel_post ?? + update.edited_message + )?.chat + if (chat === undefined) continue + const type = narrowChatType(chat.type) + if (type === null) continue + const id = String(chat.id) + if (!byId.has(id)) byId.set(id, { id, title: chatLabel(chat), type }) + } + return { status: "ok", chats: [...byId.values()] } + }), Effect.timeoutOrElse({ duration: Duration.millis(timeoutMs), orElse: () =>