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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 20 additions & 51 deletions src/alerts/dispatcher.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import type Database from "better-sqlite3";
import { getUndeliveredAlerts, markAlertDelivered, incrementRetryCount, MAX_RETRY_COUNT } from "../db/repositories.js";
import { buildAlertEvent, type AlertEvent } from "./types.js";
import { buildAlertEvent, type AlertEvent, type AlertChannel } from "./types.js";
import { sendWebhookAlert } from "./webhook.js";
import { sendSlackAlert } from "./slack.js";
import { sendPagerDutyAlert } from "./pagerduty.js";
Expand All @@ -10,26 +10,26 @@ import { getLogger } from "../logging/index.js";

const logger = getLogger().child({ component: "AlertDispatcher" });

// ─── Public contract ─────────────────────────────────────────────────────────

export interface DeliveryResult {
/** Total alerts processed (includes failed). */
attempted: number;
/** Alerts successfully sent and marked delivered = 1. */
delivered: number;
/** Alerts that threw during delivery — retry count incremented. */
failed: number;
/** Alerts that exceeded max retries and were abandoned. */
abandoned: number;
/** Error messages for each failed delivery. */
errors: string[];
}

// ─── Core implementation ──────────────────────────────────────────────────────
export const DEFAULT_CHANNELS: Record<string, AlertChannel> = {
webhook: { send: sendWebhookAlert },
slack: { send: sendSlackAlert },
pagerduty: { send: sendPagerDutyAlert },
discord: { send: sendDiscordAlert },
telegram: { send: sendTelegramAlert },
};

