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}