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});