export async function deliverPendingAlerts(
db: Database.Database,
network: string,
channels: Record<string, AlertChannel> = DEFAULT_CHANNELS,
): Promise<DeliveryResult> {
const result: DeliveryResult = {
attempted: 0,
Expand All @@ -40,7 +40,6 @@ export async function deliverPendingAlerts(
};

const pending = getUndeliveredAlerts(db, network);

if (pending.length === 0) return result;

logger.debug(`Dispatcher: ${pending.length} undelivered alert(s) for network ${network}`);
Expand All @@ -62,40 +61,36 @@ export async function deliverPendingAlerts(
});

try {
await route(alert.channelType, alert.channelTarget, event, alert.webhookSecret);
const channel = channels[alert.channelType];
if (!channel) throw new Error(`Unknown channel type: ${alert.channelType}`);
await channel.send(alert.channelTarget, event, alert.webhookSecret);
markAlertDelivered(db, alert.alertFiredId);
result.delivered++;

logger.info(
`Alert delivered — id: ${alert.alertFiredId}, ` +
`channel: ${alert.channelType}, contract: ${alert.contractId}`,
`Alert delivered — id: ${alert.alertFiredId}, channel: ${alert.channelType}, contract: ${alert.contractId}`,
);
} catch (err: unknown) {
const message = err instanceof Error ? err.message : String(err);
result.failed++;
result.errors.push(message);

incrementRetryCount(db, alert.alertFiredId);
Comment on lines +66 to 76

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift

Separate send failures from post-send persistence failures.

markAlertDelivered() is inside the same try as channel.send(). If the outbound send succeeds and the DB write throws, Line 76 increments retry_count and the same alert will be retried on the next cycle, duplicating the notification. Handle the post-send persistence path separately instead of funneling it into the resend logic.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/alerts/dispatcher.ts` around lines 66 - 76, The send and persistence
steps in dispatcher logic are currently wrapped in the same try/catch, so a
failure in markAlertDelivered() is treated like a channel.send() failure and
causes an unnecessary retry. Update the alert delivery flow in dispatcher.ts to
separate the outbound send path from the post-send database write, using
distinct error handling around channel.send() and markAlertDelivered(). Keep
retry_count increments and retry/error accounting only for send failures, and
handle persistence errors without re-queuing the alert as undelivered.

const nextRetry = alert.retryCount + 1;

if (nextRetry >= MAX_RETRY_COUNT) {
result.abandoned++;
logger.error(
`Alert abandoned after ${MAX_RETRY_COUNT} retries — id: ${alert.alertFiredId}, ` +
`channel: ${alert.channelType}, error: ${message}`,
`Alert abandoned after ${MAX_RETRY_COUNT} retries — id: ${alert.alertFiredId}, channel: ${alert.channelType}, error: ${message}`,
);
} else {
logger.warn(
`Alert delivery failed (attempt ${nextRetry}/${MAX_RETRY_COUNT}) — ` +
`id: ${alert.alertFiredId}, channel: ${alert.channelType}, error: ${message}`,
`Alert delivery failed (attempt ${nextRetry}/${MAX_RETRY_COUNT}) — id: ${alert.alertFiredId}, channel: ${alert.channelType}, error: ${message}`,
);
}
}
}

logger.debug(
`Dispatcher finished — attempted: ${result.attempted}, ` +
`delivered: ${result.delivered}, failed: ${result.failed}, abandoned: ${result.abandoned}`,
`Dispatcher finished — attempted: ${result.attempted}, delivered: ${result.delivered}, failed: ${result.failed}, abandoned: ${result.abandoned}`,
);

return result;
Expand All @@ -106,42 +101,16 @@ export async function deliverSingleAlert(
channelTarget: string,
event: AlertEvent,
webhookSecret?: string | null,
channels: Record<string, AlertChannel> = DEFAULT_CHANNELS,
): Promise<boolean> {
try {
await route(channelType, channelTarget, event, webhookSecret ?? null);
const channel = channels[channelType];
if (!channel) throw new Error(`Unknown channel type: ${channelType}`);
await channel.send(channelTarget, event, webhookSecret ?? null);
return true;
} catch (err: unknown) {
const message = err instanceof Error ? err.message : String(err);
logger.warn(`Resolution alert delivery failed — channel: ${channelType}, error: ${message}`);
return false;
}
}

// ─── Private ─────────────────────────────────────────────────────────────────

async function route(
channelType: string,
channelTarget: string,
event: AlertEvent,
webhookSecret: string | null,
): Promise<void> {
switch (channelType) {
case "webhook":
await sendWebhookAlert(channelTarget, event, webhookSecret);
break;
case "slack":
await sendSlackAlert(channelTarget, event);
break;
case "pagerduty":
await sendPagerDutyAlert(channelTarget, event);
break;
case "discord":
await sendDiscordAlert(channelTarget, event);
break;
case "telegram":
await sendTelegramAlert(channelTarget, event);
break;
default:
throw new Error(`Unknown channel type: ${channelType}`);
}
}
54 changes: 3 additions & 51 deletions src/alerts/types.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,8 @@
import { formatTimeToCloseLedger } from "../utils/formatting.js";

// ─── Core event type ─────────────────────────────────────────────────────────

export type AlertSeverity = "critical" | "warning" | "info";
export type AlertEventType = "threshold_crossed" | "alert_resolved" | "resource_alert" | "state_changed";

// ─── TTL-based alert event ──────────────────────────────────────────────────

export interface TTLAlertEvent {
type: "threshold_crossed" | "alert_resolved";
severity: AlertSeverity;
Expand All @@ -32,8 +28,6 @@ export interface TTLAlertEvent {
timestamp: string;
}

// ─── Resource-based alert event ─────────────────────────────────────────────

export interface ResourceAlertEvent {
type: "resource_alert";
severity: AlertSeverity;
Expand All @@ -57,64 +51,25 @@ export interface ResourceAlertEvent {
timestamp: string;
}

// ─── State-change alert event ───────────────────────────────────────────────
export type AlertEvent = TTLAlertEvent | ResourceAlertEvent;

export interface StateChangeAlertEvent {
type: "state_changed";
severity: AlertSeverity;
contractId: string;
contractName: string | null;
network: string;
entry: {
keyXdr: string;
type: string;
label: string | null;
};
diff: {
diffType: "created" | "updated" | "deleted";
oldValueXdr: string | null;
newValueXdr: string | null;
};
/** Ledger sequence number at the time of detection. */
detectedAtLedger: number;
/** ISO 8601 timestamp. */
timestamp: string;
export interface AlertChannel {
send(target: string, event: AlertEvent, secret?: string | null): Promise<void>;
}

// ─── Union of all alert event types ──────────────────────────────────────────

export type AlertEvent = TTLAlertEvent | ResourceAlertEvent | StateChangeAlertEvent;

// ─── Helpers ─────────────────────────────────────────────────────────────────

/**
* Compute alert severity from remaining TTL.
* - critical: less than 25% of threshold remaining
* - warning: less than threshold (but above 25%)
* - info: used for resolution events
*/
export function computeSeverity(remainingTTL: number, thresholdLedgers: number, isResolution: boolean): AlertSeverity {
if (isResolution) return "info";
if (remainingTTL <= 0) return "critical";
if (remainingTTL < thresholdLedgers * 0.25) return "critical";
return "warning";
}

/**
* Compute alert severity from resource usage percentage.
* - critical: 95% or higher
* - warning: 80-95%
* - info: not used for resource alerts
*/
export function computeResourceSeverity(usagePercent: number): AlertSeverity {
if (usagePercent >= 95) return "critical";
if (usagePercent >= 80) return "warning";
return "info";
}

/**
* Build a TTL-based AlertEvent from raw data.
*/
export function buildAlertEvent(opts: {
type: "threshold_crossed" | "alert_resolved";
contractId: string;
Expand Down Expand Up @@ -148,9 +103,6 @@ export function buildAlertEvent(opts: {
};
}

/**
* Build a resource-based AlertEvent from raw data.
*/
export function buildResourceAlertEvent(opts: {
contractId: string;
contractName: string | null;
Expand Down
Loading
Loading