monitor.test.ts
1import { env, runDurableObjectAlarm, runInDurableObject } from "cloudflare:test";
2import { afterEach, describe, expect, it, vi } from "vitest";
3import type { MonitorDO, Region } from "../src/monitor";
4
5const A = "https://a.test/healthz";
6const B = "https://b.test/healthz";
7
8/**
9 * A fresh, isolated monitor for one test. The alarm is armed in the future so it never fires on its
10 * own; runCycle() triggers it explicitly.
11 */
12async function freshMonitor(region: Region): Promise<DurableObjectStub<MonitorDO>> {
13 const ns = env.MONITOR as DurableObjectNamespace<MonitorDO>;
14 const stub = ns.get(ns.newUniqueId());
15 await runInDurableObject(stub, async (_i, state) => {
16 state.storage.kv.put("region", region);
17 await state.storage.setAlarm(Date.now() + 3600_000);
18 });
19 return stub;
20}
21
22const PUSHOVER = "https://api.pushover.net/1/messages.json";
23
24/**
25 * Routes outbound fetches: each endpoint maps to a status code, or an Error to throw. Pushover
26 * posts are recorded in `pushes`; they succeed unless `pushoverFails` is set.
27 */
28function mockEndpoints(responses: Record<string, number | Error>, pushes: URLSearchParams[] = [], pushoverFails = false) {
29 return vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => {
30 const url = input instanceof Request ? input.url : String(input);
31 if (url === PUSHOVER) {
32 if (pushoverFails) return new Response('{"status":0,"errors":["down"]}', { status: 500 });
33 pushes.push(new URLSearchParams(String(init?.body)));
34 return new Response('{"status":1,"request":"x"}');
35 }
36 const r = responses[url];
37 if (r instanceof Error) throw r;
38 return new Response(JSON.stringify({ ok: r === 200, url }), { status: r ?? 404 });
39 });
40}
41
42/** Turns the Pushover channel on for this monitor by giving its env the two secrets. */
43async function enablePushover(stub: DurableObjectStub<MonitorDO>) {
44 await runInDurableObject(stub, (instance: MonitorDO) => {
45 Object.assign((instance as unknown as { env: Env }).env, { PUSHOVER_TOKEN: "tok", PUSHOVER_USER: "usr" });
46 });
47}
48
49/** Runs one alarm cycle and returns the subjects of emails sent during it. */
50async function runCycle(stub: DurableObjectStub<MonitorDO>, send?: (msg: EmailMessage) => Promise<unknown>) {
51 const subjects: string[] = [];
52 const log = vi.spyOn(console, "log").mockImplementation(() => {});
53 const error = vi.spyOn(console, "error").mockImplementation(() => {});
54 await runInDurableObject(stub, (instance: MonitorDO) => {
55 const emailEnv = (instance as unknown as { env: Env }).env;
56 vi.spyOn(emailEnv.EMAIL, "send").mockImplementation(async (msg: any) => {
57 expect(msg).toMatchObject({ from: "[email protected]", to: "[email protected]" });
58 if (send) await send(msg);
59 return { messageId: "x" } as any;
60 });
61 });
62 expect(await runDurableObjectAlarm(stub)).toBe(true);
63 for (const [message, data] of log.mock.calls) if (message === "alert email sent") subjects.push(data.subject);
64 log.mockRestore();
65 error.mockRestore();
66 return subjects;
67}
68
69function sql<T extends Record<string, SqlStorageValue>>(stub: DurableObjectStub<MonitorDO>, query: string, ...args: SqlStorageValue[]) {
70 return runInDurableObject(stub, (_i, state) => state.storage.sql.exec<T>(query, ...args).toArray());
71}
72
73afterEach(() => vi.restoreAllMocks());
74
75describe("checks", () => {
76 it("records results and reschedules the alarm", async () => {
77 const stub = await freshMonitor("weur");
78 mockEndpoints({ [A]: 200, [B]: new Error("connection refused") });
79 const subjects = await runCycle(stub);
80
81 const snap = await stub.latest("weur");
82 const a = snap.checks.find((c) => c.endpoint === A)!;
83 const b = snap.checks.find((c) => c.endpoint === B)!;
84 expect(a).toMatchObject({ region: "weur", status: 200, error: null });
85 expect(JSON.parse(a.body!)).toEqual({ ok: true, url: A });
86 expect(b).toMatchObject({ status: 0, body: null });
87 expect(b.error).toContain("connection refused");
88 expect(snap.nextAlarm).toBeGreaterThan(Date.now());
89 expect(subjects).toEqual([expect.stringMatching(/^\[monitor\] DOWN: b\.test\/healthz \(Error: connection refused\) from weur$/)]);
90 });
91
92 it("reports timeouts", async () => {
93 const stub = await freshMonitor("apac");
94 const timeout = new DOMException("timed out", "TimeoutError");
95 mockEndpoints({ [A]: timeout, [B]: 200 });
96 await runCycle(stub);
97 const a = (await stub.latest("apac")).checks.find((c) => c.endpoint === A)!;
98 expect(a.status).toBe(0);
99 expect(a.error).toMatch(/^Timed out after \d+ ms$/);
100 });
101
102 it("deletes rows older than 14 days on insert", async () => {
103 const stub = await freshMonitor("eeur");
104 const old = Date.now() - 15 * 24 * 3600 * 1000;
105 const recent = Date.now() - 13 * 24 * 3600 * 1000;
106 await sql(stub, "INSERT INTO checks (endpoint, ts, status, duration_ms) VALUES (?, ?, 200, 1), (?, ?, 200, 1)", A, old, A, recent);
107 await sql(stub, "INSERT INTO checks (endpoint, ts, status, duration_ms) VALUES (?, ?, 200, 1)", A, Date.now());
108 const rows = await sql<{ ts: number }>(stub, "SELECT ts FROM checks ORDER BY ts");
109 expect(rows.map((r) => r.ts)).not.toContain(old);
110 expect(rows).toHaveLength(2);
111 });
112});
113
114describe("alerts", () => {
115 it("emails on failure, stays quiet, reminds, then emails on recovery", async () => {
116 const stub = await freshMonitor("enam");
117
118 mockEndpoints({ [A]: 503, [B]: 200 });
119 expect(await runCycle(stub)).toEqual([expect.stringContaining("DOWN: a.test/healthz (HTTP 503) from enam")]);
120 expect(await runCycle(stub)).toEqual([]);
121
122 // Pretend the last alert went out longer ago than REMINDER_HOURS.
123 await sql(stub, "UPDATE alert_state SET last_alert_ts = last_alert_ts - ? WHERE endpoint = ?", 7 * 3600 * 1000, A);
124 expect(await runCycle(stub)).toEqual([expect.stringContaining("STILL DOWN: a.test/healthz")]);
125
126 vi.restoreAllMocks();
127 mockEndpoints({ [A]: 200, [B]: 200 });
128 expect(await runCycle(stub)).toEqual([expect.stringContaining("RECOVERED: a.test/healthz (HTTP 200)")]);
129 expect(await runCycle(stub)).toEqual([]);
130 });
131
132 it("retries the alert on the next check if sending fails", async () => {
133 const stub = await freshMonitor("wnam");
134 mockEndpoints({ [A]: 500, [B]: 200 });
135 await runCycle(stub, async () => {
136 throw new Error("email service unavailable");
137 });
138 expect(await runCycle(stub)).toEqual([expect.stringContaining("DOWN: a.test/healthz")]);
139 expect(await runCycle(stub)).toEqual([]);
140 });
141});
142
143describe("pushover", () => {
144 const titles = (pushes: URLSearchParams[]) => pushes.map((p) => p.get("title"));
145
146 it("stays off until both secrets are set", async () => {
147 const stub = await freshMonitor("weur");
148 const pushes: URLSearchParams[] = [];
149 mockEndpoints({ [A]: 503, [B]: 200 }, pushes);
150 expect(await runCycle(stub)).toHaveLength(1);
151 expect(pushes).toEqual([]);
152 });
153
154 it("notifies alongside email: down with the failing checks, reminder, recovery", async () => {
155 const stub = await freshMonitor("enam");
156 await enablePushover(stub);
157 const pushes: URLSearchParams[] = [];
158 const body = JSON.stringify({
159 status: "degraded",
160 checks: { disk: { status: "ok", detail: "used=18%" }, backup: { status: "fail", detail: "age=27h", error: "backup stale" } },
161 });
162 vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => {
163 const url = input instanceof Request ? input.url : String(input);
164 if (url === PUSHOVER) {
165 pushes.push(new URLSearchParams(String(init?.body)));
166 return new Response('{"status":1}');
167 }
168 return new Response(url === A ? body : "{}", { status: url === A ? 503 : 200 });
169 });
170
171 expect(await runCycle(stub)).toHaveLength(1);
172 expect(pushes).toHaveLength(1);
173 expect(pushes[0].get("title")).toBe("DOWN: a.test/healthz (HTTP 503) from enam");
174 expect(pushes[0].get("message")).toMatch(/^Failing since .*\nbackup: backup stale$/);
175 expect(pushes[0].get("priority")).toBe("0");
176 expect(pushes[0].get("token")).toBe("tok");
177 expect(pushes[0].get("user")).toBe("usr");
178 expect(pushes[0].get("url")).toBe(`https://monitor.rtw.run/endpoint?url=${encodeURIComponent(A)}`);
179
180 await runCycle(stub);
181 expect(pushes).toHaveLength(1);
182
183 await sql(stub, "UPDATE alert_state SET last_alert_ts = last_alert_ts - ?, last_pushover_ts = last_pushover_ts - ? WHERE endpoint = ?",
184 7 * 3600 * 1000, 7 * 3600 * 1000, A);
185 await runCycle(stub);
186 expect(titles(pushes).at(-1)).toMatch(/^STILL DOWN: a\.test/);
187
188 vi.restoreAllMocks();
189 mockEndpoints({ [A]: 200, [B]: 200 }, pushes);
190 expect(await runCycle(stub)).toEqual([expect.stringContaining("RECOVERED")]);
191 expect(titles(pushes).at(-1)).toBe("RECOVERED: a.test/healthz (HTTP 200) from enam");
192 expect(pushes.at(-1)!.get("priority")).toBe("-1");
193 expect(pushes).toHaveLength(3);
194 });
195
196 it("retries only the channel that failed", async () => {
197 // Email down: Pushover goes out once, email is retried until it works.
198 const stub = await freshMonitor("wnam");
199 await enablePushover(stub);
200 const pushes: URLSearchParams[] = [];
201 mockEndpoints({ [A]: 500, [B]: 200 }, pushes);
202 expect(await runCycle(stub, async () => { throw new Error("email service unavailable"); })).toEqual([]);
203 expect(titles(pushes)).toEqual([expect.stringMatching(/^DOWN: a\.test/)]);
204 expect(await runCycle(stub)).toEqual([expect.stringContaining("DOWN: a.test/healthz")]);
205 expect(pushes).toHaveLength(1);
206
207 // Pushover down: email goes out once, Pushover is retried until it works.
208 const other = await freshMonitor("apac");
209 await enablePushover(other);
210 vi.restoreAllMocks();
211 mockEndpoints({ [A]: 500, [B]: 200 }, pushes, true);
212 expect(await runCycle(other)).toHaveLength(1);
213 vi.restoreAllMocks();
214 const retried: URLSearchParams[] = [];
215 mockEndpoints({ [A]: 500, [B]: 200 }, retried);
216 expect(await runCycle(other)).toEqual([]);
217 expect(titles(retried)).toEqual([expect.stringMatching(/^DOWN: a\.test.* from apac$/)]);
218 });
219});
220
221describe("history", () => {
222 it("pages newest first and filters non-200", async () => {
223 const stub = await freshMonitor("weur");
224 const now = Date.now();
225 for (let i = 0; i < 10; i++) {
226 await sql(stub, "INSERT INTO checks (endpoint, ts, status, duration_ms) VALUES (?, ?, ?, 1)", A, now - i * 1000, i % 3 ? 200 : 503);
227 }
228 const page = await stub.history(A, Number.MAX_SAFE_INTEGER, 4, false);
229 expect(page.map((c) => c.ts)).toEqual([now, now - 1000, now - 2000, now - 3000]);
230 const next = await stub.history(A, page.at(-1)!.ts, 4, false);
231 expect(next[0].ts).toBe(now - 4000);
232 const failures = await stub.history(A, Number.MAX_SAFE_INTEGER, 50, true);
233 expect(failures.map((c) => c.status)).toEqual([503, 503, 503, 503]);
234 expect(failures.every((c) => c.region === "weur")).toBe(true);
235 });
236});