From dc0d65b04d894661e57cb9572af71f4d065611a1 Mon Sep 17 00:00:00 2001 From: Dread Date: Fri, 5 Jun 2026 17:27:41 -0700 Subject: [PATCH 1/7] feat(alerts): add Bridge AlertService (PagerDuty/Slack/Discord) [ENG-361] Severity-routed best-effort alert fan-out: critical pages PagerDuty + informs Slack/Mattermost + Discord; warning informs only. Each sender no-ops when its env credential is unset. Config: ALERT_PAGERDUTY_ROUTING_KEY / ALERT_SLACK_WEBHOOK_URL / ALERT_DISCORD_WEBHOOK_URL. Co-Authored-By: Claude Opus 4.8 (1M context) --- src/config/env.ts | 8 +++++++ src/config/index.ts | 3 +++ src/services/alerts/discord.ts | 34 ++++++++++++++++++++++++++++++ src/services/alerts/index.ts | 27 ++++++++++++++++++++++++ src/services/alerts/index.types.ts | 13 ++++++++++++ src/services/alerts/pagerduty.ts | 33 +++++++++++++++++++++++++++++ src/services/alerts/slack.ts | 31 +++++++++++++++++++++++++++ 7 files changed, 149 insertions(+) create mode 100644 src/services/alerts/discord.ts create mode 100644 src/services/alerts/index.ts create mode 100644 src/services/alerts/index.types.ts create mode 100644 src/services/alerts/pagerduty.ts create mode 100644 src/services/alerts/slack.ts diff --git a/src/config/env.ts b/src/config/env.ts index 1216528c0..fe2397f8a 100644 --- a/src/config/env.ts +++ b/src/config/env.ts @@ -124,6 +124,10 @@ export const env = createEnv({ MATTERMOST_WEBHOOK_URL: z.string().min(1).optional(), + ALERT_PAGERDUTY_ROUTING_KEY: z.string().min(1).optional(), + ALERT_SLACK_WEBHOOK_URL: z.string().url().optional(), + ALERT_DISCORD_WEBHOOK_URL: z.string().url().optional(), + PROXY_CHECK_APIKEY: z.string().min(1).optional(), SVIX_SECRET: z.string().optional(), @@ -231,6 +235,10 @@ export const env = createEnv({ MATTERMOST_WEBHOOK_URL: process.env.MATTERMOST_WEBHOOK_URL, + ALERT_PAGERDUTY_ROUTING_KEY: process.env.ALERT_PAGERDUTY_ROUTING_KEY, + ALERT_SLACK_WEBHOOK_URL: process.env.ALERT_SLACK_WEBHOOK_URL, + ALERT_DISCORD_WEBHOOK_URL: process.env.ALERT_DISCORD_WEBHOOK_URL, + PROXY_CHECK_APIKEY: process.env.PROXY_CHECK_APIKEY, SVIX_SECRET: process.env.SVIX_SECRET, diff --git a/src/config/index.ts b/src/config/index.ts index 2b070e1a0..749a0933f 100644 --- a/src/config/index.ts +++ b/src/config/index.ts @@ -188,6 +188,9 @@ export const NEXTCLOUD_URL = env.NEXTCLOUD_URL export const NEXTCLOUD_USER = env.NEXTCLOUD_USER export const NEXTCLOUD_PASSWORD = env.NEXTCLOUD_PASSWORD export const MATTERMOST_WEBHOOK_URL = env.MATTERMOST_WEBHOOK_URL +export const ALERT_PAGERDUTY_ROUTING_KEY = env.ALERT_PAGERDUTY_ROUTING_KEY +export const ALERT_SLACK_WEBHOOK_URL = env.ALERT_SLACK_WEBHOOK_URL +export const ALERT_DISCORD_WEBHOOK_URL = env.ALERT_DISCORD_WEBHOOK_URL export const PROXY_CHECK_APIKEY = env.PROXY_CHECK_APIKEY export const NOSTR_PRIVATE_KEY = env.NOSTR_PRIVATE_KEY diff --git a/src/services/alerts/discord.ts b/src/services/alerts/discord.ts new file mode 100644 index 000000000..2dedb6c96 --- /dev/null +++ b/src/services/alerts/discord.ts @@ -0,0 +1,34 @@ +import { ALERT_DISCORD_WEBHOOK_URL } from "@config" +import { ErrorLevel } from "@domain/shared" +import { recordExceptionInCurrentSpan } from "@services/tracing" +import axios from "axios" + +import { BridgeAlert } from "./index.types" + +// Discord caps message content at 2000 chars; leave headroom. +const DISCORD_CONTENT_MAX = 1900 + +// Discord incoming webhook ({ content }). +export const sendDiscord = async (alert: BridgeAlert): Promise => { + if (!ALERT_DISCORD_WEBHOOK_URL) return + + const icon = alert.severity === "critical" ? "🚨" : "⚠️" + let content = `${icon} **Bridge alert** β€” ${alert.title}\nsource: \`${alert.source}\` Β· severity: \`${alert.severity}\`` + if (alert.detail) content += `\n${alert.detail}` + if (alert.context) { + content += "\n```json\n" + JSON.stringify(alert.context, null, 2) + "\n```" + } + if (content.length > DISCORD_CONTENT_MAX) { + content = content.slice(0, DISCORD_CONTENT_MAX) + "…" + } + + try { + await axios.post( + ALERT_DISCORD_WEBHOOK_URL, + { content }, + { timeout: 5000, headers: { "Content-Type": "application/json" } }, + ) + } catch (error) { + recordExceptionInCurrentSpan({ error, level: ErrorLevel.Warn }) + } +} diff --git a/src/services/alerts/index.ts b/src/services/alerts/index.ts new file mode 100644 index 000000000..733dc7e7f --- /dev/null +++ b/src/services/alerts/index.ts @@ -0,0 +1,27 @@ +import { sendPagerDuty } from "./pagerduty" +import { sendSlack } from "./slack" +import { sendDiscord } from "./discord" +import { BridgeAlert } from "./index.types" + +export * from "./index.types" + +/** + * Fire-and-forget fan-out of a Bridge alert to the configured destinations + * (ENG-361). Returns immediately; delivery is best-effort β€” each sender catches + * its own errors and no-ops when its credential/URL is unset, so it never throws + * or rejects into the caller (no need to await or handle it). + * + * Routing: + * - critical β†’ page on-call (PagerDuty) + inform (Slack/Mattermost, Discord) + * - warning β†’ inform (Slack/Mattermost, Discord) only + */ +export const alertBridge = (alert: BridgeAlert): void => { + const deliver = async () => { + const senders = [sendSlack(alert), sendDiscord(alert)] + if (alert.severity === "critical") { + senders.push(sendPagerDuty(alert)) + } + await Promise.allSettled(senders) + } + deliver().catch(() => undefined) +} diff --git a/src/services/alerts/index.types.ts b/src/services/alerts/index.types.ts new file mode 100644 index 000000000..1c6c5d651 --- /dev/null +++ b/src/services/alerts/index.types.ts @@ -0,0 +1,13 @@ +// Ops alerting for Bridge integration signals (ENG-361). + +export type AlertSeverity = "critical" | "warning" + +export type AlertSource = "bridge-webhook" | "bridge-api" | "ibex" | "erpnext-audit" + +export interface BridgeAlert { + source: AlertSource + severity: AlertSeverity + title: string + detail?: string + context?: Record +} diff --git a/src/services/alerts/pagerduty.ts b/src/services/alerts/pagerduty.ts new file mode 100644 index 000000000..da6acffcc --- /dev/null +++ b/src/services/alerts/pagerduty.ts @@ -0,0 +1,33 @@ +import { ALERT_PAGERDUTY_ROUTING_KEY } from "@config" +import { ErrorLevel } from "@domain/shared" +import { recordExceptionInCurrentSpan } from "@services/tracing" +import axios from "axios" + +import { BridgeAlert } from "./index.types" + +const PAGERDUTY_EVENTS_URL = "https://events.pagerduty.com/v2/enqueue" + +// PagerDuty Events API v2 β€” triggers a paging incident. "critical" and +// "warning" are both valid PD payload severities, so we pass them through. +export const sendPagerDuty = async (alert: BridgeAlert): Promise => { + if (!ALERT_PAGERDUTY_ROUTING_KEY) return + + try { + await axios.post( + PAGERDUTY_EVENTS_URL, + { + routing_key: ALERT_PAGERDUTY_ROUTING_KEY, + event_action: "trigger", + payload: { + summary: `[bridge:${alert.source}] ${alert.title}`, + severity: alert.severity, + source: "flash-bridge", + custom_details: { ...alert.context, detail: alert.detail }, + }, + }, + { timeout: 5000, headers: { "Content-Type": "application/json" } }, + ) + } catch (error) { + recordExceptionInCurrentSpan({ error, level: ErrorLevel.Warn }) + } +} diff --git a/src/services/alerts/slack.ts b/src/services/alerts/slack.ts new file mode 100644 index 000000000..a4c9456f4 --- /dev/null +++ b/src/services/alerts/slack.ts @@ -0,0 +1,31 @@ +import { ALERT_SLACK_WEBHOOK_URL } from "@config" +import { ErrorLevel } from "@domain/shared" +import { recordExceptionInCurrentSpan } from "@services/tracing" +import axios from "axios" + +import { BridgeAlert } from "./index.types" + +// Slack / Mattermost-compatible incoming webhook ({ text }). +export const sendSlack = async (alert: BridgeAlert): Promise => { + if (!ALERT_SLACK_WEBHOOK_URL) return + + const icon = alert.severity === "critical" ? ":rotating_light:" : ":warning:" + const lines = [ + `${icon} *Bridge alert* β€” ${alert.title}`, + `*source:* \`${alert.source}\` *severity:* \`${alert.severity}\``, + ] + if (alert.detail) lines.push(alert.detail) + if (alert.context) { + lines.push("```" + JSON.stringify(alert.context, null, 2) + "```") + } + + try { + await axios.post( + ALERT_SLACK_WEBHOOK_URL, + { text: lines.join("\n") }, + { timeout: 5000, headers: { "Content-Type": "application/json" } }, + ) + } catch (error) { + recordExceptionInCurrentSpan({ error, level: ErrorLevel.Warn }) + } +} From a33470d2394e818b399772c833b4c2afe7670e69 Mon Sep 17 00:00:00 2001 From: Dread Date: Fri, 5 Jun 2026 18:18:05 -0700 Subject: [PATCH 2/7] feat(alerts): wire Bridge alert sources to AlertService [ENG-361] Fire-and-forget alertBridge() at the Bridge failure points (alongside existing logging), all critical/page: - ERPNext audit-write failures (deposit + transfer completed/failed) - Bridge webhook processing exceptions (deposit + transfer catch) - Bridge API outage in client.request(): 5xx / timeout / network (4xx not alerted) IBEX-error source deferred: on-receive.ts is general LN/onchain receive handling, not the Bridge<->IBEX movement path; needs the exact call site (warning sev). Co-Authored-By: Claude Opus 4.8 (1M context) --- src/services/bridge/client.ts | 27 +++++++++++++++++++ .../bridge/webhook-server/routes/deposit.ts | 15 +++++++++++ .../bridge/webhook-server/routes/transfer.ts | 22 +++++++++++++++ 3 files changed, 64 insertions(+) diff --git a/src/services/bridge/client.ts b/src/services/bridge/client.ts index 363e535bf..2494b2a51 100644 --- a/src/services/bridge/client.ts +++ b/src/services/bridge/client.ts @@ -8,6 +8,7 @@ import crypto from "crypto" import { BridgeConfig } from "@config" import { BridgeCustomerId, BridgeTransferId, BridgeVirtualAccountId } from "@domain/primitives/bridge" +import { alertBridge } from "@services/alerts" import { BridgeTimeoutError } from "./errors" // ============ Error Handling ============ @@ -379,6 +380,16 @@ export class BridgeClient { const responseData = await response.json().catch(() => null) if (!response.ok) { + // Only 5xx indicates a Bridge-side outage; 4xx are normal API rejections. + if (response.status >= 500) { + alertBridge({ + source: "bridge-api", + severity: "critical", + title: `Bridge API ${response.status} on ${method} ${path}`, + detail: response.statusText, + context: { method, path, status: response.status }, + }) + } throw new BridgeApiError( `Bridge API error: ${response.status} ${response.statusText}`, response.status, @@ -389,8 +400,24 @@ export class BridgeClient { return responseData as T } catch (err) { if (err instanceof Error && err.name === "AbortError") { + alertBridge({ + source: "bridge-api", + severity: "critical", + title: `Bridge API timeout on ${method} ${path}`, + context: { method, path, timeoutMs }, + }) throw new BridgeTimeoutError() } + // Network/connectivity failures (5xx already alerted above). + if (!(err instanceof BridgeApiError)) { + alertBridge({ + source: "bridge-api", + severity: "critical", + title: `Bridge API request failed on ${method} ${path}`, + detail: err instanceof Error ? err.message : String(err), + context: { method, path }, + }) + } throw err } finally { clearTimeout(timeoutId) diff --git a/src/services/bridge/webhook-server/routes/deposit.ts b/src/services/bridge/webhook-server/routes/deposit.ts index 4fe09f432..06113999c 100644 --- a/src/services/bridge/webhook-server/routes/deposit.ts +++ b/src/services/bridge/webhook-server/routes/deposit.ts @@ -12,6 +12,7 @@ import { baseLogger } from "@services/logger" import { createBridgeDeposit } from "@services/mongoose/bridge-deposit-log" import { reconcileByTxHash } from "@services/bridge/reconciliation" import { writeBridgeDepositRequest } from "@services/frappe/BridgeTransferRequestWriter" +import { alertBridge } from "@services/alerts" export const depositHandler = async (req: Request, res: Response) => { const { event_id, event_object } = req.body @@ -94,6 +95,13 @@ export const depositHandler = async (req: Request, res: Response) => { { error: auditResult, event_id, id }, "Failed to persist Bridge deposit ERPNext audit row", ) + alertBridge({ + source: "erpnext-audit", + severity: "critical", + title: "Bridge deposit ERPNext audit write failed", + detail: auditResult.message, + context: { event_id, transfer_id: id }, + }) return res.status(500).json({ error: "Failed to persist ERPNext audit row" }) } @@ -109,6 +117,13 @@ export const depositHandler = async (req: Request, res: Response) => { return res.status(200).json({ status: "success" }) } catch (error) { baseLogger.error({ error, id, event_id }, "Error processing Bridge deposit webhook") + alertBridge({ + source: "bridge-webhook", + severity: "critical", + title: "Bridge deposit webhook processing error", + detail: error instanceof Error ? error.message : String(error), + context: { event_id, transfer_id: id }, + }) return res.status(500).json({ error: "Internal server error" }) } } diff --git a/src/services/bridge/webhook-server/routes/transfer.ts b/src/services/bridge/webhook-server/routes/transfer.ts index 88f00fad8..3ea8d3d2a 100644 --- a/src/services/bridge/webhook-server/routes/transfer.ts +++ b/src/services/bridge/webhook-server/routes/transfer.ts @@ -13,6 +13,7 @@ import { writeBridgeCashoutCompleted, writeBridgeCashoutFailed, } from "@services/frappe/BridgeTransferRequestWriter" +import { alertBridge } from "@services/alerts" const TERMINAL_FAILURE_STATES = new Set([ "undeliverable", @@ -121,6 +122,13 @@ export const transferHandler = async (req: Request, res: Response) => { { transfer_id, error: auditResult }, "Failed to persist Bridge transfer ERPNext audit row", ) + alertBridge({ + source: "erpnext-audit", + severity: "critical", + title: "Bridge transfer ERPNext audit write failed", + detail: auditResult.message, + context: { transfer_id, event }, + }) return res.status(500).json({ error: "Failed to persist ERPNext audit row" }) } @@ -196,6 +204,13 @@ export const transferHandler = async (req: Request, res: Response) => { { transfer_id, error: auditResult }, "Failed to persist Bridge transfer failure ERPNext audit row", ) + alertBridge({ + source: "erpnext-audit", + severity: "critical", + title: "Bridge transfer-failure ERPNext audit write failed", + detail: auditResult.message, + context: { transfer_id, event }, + }) return res.status(500).json({ error: "Failed to persist ERPNext audit row" }) } @@ -216,6 +231,13 @@ export const transferHandler = async (req: Request, res: Response) => { return res.status(200).json({ status: "success" }) } catch (error) { baseLogger.error({ error, transfer_id }, "Error processing Bridge transfer webhook") + alertBridge({ + source: "bridge-webhook", + severity: "critical", + title: "Bridge transfer webhook processing error", + detail: error instanceof Error ? error.message : String(error), + context: { transfer_id, event }, + }) return res.status(500).json({ error: "Internal server error" }) } } From 1030a3b341c6c2c71af20232227562645f025fc1 Mon Sep 17 00:00:00 2001 From: Dread Date: Fri, 5 Jun 2026 18:31:48 -0700 Subject: [PATCH 3/7] docs(alerts): add Bridge alerting setup guide [ENG-361] Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/bridge-integration/ALERTING.md | 69 +++++++++++++++++++++++++++++ 1 file changed, 69 insertions(+) create mode 100644 docs/bridge-integration/ALERTING.md diff --git a/docs/bridge-integration/ALERTING.md b/docs/bridge-integration/ALERTING.md new file mode 100644 index 000000000..712e8f18a --- /dev/null +++ b/docs/bridge-integration/ALERTING.md @@ -0,0 +1,69 @@ +# Bridge Alerting (ENG-361) + +Operational alerting for the Bridge integration. When a Bridge signal fails +(webhook processing, ERPNext audit write, or a Bridge API outage), the +`AlertService` (`src/services/alerts`) fans the alert out to the configured +destinations. + +## Routing + +| Severity | PagerDuty (page) | Slack / Mattermost (inform) | Discord (inform) | +| ------------ | :--------------: | :-------------------------: | :--------------: | +| **critical** | βœ… | βœ… | βœ… | +| **warning** | β€” | βœ… | βœ… | + +Delivery is best-effort and fire-and-forget β€” a failing or unconfigured +destination never blocks or fails the webhook/request path. **A destination +with no configured credential is silently skipped**, so channels can be enabled +incrementally. + +## Alert sources + +| Source | Severity | Where | +| ------------------------------------------------- | -------- | ----------------------------------------------------------- | +| ERPNext audit-write failure (deposit + transfer) | critical | `services/bridge/webhook-server/routes/{deposit,transfer}.ts` | +| Bridge webhook processing exception | critical | same routes (catch block) | +| Bridge API outage β€” 5xx / timeout / network | critical | `services/bridge/client.ts` | +| IBEX error on a Bridge↔IBEX movement | warning | _follow-up β€” not yet wired_ | + +`4xx` responses from Bridge are normal API rejections and are **not** alerted. + +## Configuration + +Three optional env vars, each gating one destination: + +| Env var | Destination | Value | +| ----------------------------- | ------------------- | ------------------------------------------- | +| `ALERT_PAGERDUTY_ROUTING_KEY` | PagerDuty | Events API v2 **integration / routing key** | +| `ALERT_SLACK_WEBHOOK_URL` | Slack or Mattermost | Incoming-webhook URL | +| `ALERT_DISCORD_WEBHOOK_URL` | Discord | Channel webhook URL | + +### How to get each value + +**PagerDuty** β€” `ALERT_PAGERDUTY_ROUTING_KEY` +1. PagerDuty β†’ **Services** β†’ pick (or create) the service that should page for Bridge. +2. **Integrations** β†’ **Add integration** β†’ **Events API v2**. +3. Copy the **Integration Key** β€” that is the routing key. + +**Slack** β€” `ALERT_SLACK_WEBHOOK_URL` +1. Create/choose a Slack app β†’ **Incoming Webhooks** β†’ **Activate**. +2. **Add New Webhook to Workspace** β†’ choose the target channel. +3. Copy the URL (`https://hooks.slack.com/services/...`). + _Mattermost works too_ β€” it accepts the same `{ text }` payload; use its incoming-webhook URL. + +**Discord** β€” `ALERT_DISCORD_WEBHOOK_URL` +1. Discord β†’ target channel β†’ **Edit Channel** β†’ **Integrations** β†’ **Webhooks**. +2. **New Webhook** β†’ name it β†’ **Copy Webhook URL**. + +### Where to set them + +- **Local dev:** add to `.env` (and `.env.ci` for CI). +- **Staging / production:** set as environment variables / secrets in the deployment β€” the same place `MATTERMOST_WEBHOOK_URL` is configured. Treat all three as **secrets**. + +> If none are set, alerting is a no-op (no errors, no delivery) β€” useful until the channels are provisioned. + +## Verifying in staging (ENG-361 acceptance) + +1. Set at least `ALERT_PAGERDUTY_ROUTING_KEY` and `ALERT_SLACK_WEBHOOK_URL` in staging. +2. Simulate a Bridge webhook failure (e.g. force an ERPNext audit-write error, or replay a malformed transfer webhook). +3. Confirm on-call is paged via PagerDuty **and** a message posts to Slack within ~1 minute. From 49d32d2896c117b775213cea8b5b5ffdff5078df Mon Sep 17 00:00:00 2001 From: Vandana Date: Mon, 8 Jun 2026 07:37:57 -0700 Subject: [PATCH 4/7] chore(alerts): clean bridge alert lint and deposit test --- src/services/bridge/client.ts | 42 +++++++++++-------- .../bridge/webhook-server/routes/deposit.ts | 4 +- .../bridge/webhook-server/deposit.spec.ts | 7 +++- 3 files changed, 34 insertions(+), 19 deletions(-) diff --git a/src/services/bridge/client.ts b/src/services/bridge/client.ts index 2494b2a51..41ed9ba5a 100644 --- a/src/services/bridge/client.ts +++ b/src/services/bridge/client.ts @@ -7,8 +7,13 @@ import crypto from "crypto" import { BridgeConfig } from "@config" -import { BridgeCustomerId, BridgeTransferId, BridgeVirtualAccountId } from "@domain/primitives/bridge" +import { + BridgeCustomerId, + BridgeTransferId, + BridgeVirtualAccountId, +} from "@domain/primitives/bridge" import { alertBridge } from "@services/alerts" + import { BridgeTimeoutError } from "./errors" // ============ Error Handling ============ @@ -68,15 +73,15 @@ export interface Customer { id: string type: "individual" | "business" status?: - | "active" - | "awaiting_questionnaire" - | "rejected" - | "paused" - | "under_review" - | "offboarded" - | "awaiting_ubo" - | "incomplete" - | "not_started" + | "active" + | "awaiting_questionnaire" + | "rejected" + | "paused" + | "under_review" + | "offboarded" + | "awaiting_ubo" + | "incomplete" + | "not_started" has_accepted_terms_of_service?: string created_at: string updated_at: string @@ -466,7 +471,9 @@ export class BridgeClient { } async getVirtualAccount( - customerId: BridgeCustomerId, virtualAccountId: BridgeVirtualAccountId, idempotencyKey?: string, + customerId: BridgeCustomerId, + virtualAccountId: BridgeVirtualAccountId, + idempotencyKey?: string, ): Promise { return this.request( "GET", @@ -476,16 +483,15 @@ export class BridgeClient { ) } - - async getVirtualAccountByCustomerId(customerId: BridgeCustomerId): Promise { + async getVirtualAccountByCustomerId( + customerId: BridgeCustomerId, + ): Promise { const response = await this.request<{ data: VirtualAccount[] }>( "GET", `/customers/${customerId}/virtual_accounts`, ) return response.data as VirtualAccount[] - - } // ============ External Accounts ============ @@ -594,8 +600,10 @@ export async function* listAllEvents( const startMs = params?.start ? new Date(params.start).getTime() : -Infinity const endMs = params?.end ? new Date(params.end).getTime() : Infinity - // Strip start/end β€” Bridge /webhook_events only supports cursor params; filter locally. - const { start: _s, end: _e, ...apiParams } = params ?? {} + // Strip start/end: Bridge /webhook_events only supports cursor params; filter locally. + const apiParams = { ...(params ?? {}) } + delete apiParams.start + delete apiParams.end let cursor: string | undefined do { diff --git a/src/services/bridge/webhook-server/routes/deposit.ts b/src/services/bridge/webhook-server/routes/deposit.ts index 06113999c..4eb0ae1ca 100644 --- a/src/services/bridge/webhook-server/routes/deposit.ts +++ b/src/services/bridge/webhook-server/routes/deposit.ts @@ -108,7 +108,9 @@ export const depositHandler = async (req: Request, res: Response) => { // Idempotency: mark processed only after local and ERPNext writes succeed, so // provider retries can recover audit gaps after transient ERPNext failures. const auditLockKey = `bridge-deposit:${event_id}` - const auditLockResult = await LockService().lockIdempotencyKey(auditLockKey as IdempotencyKey) + const auditLockResult = await LockService().lockIdempotencyKey( + auditLockKey as IdempotencyKey, + ) if (auditLockResult instanceof Error) { baseLogger.info({ event_id, id, state }, "Duplicate Bridge deposit webhook") return res.status(200).json({ status: "already_processed" }) diff --git a/test/flash/unit/services/bridge/webhook-server/deposit.spec.ts b/test/flash/unit/services/bridge/webhook-server/deposit.spec.ts index 44788f076..54094f81d 100644 --- a/test/flash/unit/services/bridge/webhook-server/deposit.spec.ts +++ b/test/flash/unit/services/bridge/webhook-server/deposit.spec.ts @@ -115,7 +115,11 @@ describe("depositHandler β€” invalid payload", () => { describe("depositHandler β€” idempotency", () => { it("returns already_processed after idempotent local and audit writes on duplicate delivery", async () => { - mockLockService(false) + const lockFn = jest + .fn() + .mockResolvedValueOnce({}) + .mockResolvedValueOnce(new Error("already locked")) + ;(LockService as jest.Mock).mockReturnValue({ lockIdempotencyKey: lockFn }) ;(DepositLog.createBridgeDeposit as jest.Mock).mockResolvedValue({ id: "log-dup" }) const res = makeRes() @@ -124,6 +128,7 @@ describe("depositHandler β€” idempotency", () => { expect(res.json as jest.Mock).toHaveBeenCalledWith({ status: "already_processed" }) expect(DepositLog.createBridgeDeposit).toHaveBeenCalledTimes(1) expect(writeBridgeDepositRequest).toHaveBeenCalledTimes(1) + expect(lockFn).toHaveBeenCalledTimes(2) }) it("locks on the event id after writing the audit row", async () => { From ba904f2eff70f6babfde3923b43090fa58a4decf Mon Sep 17 00:00:00 2001 From: heyolaniran Date: Tue, 9 Jun 2026 12:31:38 +0100 Subject: [PATCH 5/7] feat(alerts): dedupe Bridge alerts and wire IBEX movement warnings [ENG-361] Group PagerDuty triggers with dedup_key, suppress duplicate Slack/Discord informs, and alert IBEX crypto-receive and reconciliation failures as warnings. --- docs/bridge-integration/ALERTING.md | 24 +++- src/services/alerts/dedup-key.ts | 87 ++++++++++++++ src/services/alerts/ibex-bridge-movement.ts | 88 ++++++++++++++ src/services/alerts/index.ts | 25 +++- src/services/alerts/index.types.ts | 1 + src/services/alerts/inform-dedup.ts | 35 ++++++ src/services/alerts/pagerduty.ts | 1 + src/services/bridge/client.ts | 5 +- src/services/bridge/reconciliation.ts | 61 +++++++++- .../bridge/webhook-server/routes/deposit.ts | 15 ++- .../bridge/webhook-server/routes/transfer.ts | 5 +- .../webhook-server/routes/crypto-receive.ts | 64 +++++++++- .../unit/services/alerts/dedup-key.spec.ts | 109 ++++++++++++++++++ .../alerts/ibex-bridge-movement.spec.ts | 59 ++++++++++ test/flash/unit/services/alerts/index.spec.ts | 76 ++++++++++++ .../unit/services/alerts/inform-dedup.spec.ts | 34 ++++++ .../unit/services/alerts/pagerduty.spec.ts | 44 +++++++ .../services/bridge/reconciliation.spec.ts | 101 ++++++++++------ .../routes/crypto-receive.spec.ts | 17 +++ 19 files changed, 796 insertions(+), 55 deletions(-) create mode 100644 src/services/alerts/dedup-key.ts create mode 100644 src/services/alerts/ibex-bridge-movement.ts create mode 100644 src/services/alerts/inform-dedup.ts create mode 100644 test/flash/unit/services/alerts/dedup-key.spec.ts create mode 100644 test/flash/unit/services/alerts/ibex-bridge-movement.spec.ts create mode 100644 test/flash/unit/services/alerts/index.spec.ts create mode 100644 test/flash/unit/services/alerts/inform-dedup.spec.ts create mode 100644 test/flash/unit/services/alerts/pagerduty.spec.ts diff --git a/docs/bridge-integration/ALERTING.md b/docs/bridge-integration/ALERTING.md index 712e8f18a..270621a8b 100644 --- a/docs/bridge-integration/ALERTING.md +++ b/docs/bridge-integration/ALERTING.md @@ -17,6 +17,27 @@ destination never blocks or fails the webhook/request path. **A destination with no configured credential is silently skipped**, so channels can be enabled incrementally. +### Deduplication + +Alerts carry a stable `dedupKey` so repeated failures do not spam on-call or chat: + +| Destination | Behavior | +| ----------- | -------- | +| **PagerDuty** | Events API v2 `dedup_key` groups triggers into one incident | +| **Slack / Discord** | First message per `dedupKey` within TTL; duplicates are skipped | + +Key classes (see `src/services/alerts/dedup-key.ts`): + +- `bridge-api:5xx` / `bridge-api:timeout` / `bridge-api:network` β€” coarse outage keys (30 min inform TTL) +- `erpnext-audit:deposit:{transfer_id}` β€” per deposit audit failure (1 h inform TTL) +- `erpnext-audit:transfer-complete:{transfer_id}` / `transfer-failed:{transfer_id}` β€” per transfer audit failure +- `bridge-webhook:deposit:{event_id}` / `bridge-webhook:transfer:{transfer_id}:{event}` β€” per webhook processing error +- `ibex:crypto-receive:{tx_hash}` β€” per IBEX crypto receive webhook failure (1 h inform TTL) +- `ibex:reconcile:bridge-without-ibex:{tx_hash}` / `ibex:reconcile:ibex-without-bridge:{tx_hash}` β€” per reconciliation orphan +- `ibex:reconcile:failed:{tx_hash}` β€” reconciliation handler threw + +Inform dedup is in-process per pod; PagerDuty dedup is global to the service integration. + ## Alert sources | Source | Severity | Where | @@ -24,7 +45,8 @@ incrementally. | ERPNext audit-write failure (deposit + transfer) | critical | `services/bridge/webhook-server/routes/{deposit,transfer}.ts` | | Bridge webhook processing exception | critical | same routes (catch block) | | Bridge API outage β€” 5xx / timeout / network | critical | `services/bridge/client.ts` | -| IBEX error on a Bridge↔IBEX movement | warning | _follow-up β€” not yet wired_ | +| IBEX crypto receive webhook failure | warning | `services/ibex/webhook-server/routes/crypto-receive.ts` | +| Bridge↔IBEX reconciliation orphan / failure | warning | `services/bridge/reconciliation.ts`, deposit/crypto catch | `4xx` responses from Bridge are normal API rejections and are **not** alerted. diff --git a/src/services/alerts/dedup-key.ts b/src/services/alerts/dedup-key.ts new file mode 100644 index 000000000..0e2ed0837 --- /dev/null +++ b/src/services/alerts/dedup-key.ts @@ -0,0 +1,87 @@ +import { BridgeAlert } from "./index.types" + +export const PAGERDUTY_DEDUP_KEY_MAX = 255 + +const OUTAGE_TTL_MS = 30 * 60 * 1000 +const DEFAULT_TTL_MS = 60 * 60 * 1000 + +/** TTL for Slack/Discord first-alert suppression per dedup key class. */ +export const informDedupTtlMs = (dedupKey: string): number => + dedupKey.startsWith("bridge-api") ? OUTAGE_TTL_MS : DEFAULT_TTL_MS + +export const generateDedupKey = { + bridgeApi5xx: () => "bridge-api:5xx", + bridgeApiTimeout: () => "bridge-api:timeout", + bridgeApiNetwork: () => "bridge-api:network", + erpnextDepositAudit: (transferId: string) => `erpnext-audit:deposit:${transferId}`, + erpnextTransferCompletedAudit: (transferId: string) => + `erpnext-audit:transfer-complete:${transferId}`, + erpnextTransferFailedAudit: (transferId: string) => + `erpnext-audit:transfer-failed:${transferId}`, + bridgeWebhookDeposit: (eventId: string) => `bridge-webhook:deposit:${eventId}`, + bridgeWebhookTransfer: (transferId: string, event: string) => + `bridge-webhook:transfer:${transferId}:${event}`, + ibexCryptoReceive: (txHash: string) => `ibex:crypto-receive:${txHash.toLowerCase()}`, + ibexReconcileBridgeWithoutIbex: (txHash: string) => + `ibex:reconcile:bridge-without-ibex:${txHash.toLowerCase()}`, + ibexReconcileBridgeWithoutIbexTransfer: (transferId: string) => + `ibex:reconcile:bridge-without-ibex:transfer:${transferId}`, + ibexReconcileIbexWithoutBridge: (txHash: string) => + `ibex:reconcile:ibex-without-bridge:${txHash.toLowerCase()}`, + ibexReconcileFailed: (txHash: string) => `ibex:reconcile:failed:${txHash.toLowerCase()}`, +} + +const truncateDedupKey = (key: string): string => + key.length <= PAGERDUTY_DEDUP_KEY_MAX ? key : key.slice(0, PAGERDUTY_DEDUP_KEY_MAX) + +/** Stable PagerDuty / inform dedup key; prefers explicit alert.dedupKey. */ +export const resolveDedupKey = (alert: BridgeAlert): string => { + if (alert.dedupKey) return truncateDedupKey(alert.dedupKey) + + const ctx = alert.context ?? {} + + switch (alert.source) { + case "bridge-api": + if (alert.title.includes("timeout")) return generateDedupKey.bridgeApiTimeout() + if (alert.title.includes("request failed")) return generateDedupKey.bridgeApiNetwork() + return generateDedupKey.bridgeApi5xx() + case "erpnext-audit": { + const transferId = String(ctx.transfer_id ?? "unknown") + if (alert.title.includes("deposit")) return generateDedupKey.erpnextDepositAudit(transferId) + if (alert.title.includes("failure")) { + return generateDedupKey.erpnextTransferFailedAudit(transferId) + } + return generateDedupKey.erpnextTransferCompletedAudit(transferId) + } + case "bridge-webhook": { + const eventId = String(ctx.event_id ?? ctx.transfer_id ?? "unknown") + const transferId = String(ctx.transfer_id ?? "unknown") + const event = String(ctx.event ?? "unknown") + if (alert.title.includes("deposit")) { + return generateDedupKey.bridgeWebhookDeposit(eventId) + } + return generateDedupKey.bridgeWebhookTransfer(transferId, event) + } + case "ibex": { + const txHash = String(ctx.tx_hash ?? ctx.txHash ?? "unknown") + const orphanType = String(ctx.orphan_type ?? "") + if (orphanType === "ibex_without_bridge") { + return generateDedupKey.ibexReconcileIbexWithoutBridge(txHash) + } + if (orphanType === "bridge_without_ibex") { + if (txHash !== "unknown") { + return generateDedupKey.ibexReconcileBridgeWithoutIbex(txHash) + } + return generateDedupKey.ibexReconcileBridgeWithoutIbexTransfer( + String(ctx.transfer_id ?? "unknown"), + ) + } + if (alert.title.includes("reconciliation failed")) { + return generateDedupKey.ibexReconcileFailed(txHash) + } + return generateDedupKey.ibexCryptoReceive(txHash) + } + default: + return truncateDedupKey(`bridge:${alert.source}:${alert.title}`) + } +} diff --git a/src/services/alerts/ibex-bridge-movement.ts b/src/services/alerts/ibex-bridge-movement.ts new file mode 100644 index 000000000..9b254e63d --- /dev/null +++ b/src/services/alerts/ibex-bridge-movement.ts @@ -0,0 +1,88 @@ +import { alertBridge } from "./index" +import { generateDedupKey } from "./dedup-key" + +type IbexMovementAlert = { + title: string + detail?: string + context?: Record +} + +const alertIbexMovement = (dedupKey: string, alert: IbexMovementAlert): void => { + alertBridge({ + dedupKey, + source: "ibex", + severity: "warning", + ...alert, + }) +} + +export const alertIbexCryptoReceiveFailure = ({ + txHash, + code, + title, + detail, + context, +}: { + txHash: string + code: string + title: string + detail?: string + context?: Record +}): void => { + alertIbexMovement(generateDedupKey.ibexCryptoReceive(txHash), { + title, + detail, + context: { tx_hash: txHash, code, ...context }, + }) +} + +export const alertIbexReconciliationOrphan = ({ + orphanType, + txHash, + transferId, + reason, + context, +}: { + orphanType: "bridge_without_ibex" | "ibex_without_bridge" + txHash?: string + transferId?: string + reason: string + context?: Record +}): void => { + const dedupKey = + orphanType === "ibex_without_bridge" && txHash + ? generateDedupKey.ibexReconcileIbexWithoutBridge(txHash) + : txHash + ? generateDedupKey.ibexReconcileBridgeWithoutIbex(txHash) + : generateDedupKey.ibexReconcileBridgeWithoutIbexTransfer(transferId ?? "unknown") + + const title = + orphanType === "ibex_without_bridge" + ? "IBEX crypto receive without matching Bridge deposit" + : "Bridge deposit without matching IBEX crypto receive" + + alertIbexMovement(dedupKey, { + title, + detail: reason, + context: { + orphan_type: orphanType, + tx_hash: txHash, + transfer_id: transferId, + ...context, + }, + }) +} + +export const alertIbexReconciliationFailed = ({ + txHash, + detail, +}: { + txHash: string + detail: string +}): void => { + alertIbexMovement(generateDedupKey.ibexReconcileFailed(txHash), { + title: "Bridge↔IBEX reconciliation failed", + detail, + context: { tx_hash: txHash }, + }) +} diff --git a/src/services/alerts/index.ts b/src/services/alerts/index.ts index 733dc7e7f..0091505ef 100644 --- a/src/services/alerts/index.ts +++ b/src/services/alerts/index.ts @@ -1,9 +1,12 @@ import { sendPagerDuty } from "./pagerduty" import { sendSlack } from "./slack" import { sendDiscord } from "./discord" +import { resolveDedupKey } from "./dedup-key" +import { claimInformSlot } from "./inform-dedup" import { BridgeAlert } from "./index.types" export * from "./index.types" +export { generateDedupKey } from "./dedup-key" /** * Fire-and-forget fan-out of a Bridge alert to the configured destinations @@ -14,14 +17,30 @@ export * from "./index.types" * Routing: * - critical β†’ page on-call (PagerDuty) + inform (Slack/Mattermost, Discord) * - warning β†’ inform (Slack/Mattermost, Discord) only + * + * Dedup: + * - PagerDuty: Events API v2 dedup_key groups triggers into one incident. + * - Slack / Discord: first alert per dedup key within TTL only. */ export const alertBridge = (alert: BridgeAlert): void => { + const dedupKey = resolveDedupKey(alert) + const alertWithKey: BridgeAlert = { ...alert, dedupKey } + const deliver = async () => { - const senders = [sendSlack(alert), sendDiscord(alert)] + const senders: Promise[] = [] + + if (claimInformSlot(dedupKey)) { + senders.push(sendSlack(alertWithKey), sendDiscord(alertWithKey)) + } + if (alert.severity === "critical") { - senders.push(sendPagerDuty(alert)) + senders.push(sendPagerDuty(alertWithKey)) + } + + if (senders.length > 0) { + await Promise.allSettled(senders) } - await Promise.allSettled(senders) } + deliver().catch(() => undefined) } diff --git a/src/services/alerts/index.types.ts b/src/services/alerts/index.types.ts index 1c6c5d651..09d153f08 100644 --- a/src/services/alerts/index.types.ts +++ b/src/services/alerts/index.types.ts @@ -5,6 +5,7 @@ export type AlertSeverity = "critical" | "warning" export type AlertSource = "bridge-webhook" | "bridge-api" | "ibex" | "erpnext-audit" export interface BridgeAlert { + dedupKey?: string source: AlertSource severity: AlertSeverity title: string diff --git a/src/services/alerts/inform-dedup.ts b/src/services/alerts/inform-dedup.ts new file mode 100644 index 000000000..7e0bc7884 --- /dev/null +++ b/src/services/alerts/inform-dedup.ts @@ -0,0 +1,35 @@ +import { informDedupTtlMs } from "./dedup-key" + +const seenAt = new Map() + +/** + * Returns true when Slack/Discord should fire for this dedup key (first within TTL). + * Subsequent duplicates within the TTL are suppressed. + */ +export const claimInformSlot = (dedupKey: string, nowMs = Date.now()): boolean => { + const ttlMs = informDedupTtlMs(dedupKey) + const lastSentAt = seenAt.get(dedupKey) + + if (lastSentAt !== undefined && nowMs - lastSentAt < ttlMs) { + return false + } + + seenAt.set(dedupKey, nowMs) + pruneExpired(nowMs) + return true +} + +const pruneExpired = (nowMs: number): void => { + if (seenAt.size < 500) return + + for (const [key, sentAt] of seenAt) { + if (nowMs - sentAt >= informDedupTtlMs(key)) { + seenAt.delete(key) + } + } +} + +/** Test helper β€” clears the in-process inform dedup cache. */ +export const resetInformDedup = (): void => { + seenAt.clear() +} diff --git a/src/services/alerts/pagerduty.ts b/src/services/alerts/pagerduty.ts index da6acffcc..4f14827d9 100644 --- a/src/services/alerts/pagerduty.ts +++ b/src/services/alerts/pagerduty.ts @@ -18,6 +18,7 @@ export const sendPagerDuty = async (alert: BridgeAlert): Promise => { { routing_key: ALERT_PAGERDUTY_ROUTING_KEY, event_action: "trigger", + dedup_key: alert.dedupKey, payload: { summary: `[bridge:${alert.source}] ${alert.title}`, severity: alert.severity, diff --git a/src/services/bridge/client.ts b/src/services/bridge/client.ts index 41ed9ba5a..d6a8c9310 100644 --- a/src/services/bridge/client.ts +++ b/src/services/bridge/client.ts @@ -12,7 +12,7 @@ import { BridgeTransferId, BridgeVirtualAccountId, } from "@domain/primitives/bridge" -import { alertBridge } from "@services/alerts" +import { alertBridge, generateDedupKey } from "@services/alerts" import { BridgeTimeoutError } from "./errors" @@ -388,6 +388,7 @@ export class BridgeClient { // Only 5xx indicates a Bridge-side outage; 4xx are normal API rejections. if (response.status >= 500) { alertBridge({ + dedupKey: generateDedupKey.bridgeApi5xx(), source: "bridge-api", severity: "critical", title: `Bridge API ${response.status} on ${method} ${path}`, @@ -406,6 +407,7 @@ export class BridgeClient { } catch (err) { if (err instanceof Error && err.name === "AbortError") { alertBridge({ + dedupKey: generateDedupKey.bridgeApiTimeout(), source: "bridge-api", severity: "critical", title: `Bridge API timeout on ${method} ${path}`, @@ -416,6 +418,7 @@ export class BridgeClient { // Network/connectivity failures (5xx already alerted above). if (!(err instanceof BridgeApiError)) { alertBridge({ + dedupKey: generateDedupKey.bridgeApiNetwork(), source: "bridge-api", severity: "critical", title: `Bridge API request failed on ${method} ${path}`, diff --git a/src/services/bridge/reconciliation.ts b/src/services/bridge/reconciliation.ts index 851ebce90..26ef10ada 100644 --- a/src/services/bridge/reconciliation.ts +++ b/src/services/bridge/reconciliation.ts @@ -1,3 +1,4 @@ +import { alertIbexReconciliationOrphan } from "@services/alerts/ibex-bridge-movement" import { baseLogger } from "@services/logger" import { findIbexCryptoReceivesSince } from "@services/mongoose/ibex-crypto-receive-log" import { @@ -78,6 +79,7 @@ export const reconcileBridgeAndIbexDeposits = async ({ for (const deposit of bridgeDeposits) { if (!deposit.destinationTxHash) { bridgeWithoutIbex++ + const reason = "Bridge payment_processed has no destinationTxHash" await upsertBridgeReconciliationOrphan({ orphanKey: toOrphanKey("bridge-no-tx", deposit.transferId), orphanType: "bridge_without_ibex", @@ -87,13 +89,24 @@ export const reconcileBridgeAndIbexDeposits = async ({ amount: deposit.amount, currency: deposit.currency, triageContext: { - reason: "Bridge payment_processed has no destinationTxHash", + reason, windowStart: since.toISOString(), windowEnd: now.toISOString(), depositState: deposit.state, createdAt: deposit.createdAt.toISOString(), }, }) + alertIbexReconciliationOrphan({ + orphanType: "bridge_without_ibex", + transferId: deposit.transferId, + reason, + context: { + bridge_event_id: deposit.eventId, + customer_id: deposit.customerId, + amount: deposit.amount, + currency: deposit.currency, + }, + }) continue } @@ -101,6 +114,8 @@ export const reconcileBridgeAndIbexDeposits = async ({ if (matchedIbex) continue bridgeWithoutIbex++ + const reason = + "No IBEX crypto.receive found for Bridge destinationTxHash within window" await upsertBridgeReconciliationOrphan({ orphanKey: toOrphanKey("bridge", deposit.destinationTxHash), orphanType: "bridge_without_ibex", @@ -111,14 +126,25 @@ export const reconcileBridgeAndIbexDeposits = async ({ amount: deposit.amount, currency: deposit.currency, triageContext: { - reason: - "No IBEX crypto.receive found for Bridge destinationTxHash within window", + reason, windowStart: since.toISOString(), windowEnd: now.toISOString(), depositState: deposit.state, createdAt: deposit.createdAt.toISOString(), }, }) + alertIbexReconciliationOrphan({ + orphanType: "bridge_without_ibex", + txHash: deposit.destinationTxHash, + transferId: deposit.transferId, + reason, + context: { + bridge_event_id: deposit.eventId, + customer_id: deposit.customerId, + amount: deposit.amount, + currency: deposit.currency, + }, + }) } for (const receive of ibexReceives) { @@ -126,6 +152,8 @@ export const reconcileBridgeAndIbexDeposits = async ({ if (matchedBridge) continue ibexWithoutBridge++ + const reason = + "No Bridge deposit payment_processed found for IBEX tx hash within window" await upsertBridgeReconciliationOrphan({ orphanKey: toOrphanKey("ibex", receive.txHash), orphanType: "ibex_without_bridge", @@ -133,8 +161,7 @@ export const reconcileBridgeAndIbexDeposits = async ({ amount: receive.amount, currency: receive.currency, triageContext: { - reason: - "No Bridge deposit payment_processed found for IBEX tx hash within window", + reason, windowStart: since.toISOString(), windowEnd: now.toISOString(), address: receive.address, @@ -143,6 +170,18 @@ export const reconcileBridgeAndIbexDeposits = async ({ receivedAt: receive.receivedAt.toISOString(), }, }) + alertIbexReconciliationOrphan({ + orphanType: "ibex_without_bridge", + txHash: receive.txHash, + reason, + context: { + amount: receive.amount, + currency: receive.currency, + address: receive.address, + network: receive.network, + account_id: receive.accountId, + }, + }) } const summary = { @@ -263,6 +302,18 @@ export const reconcileByTxHash = async ({ triageContext, }) + alertIbexReconciliationOrphan({ + orphanType, + txHash: normalizedHash, + transferId, + reason: String(triageContext.reason), + context: { + customer_id: customerId, + amount, + currency, + }, + }) + const event: ReconcileByTxHashResult = { txHash: normalizedHash, status: "unmatched", diff --git a/src/services/bridge/webhook-server/routes/deposit.ts b/src/services/bridge/webhook-server/routes/deposit.ts index 4eb0ae1ca..33e163147 100644 --- a/src/services/bridge/webhook-server/routes/deposit.ts +++ b/src/services/bridge/webhook-server/routes/deposit.ts @@ -12,7 +12,8 @@ import { baseLogger } from "@services/logger" import { createBridgeDeposit } from "@services/mongoose/bridge-deposit-log" import { reconcileByTxHash } from "@services/bridge/reconciliation" import { writeBridgeDepositRequest } from "@services/frappe/BridgeTransferRequestWriter" -import { alertBridge } from "@services/alerts" +import { alertBridge, generateDedupKey } from "@services/alerts" +import { alertIbexReconciliationFailed } from "@services/alerts/ibex-bridge-movement" export const depositHandler = async (req: Request, res: Response) => { const { event_id, event_object } = req.body @@ -80,9 +81,13 @@ export const depositHandler = async (req: Request, res: Response) => { } if (state === "payment_processed" && receipt?.destination_tx_hash) { - reconcileByTxHash({ txHash: receipt.destination_tx_hash }).catch((err) => - baseLogger.error({ err, event_id, id }, "Real-time reconciliation failed"), - ) + reconcileByTxHash({ txHash: receipt.destination_tx_hash }).catch((err) => { + baseLogger.error({ err, event_id, id }, "Real-time reconciliation failed") + alertIbexReconciliationFailed({ + txHash: receipt.destination_tx_hash, + detail: err instanceof Error ? err.message : String(err), + }) + }) } const auditResult = await writeBridgeDepositRequest({ @@ -96,6 +101,7 @@ export const depositHandler = async (req: Request, res: Response) => { "Failed to persist Bridge deposit ERPNext audit row", ) alertBridge({ + dedupKey: generateDedupKey.erpnextDepositAudit(id), source: "erpnext-audit", severity: "critical", title: "Bridge deposit ERPNext audit write failed", @@ -120,6 +126,7 @@ export const depositHandler = async (req: Request, res: Response) => { } catch (error) { baseLogger.error({ error, id, event_id }, "Error processing Bridge deposit webhook") alertBridge({ + dedupKey: generateDedupKey.bridgeWebhookDeposit(event_id), source: "bridge-webhook", severity: "critical", title: "Bridge deposit webhook processing error", diff --git a/src/services/bridge/webhook-server/routes/transfer.ts b/src/services/bridge/webhook-server/routes/transfer.ts index 3ea8d3d2a..fd6410e87 100644 --- a/src/services/bridge/webhook-server/routes/transfer.ts +++ b/src/services/bridge/webhook-server/routes/transfer.ts @@ -13,7 +13,7 @@ import { writeBridgeCashoutCompleted, writeBridgeCashoutFailed, } from "@services/frappe/BridgeTransferRequestWriter" -import { alertBridge } from "@services/alerts" +import { alertBridge, generateDedupKey } from "@services/alerts" const TERMINAL_FAILURE_STATES = new Set([ "undeliverable", @@ -123,6 +123,7 @@ export const transferHandler = async (req: Request, res: Response) => { "Failed to persist Bridge transfer ERPNext audit row", ) alertBridge({ + dedupKey: generateDedupKey.erpnextTransferCompletedAudit(transfer_id), source: "erpnext-audit", severity: "critical", title: "Bridge transfer ERPNext audit write failed", @@ -205,6 +206,7 @@ export const transferHandler = async (req: Request, res: Response) => { "Failed to persist Bridge transfer failure ERPNext audit row", ) alertBridge({ + dedupKey: generateDedupKey.erpnextTransferFailedAudit(transfer_id), source: "erpnext-audit", severity: "critical", title: "Bridge transfer-failure ERPNext audit write failed", @@ -232,6 +234,7 @@ export const transferHandler = async (req: Request, res: Response) => { } catch (error) { baseLogger.error({ error, transfer_id }, "Error processing Bridge transfer webhook") alertBridge({ + dedupKey: generateDedupKey.bridgeWebhookTransfer(transfer_id, event), source: "bridge-webhook", severity: "critical", title: "Bridge transfer webhook processing error", diff --git a/src/services/ibex/webhook-server/routes/crypto-receive.ts b/src/services/ibex/webhook-server/routes/crypto-receive.ts index 7a45d129c..939ffd0dc 100644 --- a/src/services/ibex/webhook-server/routes/crypto-receive.ts +++ b/src/services/ibex/webhook-server/routes/crypto-receive.ts @@ -7,6 +7,10 @@ import { WalletCurrency, USDTAmount } from "@domain/shared" import { baseLogger } from "@services/logger" import { LockService } from "@services/lock" import { reconcileByTxHash } from "@services/bridge/reconciliation" +import { + alertIbexCryptoReceiveFailure, + alertIbexReconciliationFailed, +} from "@services/alerts/ibex-bridge-movement" import { writeIbexCryptoReceiveRequest } from "@services/frappe/BridgeTransferRequestWriter" import { authenticate, logRequest } from "../middleware" @@ -48,6 +52,12 @@ const cryptoReceiveHandler = async (req: Request, res: Response) => { const account = await AccountsRepository().findByBridgeEthereumAddress(address) if (account instanceof Error) { baseLogger.error({ address, tx_hash }, "Account not found for Ethereum address") + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "account_not_found", + title: "IBEX crypto receive: account not found for Bridge Ethereum address", + context: { address }, + }) return { status: "error", code: "account_not_found" } as CryptoReceiveResult } @@ -64,12 +74,23 @@ const cryptoReceiveHandler = async (req: Request, res: Response) => { { error: ibexLog, tx_hash }, "Failed to persist IBEX crypto receive log", ) + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "persist_failed", + title: "IBEX crypto receive log persistence failed", + detail: ibexLog.message, + context: { address }, + }) return { status: "error", code: "internal_error" } as CryptoReceiveResult } - reconcileByTxHash({ txHash: String(tx_hash) }).catch((err) => - baseLogger.error({ err, tx_hash }, "Real-time reconciliation failed"), - ) + reconcileByTxHash({ txHash: String(tx_hash) }).catch((err) => { + baseLogger.error({ err, tx_hash }, "Real-time reconciliation failed") + alertIbexReconciliationFailed({ + txHash: String(tx_hash), + detail: err instanceof Error ? err.message : String(err), + }) + }) const wallets = await listWalletsByAccountId(account.id) if (wallets instanceof Error) { @@ -77,18 +98,38 @@ const cryptoReceiveHandler = async (req: Request, res: Response) => { { accountId: account.id, error: wallets }, "Failed to list wallets", ) + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "wallet_list_failed", + title: "IBEX crypto receive: wallet list failed", + detail: wallets.message, + context: { accountId: account.id, address }, + }) return { status: "error", code: "wallet_list_failed" } as CryptoReceiveResult } const usdtWallet = wallets.find((w) => w.currency === WalletCurrency.Usdt) if (!usdtWallet) { baseLogger.error({ accountId: account.id }, "USDT wallet not found") + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "usdt_wallet_not_found", + title: "IBEX crypto receive: USDT wallet not found", + context: { accountId: account.id, address }, + }) return { status: "error", code: "usdt_wallet_not_found" } as CryptoReceiveResult } const usdtAmount = USDTAmount.fromNumber(amount) if (usdtAmount instanceof Error) { baseLogger.error({ amount, error: usdtAmount }, "Invalid USDT amount") + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "invalid_amount", + title: "IBEX crypto receive: invalid USDT amount", + detail: usdtAmount.message, + context: { accountId: account.id, address, amount }, + }) return { status: "error", code: "invalid_amount" } as CryptoReceiveResult } @@ -123,6 +164,17 @@ const cryptoReceiveHandler = async (req: Request, res: Response) => { }, "Failed to persist IBEX crypto receive ERPNext audit row", ) + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "erpnext_audit_failed", + title: "IBEX crypto receive ERPNext audit write failed", + detail: auditResult.message, + context: { + accountId: account.id, + walletId: usdtWallet.id, + address, + }, + }) return { status: "error", code: "erpnext_audit_failed" } as CryptoReceiveResult } @@ -135,6 +187,12 @@ const cryptoReceiveHandler = async (req: Request, res: Response) => { return { status: "success" } as CryptoReceiveResult } catch (error) { baseLogger.error({ error, tx_hash }, "Error processing crypto receive webhook") + alertIbexCryptoReceiveFailure({ + txHash: String(tx_hash), + code: "internal_error", + title: "IBEX crypto receive webhook processing error", + detail: error instanceof Error ? error.message : String(error), + }) return { status: "error", code: "internal_error" } as CryptoReceiveResult } }, diff --git a/test/flash/unit/services/alerts/dedup-key.spec.ts b/test/flash/unit/services/alerts/dedup-key.spec.ts new file mode 100644 index 000000000..d158bb62c --- /dev/null +++ b/test/flash/unit/services/alerts/dedup-key.spec.ts @@ -0,0 +1,109 @@ +import { + generateDedupKey, + informDedupTtlMs, + resolveDedupKey, +} from "@services/alerts/dedup-key" + +describe("generateDedupKey", () => { + it("uses coarse keys for Bridge API outage classes", () => { + expect(generateDedupKey.bridgeApi5xx()).toBe("bridge-api:5xx") + expect(generateDedupKey.bridgeApiTimeout()).toBe("bridge-api:timeout") + expect(generateDedupKey.bridgeApiNetwork()).toBe("bridge-api:network") + }) + + it("scopes ERPNext and webhook keys per resource", () => { + expect(generateDedupKey.erpnextDepositAudit("tr_1")).toBe( + "erpnext-audit:deposit:tr_1", + ) + expect(generateDedupKey.erpnextTransferCompletedAudit("tr_2")).toBe( + "erpnext-audit:transfer-complete:tr_2", + ) + expect(generateDedupKey.bridgeWebhookDeposit("wh_1")).toBe( + "bridge-webhook:deposit:wh_1", + ) + expect(generateDedupKey.bridgeWebhookTransfer("tr_3", "transfer.completed")).toBe( + "bridge-webhook:transfer:tr_3:transfer.completed", + ) + }) + + it("scopes IBEX Bridge movement keys per tx hash or transfer", () => { + expect(generateDedupKey.ibexCryptoReceive("0XABC")).toBe( + "ibex:crypto-receive:0xabc", + ) + expect(generateDedupKey.ibexReconcileBridgeWithoutIbex("0xabc")).toBe( + "ibex:reconcile:bridge-without-ibex:0xabc", + ) + expect(generateDedupKey.ibexReconcileIbexWithoutBridge("0xabc")).toBe( + "ibex:reconcile:ibex-without-bridge:0xabc", + ) + expect(generateDedupKey.ibexReconcileBridgeWithoutIbexTransfer("tr_1")).toBe( + "ibex:reconcile:bridge-without-ibex:transfer:tr_1", + ) + }) +}) + +describe("informDedupTtlMs", () => { + it("uses a shorter TTL for Bridge API outage keys", () => { + expect(informDedupTtlMs("bridge-api:5xx")).toBe(30 * 60 * 1000) + expect(informDedupTtlMs("erpnext-audit:deposit:tr_1")).toBe(60 * 60 * 1000) + }) +}) + +describe("resolveDedupKey", () => { + it("prefers an explicit dedupKey", () => { + expect( + resolveDedupKey({ + dedupKey: "custom-key", + source: "bridge-api", + severity: "critical", + title: "anything", + }), + ).toBe("custom-key") + }) + + it("falls back to outage keys for bridge-api alerts", () => { + expect( + resolveDedupKey({ + source: "bridge-api", + severity: "critical", + title: "Bridge API timeout on GET /transfers", + }), + ).toBe("bridge-api:timeout") + + expect( + resolveDedupKey({ + source: "bridge-api", + severity: "critical", + title: "Bridge API request failed on POST /customers", + }), + ).toBe("bridge-api:network") + + expect( + resolveDedupKey({ + source: "bridge-api", + severity: "critical", + title: "Bridge API 502 on GET /transfers", + }), + ).toBe("bridge-api:5xx") + }) + + it("falls back to IBEX movement keys", () => { + expect( + resolveDedupKey({ + source: "ibex", + severity: "warning", + title: "IBEX crypto receive ERPNext audit write failed", + context: { tx_hash: "0xabc" }, + }), + ).toBe("ibex:crypto-receive:0xabc") + + expect( + resolveDedupKey({ + source: "ibex", + severity: "warning", + title: "Bridge deposit without matching IBEX crypto receive", + context: { orphan_type: "bridge_without_ibex", tx_hash: "0xabc" }, + }), + ).toBe("ibex:reconcile:bridge-without-ibex:0xabc") + }) +}) diff --git a/test/flash/unit/services/alerts/ibex-bridge-movement.spec.ts b/test/flash/unit/services/alerts/ibex-bridge-movement.spec.ts new file mode 100644 index 000000000..aaa1ff20d --- /dev/null +++ b/test/flash/unit/services/alerts/ibex-bridge-movement.spec.ts @@ -0,0 +1,59 @@ +jest.mock("@services/alerts", () => ({ + alertBridge: jest.fn(), +})) + +import { alertBridge } from "@services/alerts" +import { + alertIbexCryptoReceiveFailure, + alertIbexReconciliationOrphan, +} from "@services/alerts/ibex-bridge-movement" + +describe("ibex bridge movement alerts", () => { + beforeEach(() => { + jest.clearAllMocks() + }) + + it("routes crypto receive failures as IBEX warnings", () => { + alertIbexCryptoReceiveFailure({ + txHash: "0xabc", + code: "erpnext_audit_failed", + title: "IBEX crypto receive ERPNext audit write failed", + detail: "timeout", + context: { accountId: "acc_1" }, + }) + + expect(alertBridge).toHaveBeenCalledWith({ + dedupKey: "ibex:crypto-receive:0xabc", + source: "ibex", + severity: "warning", + title: "IBEX crypto receive ERPNext audit write failed", + detail: "timeout", + context: { + tx_hash: "0xabc", + code: "erpnext_audit_failed", + accountId: "acc_1", + }, + }) + }) + + it("routes reconciliation orphans as IBEX warnings", () => { + alertIbexReconciliationOrphan({ + orphanType: "ibex_without_bridge", + txHash: "0xdef", + reason: "No Bridge deposit payment_processed found for IBEX tx hash within window", + }) + + expect(alertBridge).toHaveBeenCalledWith({ + dedupKey: "ibex:reconcile:ibex-without-bridge:0xdef", + source: "ibex", + severity: "warning", + title: "IBEX crypto receive without matching Bridge deposit", + detail: "No Bridge deposit payment_processed found for IBEX tx hash within window", + context: { + orphan_type: "ibex_without_bridge", + tx_hash: "0xdef", + transfer_id: undefined, + }, + }) + }) +}) diff --git a/test/flash/unit/services/alerts/index.spec.ts b/test/flash/unit/services/alerts/index.spec.ts new file mode 100644 index 000000000..675185119 --- /dev/null +++ b/test/flash/unit/services/alerts/index.spec.ts @@ -0,0 +1,76 @@ +jest.mock("@services/alerts/slack", () => ({ + sendSlack: jest.fn().mockResolvedValue(undefined), +})) + +jest.mock("@services/alerts/discord", () => ({ + sendDiscord: jest.fn().mockResolvedValue(undefined), +})) + +jest.mock("@services/alerts/pagerduty", () => ({ + sendPagerDuty: jest.fn().mockResolvedValue(undefined), +})) + +import { alertBridge } from "@services/alerts" +import { resetInformDedup } from "@services/alerts/inform-dedup" +import { sendDiscord } from "@services/alerts/discord" +import { sendPagerDuty } from "@services/alerts/pagerduty" +import { sendSlack } from "@services/alerts/slack" + +describe("alertBridge", () => { + beforeEach(() => { + jest.clearAllMocks() + resetInformDedup() + }) + + it("fans out critical alerts to inform channels and PagerDuty", async () => { + alertBridge({ + dedupKey: "bridge-api:5xx", + source: "bridge-api", + severity: "critical", + title: "Bridge API 502 on GET /transfers", + }) + + await Promise.resolve() + + expect(sendSlack).toHaveBeenCalledTimes(1) + expect(sendDiscord).toHaveBeenCalledTimes(1) + expect(sendPagerDuty).toHaveBeenCalledTimes(1) + expect(sendPagerDuty).toHaveBeenCalledWith( + expect.objectContaining({ dedupKey: "bridge-api:5xx" }), + ) + }) + + it("does not page PagerDuty for warning alerts", async () => { + alertBridge({ + dedupKey: "ibex:warning:tx_1", + source: "ibex", + severity: "warning", + title: "IBEX movement failed", + }) + + await Promise.resolve() + + expect(sendSlack).toHaveBeenCalledTimes(1) + expect(sendDiscord).toHaveBeenCalledTimes(1) + expect(sendPagerDuty).not.toHaveBeenCalled() + }) + + it("suppresses duplicate Slack and Discord alerts for the same dedup key", async () => { + const alert = { + dedupKey: "bridge-api:5xx", + source: "bridge-api" as const, + severity: "critical" as const, + title: "Bridge API 502 on GET /transfers", + } + + alertBridge(alert) + alertBridge(alert) + alertBridge(alert) + + await Promise.resolve() + + expect(sendSlack).toHaveBeenCalledTimes(1) + expect(sendDiscord).toHaveBeenCalledTimes(1) + expect(sendPagerDuty).toHaveBeenCalledTimes(3) + }) +}) diff --git a/test/flash/unit/services/alerts/inform-dedup.spec.ts b/test/flash/unit/services/alerts/inform-dedup.spec.ts new file mode 100644 index 000000000..c1c0fcc16 --- /dev/null +++ b/test/flash/unit/services/alerts/inform-dedup.spec.ts @@ -0,0 +1,34 @@ +import { claimInformSlot, resetInformDedup } from "@services/alerts/inform-dedup" + +describe("claimInformSlot", () => { + beforeEach(() => { + resetInformDedup() + }) + + it("allows the first inform for a dedup key", () => { + expect(claimInformSlot("erpnext-audit:deposit:tr_1", 1_000)).toBe(true) + }) + + it("suppresses duplicate informs within the TTL", () => { + const key = "bridge-api:5xx" + const start = 10_000 + + expect(claimInformSlot(key, start)).toBe(true) + expect(claimInformSlot(key, start + 1_000)).toBe(false) + expect(claimInformSlot(key, start + 29 * 60 * 1000)).toBe(false) + }) + + it("allows a new inform after the TTL expires", () => { + const key = "bridge-api:5xx" + const start = 10_000 + + expect(claimInformSlot(key, start)).toBe(true) + expect(claimInformSlot(key, start + 30 * 60 * 1000)).toBe(true) + }) + + it("tracks different dedup keys independently", () => { + expect(claimInformSlot("erpnext-audit:deposit:tr_a", 1_000)).toBe(true) + expect(claimInformSlot("erpnext-audit:deposit:tr_b", 1_000)).toBe(true) + expect(claimInformSlot("erpnext-audit:deposit:tr_a", 2_000)).toBe(false) + }) +}) diff --git a/test/flash/unit/services/alerts/pagerduty.spec.ts b/test/flash/unit/services/alerts/pagerduty.spec.ts new file mode 100644 index 000000000..d95b45199 --- /dev/null +++ b/test/flash/unit/services/alerts/pagerduty.spec.ts @@ -0,0 +1,44 @@ +jest.mock("@config", () => ({ + ALERT_PAGERDUTY_ROUTING_KEY: "test-routing-key", +})) + +jest.mock("@services/tracing", () => ({ + recordExceptionInCurrentSpan: jest.fn(), +})) + +jest.mock("axios", () => ({ + post: jest.fn().mockResolvedValue({ status: 202 }), +})) + +import axios from "axios" +import { sendPagerDuty } from "@services/alerts/pagerduty" + +describe("sendPagerDuty", () => { + beforeEach(() => { + jest.clearAllMocks() + }) + + it("includes dedup_key in the Events API v2 payload", async () => { + await sendPagerDuty({ + dedupKey: "bridge-api:5xx", + source: "bridge-api", + severity: "critical", + title: "Bridge API 502 on GET /transfers", + context: { method: "GET", path: "/transfers" }, + }) + + expect(axios.post).toHaveBeenCalledWith( + "https://events.pagerduty.com/v2/enqueue", + expect.objectContaining({ + routing_key: "test-routing-key", + event_action: "trigger", + dedup_key: "bridge-api:5xx", + payload: expect.objectContaining({ + summary: "[bridge:bridge-api] Bridge API 502 on GET /transfers", + severity: "critical", + }), + }), + expect.any(Object), + ) + }) +}) diff --git a/test/flash/unit/services/bridge/reconciliation.spec.ts b/test/flash/unit/services/bridge/reconciliation.spec.ts index 2c4c3a6d6..85c87eeb0 100644 --- a/test/flash/unit/services/bridge/reconciliation.spec.ts +++ b/test/flash/unit/services/bridge/reconciliation.spec.ts @@ -10,12 +10,12 @@ jest.mock("@services/logger", () => ({ })) jest.mock("@services/mongoose/schema", () => ({ - BridgeDepositLog: { findOne: jest.fn(), find: jest.fn() }, - IbexCryptoReceiveLog: { findOne: jest.fn() }, + BridgeDeposits: { findOne: jest.fn(), find: jest.fn() }, + IbexCryptoReceive: { findOne: jest.fn() }, })) jest.mock("@services/mongoose/ibex-crypto-receive-log", () => ({ - findIbexCryptoReceiveLogsSince: jest.fn(), + findIbexCryptoReceivesSince: jest.fn(), })) jest.mock("@services/mongoose/bridge-reconciliation-orphan", () => ({ @@ -33,13 +33,19 @@ jest.mock("@domain/pubsub", () => ({ }, })) -import { BridgeDepositLog, IbexCryptoReceiveLog } from "@services/mongoose/schema" -import { findIbexCryptoReceiveLogsSince } from "@services/mongoose/ibex-crypto-receive-log" +jest.mock("@services/alerts/ibex-bridge-movement", () => ({ + alertIbexReconciliationOrphan: jest.fn(), + alertIbexReconciliationFailed: jest.fn(), +})) + +import { BridgeDeposits, IbexCryptoReceive } from "@services/mongoose/schema" +import { findIbexCryptoReceivesSince } from "@services/mongoose/ibex-crypto-receive-log" import { upsertBridgeReconciliationOrphan, resolveOrphansByTxHash, } from "@services/mongoose/bridge-reconciliation-orphan" import { PubSubService } from "@services/pubsub" +import { alertIbexReconciliationOrphan } from "@services/alerts/ibex-bridge-movement" import { reconcileByTxHash, reconcileBridgeAndIbexDeposits, @@ -91,8 +97,8 @@ beforeEach(() => { describe("reconcileByTxHash", () => { describe("both sides found β†’ matched", () => { beforeEach(() => { - ;(BridgeDepositLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) - ;(IbexCryptoReceiveLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) }) it("returns status matched", async () => { @@ -134,15 +140,15 @@ describe("reconcileByTxHash", () => { expect(result).not.toBeInstanceOf(Error) if (result instanceof Error) return expect(result.txHash).toBe(NORM_HASH) - const [bridgeCall] = (BridgeDepositLog.findOne as jest.Mock).mock.calls + const [bridgeCall] = (BridgeDeposits.findOne as jest.Mock).mock.calls expect(bridgeCall[0].destinationTxHash.$regex.flags).toContain("i") }) }) describe("only Bridge found β†’ bridge_without_ibex", () => { beforeEach(() => { - ;(BridgeDepositLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) - ;(IbexCryptoReceiveLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) }) it("returns status unmatched with correct orphanType", async () => { @@ -182,12 +188,23 @@ describe("reconcileByTxHash", () => { }), ) }) + + it("alerts ops when Bridge has no matching IBEX receive", async () => { + await reconcileByTxHash({ txHash: TX_HASH }) + expect(alertIbexReconciliationOrphan).toHaveBeenCalledWith( + expect.objectContaining({ + orphanType: "bridge_without_ibex", + txHash: NORM_HASH, + transferId: BRIDGE_DEPOSIT.transferId, + }), + ) + }) }) describe("only IBEX found β†’ ibex_without_bridge", () => { beforeEach(() => { - ;(BridgeDepositLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) - ;(IbexCryptoReceiveLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) }) it("returns status unmatched with correct orphanType", async () => { @@ -208,13 +225,23 @@ describe("reconcileByTxHash", () => { }), ) }) + + it("alerts ops when IBEX has no matching Bridge deposit", async () => { + await reconcileByTxHash({ txHash: TX_HASH }) + expect(alertIbexReconciliationOrphan).toHaveBeenCalledWith( + expect.objectContaining({ + orphanType: "ibex_without_bridge", + txHash: NORM_HASH, + }), + ) + }) }) describe("self-healing: second call with both sides resolves orphan", () => { it("resolves orphan when called again after missing side arrives", async () => { // First call: only Bridge - ;(BridgeDepositLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) - ;(IbexCryptoReceiveLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) await reconcileByTxHash({ txHash: TX_HASH }) expect(upsertBridgeReconciliationOrphan).toHaveBeenCalledTimes(1) @@ -223,8 +250,8 @@ describe("reconcileByTxHash", () => { ;(resolveOrphansByTxHash as jest.Mock).mockResolvedValue({ resolvedCount: 1 }) // Second call: both sides present (IBEX webhook arrived) - ;(BridgeDepositLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) - ;(IbexCryptoReceiveLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) const result = await reconcileByTxHash({ txHash: TX_HASH }) expect(result).not.toBeInstanceOf(Error) @@ -236,11 +263,11 @@ describe("reconcileByTxHash", () => { }) describe("Bridge query uses payment_processed state filter", () => { - it("passes state: payment_processed to BridgeDepositLog.findOne", async () => { - ;(BridgeDepositLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) - ;(IbexCryptoReceiveLog.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) + it("passes state: payment_processed to BridgeDeposits.findOne", async () => { + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) await reconcileByTxHash({ txHash: TX_HASH }) - expect(BridgeDepositLog.findOne).toHaveBeenCalledWith( + expect(BridgeDeposits.findOne).toHaveBeenCalledWith( expect.objectContaining({ state: "payment_processed" }), ) }) @@ -256,8 +283,8 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("all deposits matched", () => { it("returns zero orphans when every Bridge deposit has a matching IBEX receive", async () => { - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue(makeBridgeFind([BRIDGE_DEPOSIT])) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE]) + ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([BRIDGE_DEPOSIT])) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE]) const result = await reconcileBridgeAndIbexDeposits() expect(result).not.toBeInstanceOf(Error) @@ -272,8 +299,8 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("Bridge deposit with no matching IBEX receive", () => { it("flags as bridge_without_ibex orphan", async () => { - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue(makeBridgeFind([BRIDGE_DEPOSIT])) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue([]) + ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([BRIDGE_DEPOSIT])) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([]) const result = await reconcileBridgeAndIbexDeposits() expect(result).not.toBeInstanceOf(Error) @@ -293,8 +320,8 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("Bridge deposit with no destinationTxHash", () => { it("flags as bridge-no-tx:{transferId} orphan", async () => { const depositNoHash = { ...BRIDGE_DEPOSIT, destinationTxHash: undefined } - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue(makeBridgeFind([depositNoHash])) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue([]) + ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([depositNoHash])) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([]) const result = await reconcileBridgeAndIbexDeposits() expect(result).not.toBeInstanceOf(Error) @@ -311,8 +338,8 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("IBEX receive with no matching Bridge deposit", () => { it("flags as ibex_without_bridge orphan", async () => { - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue(makeBridgeFind([])) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE]) + ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([])) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE]) const result = await reconcileBridgeAndIbexDeposits() expect(result).not.toBeInstanceOf(Error) @@ -329,12 +356,12 @@ describe("reconcileBridgeAndIbexDeposits", () => { }) describe("batch uses payment_processed state filter", () => { - it("passes state: payment_processed to BridgeDepositLog.find", async () => { - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue(makeBridgeFind([])) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue([]) + it("passes state: payment_processed to BridgeDeposits.find", async () => { + ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([])) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([]) await reconcileBridgeAndIbexDeposits() - expect(BridgeDepositLog.find).toHaveBeenCalledWith( + expect(BridgeDeposits.find).toHaveBeenCalledWith( expect.objectContaining({ state: "payment_processed" }), ) }) @@ -345,10 +372,10 @@ describe("reconcileBridgeAndIbexDeposits", () => { const deposit2 = { ...BRIDGE_DEPOSIT, transferId: "tr_002", destinationTxHash: "0xother" } const ibex2 = { ...IBEX_RECEIVE, txHash: "0xorphan_ibex" } - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue( + ;(BridgeDeposits.find as jest.Mock).mockReturnValue( makeBridgeFind([BRIDGE_DEPOSIT, deposit2]), ) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE, ibex2]) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE, ibex2]) const result = await reconcileBridgeAndIbexDeposits() expect(result).not.toBeInstanceOf(Error) @@ -365,9 +392,9 @@ describe("reconcileBridgeAndIbexDeposits", () => { }) describe("error handling", () => { - it("returns an Error when findIbexCryptoReceiveLogsSince fails", async () => { - ;(BridgeDepositLog.find as jest.Mock).mockReturnValue(makeBridgeFind([])) - ;(findIbexCryptoReceiveLogsSince as jest.Mock).mockResolvedValue( + it("returns an Error when findIbexCryptoReceivesSince fails", async () => { + ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([])) + ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue( new Error("mongo connection lost"), ) diff --git a/test/flash/unit/services/ibex/webhook-server/routes/crypto-receive.spec.ts b/test/flash/unit/services/ibex/webhook-server/routes/crypto-receive.spec.ts index e4e42d929..af7830a4c 100644 --- a/test/flash/unit/services/ibex/webhook-server/routes/crypto-receive.spec.ts +++ b/test/flash/unit/services/ibex/webhook-server/routes/crypto-receive.spec.ts @@ -27,10 +27,19 @@ jest.mock("@services/bridge/reconciliation", () => ({ reconcileByTxHash: jest.fn().mockResolvedValue({ status: "matched" }), })) +jest.mock("@app/bridge/send-deposit-notification", () => ({ + sendBridgeDepositNotificationBestEffort: jest.fn().mockResolvedValue(undefined), +})) + jest.mock("@services/frappe/BridgeTransferRequestWriter", () => ({ writeIbexCryptoReceiveRequest: jest.fn(), })) +jest.mock("@services/alerts/ibex-bridge-movement", () => ({ + alertIbexCryptoReceiveFailure: jest.fn(), + alertIbexReconciliationFailed: jest.fn(), +})) + import { cryptoReceiveHandler } from "@services/ibex/webhook-server/routes/crypto-receive" import { AccountsRepository } from "@services/mongoose/accounts" import { createIbexCryptoReceive } from "@services/mongoose/ibex-crypto-receive-log" @@ -38,6 +47,7 @@ import { listWalletsByAccountId } from "@app/wallets" import { LockService } from "@services/lock" import { WalletCurrency } from "@domain/shared" import { writeIbexCryptoReceiveRequest } from "@services/frappe/BridgeTransferRequestWriter" +import { alertIbexCryptoReceiveFailure } from "@services/alerts/ibex-bridge-movement" const ACCOUNT_ID = "account-001" as AccountId const WALLET_ID = "wallet-usdt-001" as WalletId @@ -155,6 +165,13 @@ describe("cryptoReceiveHandler", () => { expect(res.status).toHaveBeenCalledWith(500) expect(res.json).toHaveBeenCalledWith({ error: "erpnext_audit_failed" }) + expect(alertIbexCryptoReceiveFailure).toHaveBeenCalledWith( + expect.objectContaining({ + txHash: TX_HASH, + code: "erpnext_audit_failed", + title: "IBEX crypto receive ERPNext audit write failed", + }), + ) }) it("rejects legacy Tron USDT receive webhooks for the ETH-USDT Cash Wallet path", async () => { From 160aad79d626ebe8441edb35762ede8e938ad3ac Mon Sep 17 00:00:00 2001 From: Vandana Date: Tue, 9 Jun 2026 10:59:14 -0700 Subject: [PATCH 6/7] chore(alerts): format bridge alert updates --- src/services/alerts/dedup-key.ts | 9 ++-- src/services/alerts/ibex-bridge-movement.ts | 3 +- .../unit/services/alerts/dedup-key.spec.ts | 4 +- .../services/bridge/reconciliation.spec.ts | 42 ++++++++++++++----- 4 files changed, 41 insertions(+), 17 deletions(-) diff --git a/src/services/alerts/dedup-key.ts b/src/services/alerts/dedup-key.ts index 0e2ed0837..54287770d 100644 --- a/src/services/alerts/dedup-key.ts +++ b/src/services/alerts/dedup-key.ts @@ -28,7 +28,8 @@ export const generateDedupKey = { `ibex:reconcile:bridge-without-ibex:transfer:${transferId}`, ibexReconcileIbexWithoutBridge: (txHash: string) => `ibex:reconcile:ibex-without-bridge:${txHash.toLowerCase()}`, - ibexReconcileFailed: (txHash: string) => `ibex:reconcile:failed:${txHash.toLowerCase()}`, + ibexReconcileFailed: (txHash: string) => + `ibex:reconcile:failed:${txHash.toLowerCase()}`, } const truncateDedupKey = (key: string): string => @@ -43,11 +44,13 @@ export const resolveDedupKey = (alert: BridgeAlert): string => { switch (alert.source) { case "bridge-api": if (alert.title.includes("timeout")) return generateDedupKey.bridgeApiTimeout() - if (alert.title.includes("request failed")) return generateDedupKey.bridgeApiNetwork() + if (alert.title.includes("request failed")) + return generateDedupKey.bridgeApiNetwork() return generateDedupKey.bridgeApi5xx() case "erpnext-audit": { const transferId = String(ctx.transfer_id ?? "unknown") - if (alert.title.includes("deposit")) return generateDedupKey.erpnextDepositAudit(transferId) + if (alert.title.includes("deposit")) + return generateDedupKey.erpnextDepositAudit(transferId) if (alert.title.includes("failure")) { return generateDedupKey.erpnextTransferFailedAudit(transferId) } diff --git a/src/services/alerts/ibex-bridge-movement.ts b/src/services/alerts/ibex-bridge-movement.ts index 9b254e63d..a8bfd8ab8 100644 --- a/src/services/alerts/ibex-bridge-movement.ts +++ b/src/services/alerts/ibex-bridge-movement.ts @@ -1,6 +1,7 @@ -import { alertBridge } from "./index" import { generateDedupKey } from "./dedup-key" +import { alertBridge } from "./index" + type IbexMovementAlert = { title: string detail?: string diff --git a/test/flash/unit/services/alerts/dedup-key.spec.ts b/test/flash/unit/services/alerts/dedup-key.spec.ts index d158bb62c..4a8d35221 100644 --- a/test/flash/unit/services/alerts/dedup-key.spec.ts +++ b/test/flash/unit/services/alerts/dedup-key.spec.ts @@ -27,9 +27,7 @@ describe("generateDedupKey", () => { }) it("scopes IBEX Bridge movement keys per tx hash or transfer", () => { - expect(generateDedupKey.ibexCryptoReceive("0XABC")).toBe( - "ibex:crypto-receive:0xabc", - ) + expect(generateDedupKey.ibexCryptoReceive("0XABC")).toBe("ibex:crypto-receive:0xabc") expect(generateDedupKey.ibexReconcileBridgeWithoutIbex("0xabc")).toBe( "ibex:reconcile:bridge-without-ibex:0xabc", ) diff --git a/test/flash/unit/services/bridge/reconciliation.spec.ts b/test/flash/unit/services/bridge/reconciliation.spec.ts index 85c87eeb0..9e8a91092 100644 --- a/test/flash/unit/services/bridge/reconciliation.spec.ts +++ b/test/flash/unit/services/bridge/reconciliation.spec.ts @@ -97,8 +97,12 @@ beforeEach(() => { describe("reconcileByTxHash", () => { describe("both sides found β†’ matched", () => { beforeEach(() => { - ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) - ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(BRIDGE_DEPOSIT), + ) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(IBEX_RECEIVE), + ) }) it("returns status matched", async () => { @@ -147,7 +151,9 @@ describe("reconcileByTxHash", () => { describe("only Bridge found β†’ bridge_without_ibex", () => { beforeEach(() => { - ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(BRIDGE_DEPOSIT), + ) ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) }) @@ -204,7 +210,9 @@ describe("reconcileByTxHash", () => { describe("only IBEX found β†’ ibex_without_bridge", () => { beforeEach(() => { ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) - ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(IBEX_RECEIVE), + ) }) it("returns status unmatched with correct orphanType", async () => { @@ -240,7 +248,9 @@ describe("reconcileByTxHash", () => { describe("self-healing: second call with both sides resolves orphan", () => { it("resolves orphan when called again after missing side arrives", async () => { // First call: only Bridge - ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(BRIDGE_DEPOSIT), + ) ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(null)) await reconcileByTxHash({ txHash: TX_HASH }) expect(upsertBridgeReconciliationOrphan).toHaveBeenCalledTimes(1) @@ -250,8 +260,12 @@ describe("reconcileByTxHash", () => { ;(resolveOrphansByTxHash as jest.Mock).mockResolvedValue({ resolvedCount: 1 }) // Second call: both sides present (IBEX webhook arrived) - ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue(makeLeanQuery(BRIDGE_DEPOSIT)) - ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue(makeLeanQuery(IBEX_RECEIVE)) + ;(BridgeDeposits.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(BRIDGE_DEPOSIT), + ) + ;(IbexCryptoReceive.findOne as jest.Mock).mockReturnValue( + makeLeanQuery(IBEX_RECEIVE), + ) const result = await reconcileByTxHash({ txHash: TX_HASH }) expect(result).not.toBeInstanceOf(Error) @@ -283,7 +297,9 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("all deposits matched", () => { it("returns zero orphans when every Bridge deposit has a matching IBEX receive", async () => { - ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([BRIDGE_DEPOSIT])) + ;(BridgeDeposits.find as jest.Mock).mockReturnValue( + makeBridgeFind([BRIDGE_DEPOSIT]), + ) ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([IBEX_RECEIVE]) const result = await reconcileBridgeAndIbexDeposits() @@ -299,7 +315,9 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("Bridge deposit with no matching IBEX receive", () => { it("flags as bridge_without_ibex orphan", async () => { - ;(BridgeDeposits.find as jest.Mock).mockReturnValue(makeBridgeFind([BRIDGE_DEPOSIT])) + ;(BridgeDeposits.find as jest.Mock).mockReturnValue( + makeBridgeFind([BRIDGE_DEPOSIT]), + ) ;(findIbexCryptoReceivesSince as jest.Mock).mockResolvedValue([]) const result = await reconcileBridgeAndIbexDeposits() @@ -369,7 +387,11 @@ describe("reconcileBridgeAndIbexDeposits", () => { describe("mixed scenario", () => { it("counts matched and unmatched independently", async () => { - const deposit2 = { ...BRIDGE_DEPOSIT, transferId: "tr_002", destinationTxHash: "0xother" } + const deposit2 = { + ...BRIDGE_DEPOSIT, + transferId: "tr_002", + destinationTxHash: "0xother", + } const ibex2 = { ...IBEX_RECEIVE, txHash: "0xorphan_ibex" } ;(BridgeDeposits.find as jest.Mock).mockReturnValue( From 48848301030dda3d41078151ab722b659529fe58 Mon Sep 17 00:00:00 2001 From: Vandana Date: Tue, 9 Jun 2026 11:05:27 -0700 Subject: [PATCH 7/7] chore(alerts): trim dedup cleanup noise --- src/services/alerts/dedup-key.ts | 58 +----------------- src/services/alerts/discord.ts | 6 +- src/services/alerts/index.ts | 10 +-- src/services/alerts/index.types.ts | 2 +- src/services/alerts/pagerduty.ts | 2 +- src/services/alerts/slack.ts | 2 +- .../unit/services/alerts/dedup-key.spec.ts | 61 ++----------------- 7 files changed, 16 insertions(+), 125 deletions(-) diff --git a/src/services/alerts/dedup-key.ts b/src/services/alerts/dedup-key.ts index 54287770d..2a1cb2584 100644 --- a/src/services/alerts/dedup-key.ts +++ b/src/services/alerts/dedup-key.ts @@ -1,5 +1,3 @@ -import { BridgeAlert } from "./index.types" - export const PAGERDUTY_DEDUP_KEY_MAX = 255 const OUTAGE_TTL_MS = 30 * 60 * 1000 @@ -32,59 +30,5 @@ export const generateDedupKey = { `ibex:reconcile:failed:${txHash.toLowerCase()}`, } -const truncateDedupKey = (key: string): string => +export const normalizeDedupKey = (key: string): string => key.length <= PAGERDUTY_DEDUP_KEY_MAX ? key : key.slice(0, PAGERDUTY_DEDUP_KEY_MAX) - -/** Stable PagerDuty / inform dedup key; prefers explicit alert.dedupKey. */ -export const resolveDedupKey = (alert: BridgeAlert): string => { - if (alert.dedupKey) return truncateDedupKey(alert.dedupKey) - - const ctx = alert.context ?? {} - - switch (alert.source) { - case "bridge-api": - if (alert.title.includes("timeout")) return generateDedupKey.bridgeApiTimeout() - if (alert.title.includes("request failed")) - return generateDedupKey.bridgeApiNetwork() - return generateDedupKey.bridgeApi5xx() - case "erpnext-audit": { - const transferId = String(ctx.transfer_id ?? "unknown") - if (alert.title.includes("deposit")) - return generateDedupKey.erpnextDepositAudit(transferId) - if (alert.title.includes("failure")) { - return generateDedupKey.erpnextTransferFailedAudit(transferId) - } - return generateDedupKey.erpnextTransferCompletedAudit(transferId) - } - case "bridge-webhook": { - const eventId = String(ctx.event_id ?? ctx.transfer_id ?? "unknown") - const transferId = String(ctx.transfer_id ?? "unknown") - const event = String(ctx.event ?? "unknown") - if (alert.title.includes("deposit")) { - return generateDedupKey.bridgeWebhookDeposit(eventId) - } - return generateDedupKey.bridgeWebhookTransfer(transferId, event) - } - case "ibex": { - const txHash = String(ctx.tx_hash ?? ctx.txHash ?? "unknown") - const orphanType = String(ctx.orphan_type ?? "") - if (orphanType === "ibex_without_bridge") { - return generateDedupKey.ibexReconcileIbexWithoutBridge(txHash) - } - if (orphanType === "bridge_without_ibex") { - if (txHash !== "unknown") { - return generateDedupKey.ibexReconcileBridgeWithoutIbex(txHash) - } - return generateDedupKey.ibexReconcileBridgeWithoutIbexTransfer( - String(ctx.transfer_id ?? "unknown"), - ) - } - if (alert.title.includes("reconciliation failed")) { - return generateDedupKey.ibexReconcileFailed(txHash) - } - return generateDedupKey.ibexCryptoReceive(txHash) - } - default: - return truncateDedupKey(`bridge:${alert.source}:${alert.title}`) - } -} diff --git a/src/services/alerts/discord.ts b/src/services/alerts/discord.ts index 2dedb6c96..df9ba0bac 100644 --- a/src/services/alerts/discord.ts +++ b/src/services/alerts/discord.ts @@ -12,14 +12,14 @@ const DISCORD_CONTENT_MAX = 1900 export const sendDiscord = async (alert: BridgeAlert): Promise => { if (!ALERT_DISCORD_WEBHOOK_URL) return - const icon = alert.severity === "critical" ? "🚨" : "⚠️" - let content = `${icon} **Bridge alert** β€” ${alert.title}\nsource: \`${alert.source}\` Β· severity: \`${alert.severity}\`` + const label = alert.severity === "critical" ? "[CRITICAL]" : "[WARNING]" + let content = `${label} **Bridge alert** - ${alert.title}\nsource: \`${alert.source}\` | severity: \`${alert.severity}\`` if (alert.detail) content += `\n${alert.detail}` if (alert.context) { content += "\n```json\n" + JSON.stringify(alert.context, null, 2) + "\n```" } if (content.length > DISCORD_CONTENT_MAX) { - content = content.slice(0, DISCORD_CONTENT_MAX) + "…" + content = content.slice(0, DISCORD_CONTENT_MAX) + "..." } try { diff --git a/src/services/alerts/index.ts b/src/services/alerts/index.ts index 0091505ef..95f7a8255 100644 --- a/src/services/alerts/index.ts +++ b/src/services/alerts/index.ts @@ -1,7 +1,7 @@ import { sendPagerDuty } from "./pagerduty" import { sendSlack } from "./slack" import { sendDiscord } from "./discord" -import { resolveDedupKey } from "./dedup-key" +import { normalizeDedupKey } from "./dedup-key" import { claimInformSlot } from "./inform-dedup" import { BridgeAlert } from "./index.types" @@ -10,20 +10,20 @@ export { generateDedupKey } from "./dedup-key" /** * Fire-and-forget fan-out of a Bridge alert to the configured destinations - * (ENG-361). Returns immediately; delivery is best-effort β€” each sender catches + * (ENG-361). Returns immediately; delivery is best-effort: each sender catches * its own errors and no-ops when its credential/URL is unset, so it never throws * or rejects into the caller (no need to await or handle it). * * Routing: - * - critical β†’ page on-call (PagerDuty) + inform (Slack/Mattermost, Discord) - * - warning β†’ inform (Slack/Mattermost, Discord) only + * - critical: page on-call (PagerDuty) + inform (Slack/Mattermost, Discord) + * - warning: inform (Slack/Mattermost, Discord) only * * Dedup: * - PagerDuty: Events API v2 dedup_key groups triggers into one incident. * - Slack / Discord: first alert per dedup key within TTL only. */ export const alertBridge = (alert: BridgeAlert): void => { - const dedupKey = resolveDedupKey(alert) + const dedupKey = normalizeDedupKey(alert.dedupKey) const alertWithKey: BridgeAlert = { ...alert, dedupKey } const deliver = async () => { diff --git a/src/services/alerts/index.types.ts b/src/services/alerts/index.types.ts index 09d153f08..b93b6c91c 100644 --- a/src/services/alerts/index.types.ts +++ b/src/services/alerts/index.types.ts @@ -5,7 +5,7 @@ export type AlertSeverity = "critical" | "warning" export type AlertSource = "bridge-webhook" | "bridge-api" | "ibex" | "erpnext-audit" export interface BridgeAlert { - dedupKey?: string + dedupKey: string source: AlertSource severity: AlertSeverity title: string diff --git a/src/services/alerts/pagerduty.ts b/src/services/alerts/pagerduty.ts index 4f14827d9..46d8317e5 100644 --- a/src/services/alerts/pagerduty.ts +++ b/src/services/alerts/pagerduty.ts @@ -7,7 +7,7 @@ import { BridgeAlert } from "./index.types" const PAGERDUTY_EVENTS_URL = "https://events.pagerduty.com/v2/enqueue" -// PagerDuty Events API v2 β€” triggers a paging incident. "critical" and +// PagerDuty Events API v2 triggers a paging incident. "critical" and // "warning" are both valid PD payload severities, so we pass them through. export const sendPagerDuty = async (alert: BridgeAlert): Promise => { if (!ALERT_PAGERDUTY_ROUTING_KEY) return diff --git a/src/services/alerts/slack.ts b/src/services/alerts/slack.ts index a4c9456f4..10a260088 100644 --- a/src/services/alerts/slack.ts +++ b/src/services/alerts/slack.ts @@ -11,7 +11,7 @@ export const sendSlack = async (alert: BridgeAlert): Promise => { const icon = alert.severity === "critical" ? ":rotating_light:" : ":warning:" const lines = [ - `${icon} *Bridge alert* β€” ${alert.title}`, + `${icon} *Bridge alert* - ${alert.title}`, `*source:* \`${alert.source}\` *severity:* \`${alert.severity}\``, ] if (alert.detail) lines.push(alert.detail) diff --git a/test/flash/unit/services/alerts/dedup-key.spec.ts b/test/flash/unit/services/alerts/dedup-key.spec.ts index 4a8d35221..f06f18c1d 100644 --- a/test/flash/unit/services/alerts/dedup-key.spec.ts +++ b/test/flash/unit/services/alerts/dedup-key.spec.ts @@ -1,7 +1,7 @@ import { generateDedupKey, informDedupTtlMs, - resolveDedupKey, + normalizeDedupKey, } from "@services/alerts/dedup-key" describe("generateDedupKey", () => { @@ -47,61 +47,8 @@ describe("informDedupTtlMs", () => { }) }) -describe("resolveDedupKey", () => { - it("prefers an explicit dedupKey", () => { - expect( - resolveDedupKey({ - dedupKey: "custom-key", - source: "bridge-api", - severity: "critical", - title: "anything", - }), - ).toBe("custom-key") - }) - - it("falls back to outage keys for bridge-api alerts", () => { - expect( - resolveDedupKey({ - source: "bridge-api", - severity: "critical", - title: "Bridge API timeout on GET /transfers", - }), - ).toBe("bridge-api:timeout") - - expect( - resolveDedupKey({ - source: "bridge-api", - severity: "critical", - title: "Bridge API request failed on POST /customers", - }), - ).toBe("bridge-api:network") - - expect( - resolveDedupKey({ - source: "bridge-api", - severity: "critical", - title: "Bridge API 502 on GET /transfers", - }), - ).toBe("bridge-api:5xx") - }) - - it("falls back to IBEX movement keys", () => { - expect( - resolveDedupKey({ - source: "ibex", - severity: "warning", - title: "IBEX crypto receive ERPNext audit write failed", - context: { tx_hash: "0xabc" }, - }), - ).toBe("ibex:crypto-receive:0xabc") - - expect( - resolveDedupKey({ - source: "ibex", - severity: "warning", - title: "Bridge deposit without matching IBEX crypto receive", - context: { orphan_type: "bridge_without_ibex", tx_hash: "0xabc" }, - }), - ).toBe("ibex:reconcile:bridge-without-ibex:0xabc") +describe("normalizeDedupKey", () => { + it("truncates keys to PagerDuty's maximum dedup_key length", () => { + expect(normalizeDedupKey("a".repeat(300))).toHaveLength(255) }) })