monitor.ts

  1import { DurableObject } from "cloudflare:workers";
  2import { sendAlert, type AlertKind } from "./email";
  3import { pushoverEnabled, sendPushover } from "./pushover";
  4
  5export const REGIONS = ["wnam", "enam", "weur", "eeur", "apac"] as const;
  6export type Region = (typeof REGIONS)[number];
  7
  8export const REGION_NAMES: Record<Region, string> = {
  9  wnam: "Western North America",
 10  enam: "Eastern North America",
 11  weur: "Western Europe",
 12  eeur: "Eastern Europe",
 13  apac: "Asia Pacific",
 14};
 15
 16/** One health check result. status 0 means no HTTP response (timeout / network error). */
 17export interface Check {
 18  region: Region;
 19  endpoint: string;
 20  ts: number;
 21  status: number;
 22  duration_ms: number;
 23  body: string | null;
 24  error: string | null;
 25}
 26
 27export interface RegionSnapshot {
 28  region: Region;
 29  nextAlarm: number | null;
 30  checks: Check[];
 31}
 32
 33const RETENTION_MS = 14 * 24 * 60 * 60 * 1000;
 34const MAX_BODY_CHARS = 64 * 1024;
 35
 36type CheckRow = Omit<Check, "region">;
 37type AlertRow = {
 38  endpoint: string;
 39  failing_since: number | null;
 40  last_alert_ts: number | null;
 41  last_pushover_ts: number | null;
 42};
 43
 44/** The alert channels. Each keeps its own last-sent time, so one failing never holds up the other. */
 45const CHANNELS = [
 46  { column: "last_alert_ts", enabled: (_env: Env) => true, send: sendAlert },
 47  { column: "last_pushover_ts", enabled: pushoverEnabled, send: sendPushover },
 48] as const;
 49
 50export class MonitorDO extends DurableObject<Env> {
 51  private sql: SqlStorage;
 52
 53  constructor(ctx: DurableObjectState, env: Env) {
 54    super(ctx, env);
 55    this.sql = ctx.storage.sql;
 56    this.sql.exec(`
 57      CREATE TABLE IF NOT EXISTS checks (
 58        id INTEGER PRIMARY KEY,
 59        endpoint TEXT NOT NULL,
 60        ts INTEGER NOT NULL,
 61        status INTEGER NOT NULL,
 62        duration_ms INTEGER NOT NULL,
 63        body TEXT,
 64        error TEXT
 65      );
 66      CREATE INDEX IF NOT EXISTS checks_endpoint_ts ON checks (endpoint, ts);
 67      CREATE INDEX IF NOT EXISTS checks_ts ON checks (ts);
 68      CREATE TRIGGER IF NOT EXISTS checks_retention AFTER INSERT ON checks BEGIN
 69        DELETE FROM checks WHERE ts < NEW.ts - ${RETENTION_MS};
 70      END;
 71      CREATE TABLE IF NOT EXISTS alert_state (
 72        endpoint TEXT PRIMARY KEY,
 73        failing_since INTEGER,
 74        last_alert_ts INTEGER
 75      );
 76    `);
 77    const columns = this.sql.exec<{ name: string }>(`SELECT name FROM pragma_table_info('alert_state')`).toArray();
 78    if (!columns.some((c) => c.name === "last_pushover_ts")) {
 79      this.sql.exec(`ALTER TABLE alert_state ADD COLUMN last_pushover_ts INTEGER`);
 80    }
 81  }
 82
 83  /** The region this instance monitors from; recorded by ensureAlarm(). */
 84  private get region(): Region {
 85    return (this.ctx.storage.kv.get<Region>("region") ?? "wnam") as Region;
 86  }
 87
 88  /** Records this instance's region and starts the check loop if it isn't running. */
 89  async ensureAlarm(region: Region): Promise<void> {
 90    this.ctx.storage.kv.put("region", region);
 91    if ((await this.ctx.storage.getAlarm()) === null) {
 92      await this.ctx.storage.setAlarm(Date.now());
 93    }
 94  }
 95
 96  async alarm(): Promise<void> {
 97    // Schedule the next run first so a failure below can never stop the loop.
 98    await this.ctx.storage.setAlarm(Date.now() + Number(this.env.CHECK_INTERVAL_SECONDS) * 1000);
 99    await Promise.all(endpoints(this.env).map((url) => this.check(url)));
100  }
101
102  /** Latest check for each configured endpoint. Also bootstraps the alarm loop. */
103  async latest(region: Region): Promise<RegionSnapshot> {
104    await this.ensureAlarm(region);
105    const checks: Check[] = [];
106    for (const url of endpoints(this.env)) {
107      const row = this.sql
108        .exec<CheckRow>(`SELECT endpoint, ts, status, duration_ms, body, error FROM checks
109                         WHERE endpoint = ? ORDER BY ts DESC LIMIT 1`, url)
110        .toArray()[0];
111      if (row) checks.push({ region, ...row });
112    }
113    return { region, nextAlarm: await this.ctx.storage.getAlarm(), checks };
114  }
115
116  /** Checks for one endpoint older than `before`, newest first. */
117  async history(endpoint: string, before: number, limit: number, failuresOnly: boolean): Promise<Check[]> {
118    const region = this.region;
119    return this.sql
120      .exec<CheckRow>(
121        `SELECT endpoint, ts, status, duration_ms, body, error FROM checks
122         WHERE endpoint = ? AND ts < ? ${failuresOnly ? "AND status != 200" : ""}
123         ORDER BY ts DESC LIMIT ?`,
124        endpoint, before, limit,
125      )
126      .toArray()
127      .map((row) => ({ region, ...row }));
128  }
129
130  private async check(url: string): Promise<void> {
131    const ts = Date.now();
132    let status = 0;
133    let body: string | null = null;
134    let error: string | null = null;
135    try {
136      const res = await fetch(url, {
137        signal: AbortSignal.timeout(Number(this.env.TIMEOUT_MS)),
138        headers: { accept: "application/json", "user-agent": "rtw-monitor" },
139        cache: "no-store",
140      });
141      status = res.status;
142      body = (await res.text()).slice(0, MAX_BODY_CHARS);
143    } catch (e) {
144      error = e instanceof Error && e.name === "TimeoutError"
145        ? `Timed out after ${this.env.TIMEOUT_MS} ms`
146        : String(e);
147    }
148    const check: Check = { region: this.region, endpoint: url, ts, status, duration_ms: Date.now() - ts, body, error };
149    this.sql.exec(
150      `INSERT INTO checks (endpoint, ts, status, duration_ms, body, error) VALUES (?, ?, ?, ?, ?, ?)`,
151      url, ts, status, check.duration_ms, body, error,
152    );
153    await this.updateAlertState(check);
154  }
155
156  /**
157   * Alerts on every enabled channel (email; Pushover once its secrets are set) on healthy→failing,
158   * every REMINDER_HOURS while still failing, and on recovery. A failed send leaves that channel's
159   * last-sent time unchanged so the next check retries it, on that channel only.
160   */
161  private async updateAlertState(check: Check): Promise<void> {
162    const state = this.sql
163      .exec<AlertRow>(`SELECT * FROM alert_state WHERE endpoint = ?`, check.endpoint)
164      .toArray()[0] ?? { endpoint: check.endpoint, failing_since: null, last_alert_ts: null, last_pushover_ts: null };
165    const reminderMs = Number(this.env.REMINDER_HOURS) * 60 * 60 * 1000;
166    const failing = check.status !== 200;
167    const recovered = !failing && state.failing_since !== null;
168    if (failing && state.failing_since === null) state.failing_since = check.ts;
169
170    await Promise.all(
171      CHANNELS.filter((ch) => ch.enabled(this.env)).map(async (ch) => {
172        const last = state[ch.column];
173        let kind: AlertKind | null = null;
174        if (recovered) kind = "recovered";
175        else if (failing && last === null) kind = "down";
176        else if (failing && last !== null && check.ts - last >= reminderMs) kind = "still down";
177        if (!kind) return;
178        const sent = await ch.send(this.env, kind, check, state.failing_since);
179        if (sent && kind !== "recovered") state[ch.column] = check.ts;
180      }),
181    );
182    if (recovered) {
183      state.failing_since = null;
184      state.last_alert_ts = null;
185      state.last_pushover_ts = null;
186    }
187
188    this.sql.exec(
189      `INSERT OR REPLACE INTO alert_state (endpoint, failing_since, last_alert_ts, last_pushover_ts) VALUES (?, ?, ?, ?)`,
190      check.endpoint, state.failing_since, state.last_alert_ts, state.last_pushover_ts,
191    );
192  }
193}
194
195export function endpoints(env: Env): string[] {
196  return env.ENDPOINTS as unknown as string[];
197}