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
108 changes: 79 additions & 29 deletions README.md

Large diffs are not rendered by default.

8 changes: 8 additions & 0 deletions package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
{
"name": "apiwatch",
"private": true,
"type": "module",
"scripts": {
"test": "node --test test/bodies.test.js test/alerts.test.js"
}
}
176 changes: 144 additions & 32 deletions src/lib/alerts.js
Original file line number Diff line number Diff line change
@@ -1,14 +1,22 @@
/* When an API is broken, and the messages that say so.
*
* Two signals, both measured, neither guessed:
* Three signals, all measured, none guessed:
*
* 5xx an API answered `minErrors` or more 5xx responses within
* the last `window` seconds. Served APIs are counted from the
* server's side, called APIs from the caller's.
* 4xx an API's share of 4xx answers in the last `window` seconds
* is far above its own normal share (`clientErrors:
* "baseline"`), or any 4xx at all (`"all"`). A 404 for a
* missing record is normal traffic for many APIs, so the
* default compares an API with itself rather than with zero.
* down a served API's port stopped listening (the process exited
* or closed it) for `downAfter` seconds. Only ports that have
* served HTTP, or were named with --ports, are watched.
*
* Every 5xx and 4xx message carries the latest failing exchange, request
* and response, through the injected `sample(api, cls)`, which redacts.
*
* Each API is latched: one message when it breaks, one when it recovers
* (no 5xx for `recover` seconds, or the port is back), and a reminder
* every `remind` seconds while it stays broken. Nothing in between.
Expand All @@ -25,10 +33,25 @@ const span = (ms) => {
return m < 90 ? `${m} min` : `${(m / 60).toFixed(1)} h`;
};
const plural = (n, one, many = one + "s") => `${n} ${n === 1 ? one : many}`;
const pct = (x) => `${Math.round(x * 100)}%`;

/* How far above normal a 4xx share must be to count as a spike: three
* times the normal share and ten points above it, so an API that is 2%
* 404s normally needs 12%, and one that is 30% needs 90%. Back to normal
* once it falls under one and a half times, and five points above. */
const spikeLevel = (base) => Math.max(base * 3, base + 0.1);
const normalLevel = (base) => Math.max(base * 1.5, base + 0.05);

/* The span a baseline is learned from, before the current window. */
const BASELINE_MS = 30 * 60_000;
const BASELINE_MIN_REQUESTS = 50;

export class Alerts {
constructor({ apis, listeners, send, log = () => {}, host = "this host", window = 60, minErrors = 1, recover = 120, remind = 1800, downAfter = 5, ports = [], ignore = [] }) {
Object.assign(this, { apis, listeners, send, log, host, window, minErrors, recover, remind, downAfter });
constructor({
apis, listeners, send, log = () => {}, host = "this host", window = 60, minErrors = 1, recover = 120, remind = 1800,
downAfter = 5, ports = [], ignore = [], clientErrors = "baseline", minClient = 5, learnMinutes = 5, sample = () => null,
}) {
Object.assign(this, { apis, listeners, send, log, host, window, minErrors, recover, remind, downAfter, clientErrors, minClient, learnMinutes, sample });
this.state = new Map(); // api key -> { state, since, sentAt, total, missingSince }
this.watchPorts = new Set(ports.map(Number));
this.ignore = new Set(ignore.map((x) => String(x).toLowerCase()));
Expand Down Expand Up @@ -85,36 +108,120 @@ export class Alerts {
s.missingSince = null;
}

/* 5xx */
const recent = this.apis.errorsWithin(api, this.window * 1000, now);
if (s.state === "ok" && recent >= this.minErrors) {
s.state = "failing";
s.since = api.errorTimes[api.errorTimes.length - recent] ?? now;
s.sentAt = now;
s.total = recent;
s.errorsSeen = api.errorTimes.length;
this.emit("failing", row, () => this.failing(api, this.apis.row(api), this.apis.errorsWithin(api, this.window * 1000), Date.now()));
continue;
this.check5xx(api, row, s, now);
this.check4xx(api, row, now);
}
}

check5xx(api, row, s, now) {
const recent = this.apis.errorsWithin(api, this.window * 1000, now);
if (s.state === "ok" && recent >= this.minErrors) {
s.state = "failing";
s.since = api.errorTimes[api.errorTimes.length - recent] ?? now;
s.sentAt = now;
s.total = recent;
s.errorsSeen = api.errorTimes.length;
this.emit("failing", row, () => this.failing(api, this.apis.row(api), this.apis.errorsWithin(api, this.window * 1000), Date.now()));
return;
}
if (s.state !== "failing") return;
s.total += api.errorTimes.length - s.errorsSeen;
s.errorsSeen = api.errorTimes.length;
const last = api.errorTimes[api.errorTimes.length - 1] ?? s.since;
if (now - last >= this.recover * 1000) {
s.state = "ok";
this.emit("recovered", row, {
title: `${row.name} recovered`,
body: `No 5xx from *${row.name}* for ${span(now - last)}. It returned ${plural(s.total, "5xx response")} over ${span(last - s.since)}, starting ${clock(s.since)}.`,
});
} else if (now - s.sentAt >= this.remind * 1000) {
s.sentAt = now;
this.emit("reminder", row, {
title: `${row.name} is still returning 5xx`,
body: `*${row.name}* has returned ${plural(s.total, "5xx response")} since ${clock(s.since)} (${span(now - s.since)}), ${recent} in the last ${this.window}s.`,
});
}
}

/* 4xx, against the API's own normal share ("baseline") or any at all
* ("all"). Its state is kept apart from the 5xx state, so an API can be
* both failing and drowning in 401s, and each recovers on its own. */
check4xx(api, row, now) {
if (this.clientErrors === "off") return;
const s = this.st(`${api.key}#4xx`);
const w = this.window * 1000;
const cur = this.apis.counts(api, now - w, now + 1);
const share = cur.req ? cur.c4 / cur.req : 0;

if (this.clientErrors === "all") {
const last = api.recent4xx[api.recent4xx.length - 1]?.at ?? 0;
if (s.state === "ok" && cur.c4 >= this.minErrors && now - last < w) {
Object.assign(s, { state: "failing", since: now, sentAt: now, last });
this.emit("failing_4xx", row, () => this.failing4xx(api, this.apis.row(api), null, Date.now()));
} else if (s.state === "failing" && now - last >= this.recover * 1000) {
s.state = "ok";
this.emit("recovered_4xx", row, {
title: `${row.name} stopped returning 4xx`,
body: `No 4xx from *${row.name}* for ${span(now - last)}, since it started at ${clock(s.since)}.`,
});
}
if (s.state === "failing") {
s.total += api.errorTimes.length - s.errorsSeen;
s.errorsSeen = api.errorTimes.length;
const last = api.errorTimes[api.errorTimes.length - 1] ?? s.since;
if (now - last >= this.recover * 1000) {
s.state = "ok";
this.emit("recovered", row, {
title: `${row.name} recovered`,
body: `No 5xx from *${row.name}* for ${span(now - last)}. It returned ${plural(s.total, "5xx response")} over ${span(last - s.since)}, starting ${clock(s.since)}.`,
});
} else if (now - s.sentAt >= this.remind * 1000) {
s.sentAt = now;
this.emit("reminder", row, {
title: `${row.name} is still returning 5xx`,
body: `*${row.name}* has returned ${plural(s.total, "5xx response")} since ${clock(s.since)} (${span(now - s.since)}), ${recent} in the last ${this.window}s.`,
});
}
return;
}

/* baseline: armed once the API has been watched long enough to know
* what normal is, and only against requests outside this window. */
if (s.state === "ok") {
if (now - api.firstAt < this.learnMinutes * 60_000) return;
const base = this.apis.counts(api, now - BASELINE_MS, now - w);
if (base.req < BASELINE_MIN_REQUESTS) return;
const baseShare = base.c4 / base.req;
if (cur.c4 >= this.minClient && share >= spikeLevel(baseShare)) {
Object.assign(s, { state: "failing", since: now, sentAt: now, lastSpike: now, baseline: baseShare });
this.emit("failing_4xx", row, () => this.failing4xx(api, this.apis.row(api), baseShare, Date.now()));
}
return;
}
/* The baseline is frozen while it is firing: a long spike would
* otherwise teach the window that the spike is normal. */
if (share > normalLevel(s.baseline)) s.lastSpike = now;
if (now - s.lastSpike >= this.recover * 1000) {
s.state = "ok";
this.emit("recovered_4xx", row, {
title: `${row.name} 4xx back to normal`,
body: `*${row.name}* answered ${pct(share)} of the last ${this.window}s with a 4xx, against ${pct(s.baseline)} normally. The spike started ${clock(s.since)} and lasted ${span(s.lastSpike - s.since)}.`,
});
} else if (now - s.sentAt >= this.remind * 1000) {
s.sentAt = now;
this.emit("reminder_4xx", row, {
title: `${row.name} 4xx still above normal`,
body: `*${row.name}* answered ${pct(share)} of the last ${this.window}s with a 4xx, against ${pct(s.baseline)} normally, since ${clock(s.since)} (${span(now - s.since)}).`,
});
}
}

failing4xx(api, row, baseShare, now) {
const w = this.window * 1000;
const cur = this.apis.counts(api, now - w, now + 1);
const errs = api.recent4xx.filter((e) => e.at >= now - w);
const codes = [...new Set(errs.map((e) => e.status))].sort();
const byEndpoint = new Map();
for (const e of errs) {
const k = `${e.method} ${e.path} → ${e.status}`;
byEndpoint.set(k, (byEndpoint.get(k) ?? 0) + 1);
}
const lines = [...byEndpoint].sort((a, b) => b[1] - a[1]).slice(0, 3).map(([k, n]) => `\`${k}\` ×${n}`);
const share = cur.req ? cur.c4 / cur.req : 0;
const who = api.kind === "served"
? `*${row.name}* (port ${api.port} on ${this.host}) answered ${cur.c4} of ${cur.req} requests with a 4xx in the last ${this.window}s`
: `*${row.name}* returned a 4xx to ${cur.c4} of ${cur.req} calls from ${row.process ?? this.host} in the last ${this.window}s`;
const vs = baseShare == null ? "." : ` (${pct(share)}), against ${pct(baseShare)} normally.`;
return {
title: baseShare == null
? `${row.name} is returning ${codes.length ? codes.join(" and ") : "4xx"}`
: `${row.name} 4xx jumped to ${pct(share)}${codes.length ? ` (${codes.join(", ")})` : ""}`,
body: [who + vs, lines.length ? `Latest (sampled): ${lines.join(", ")}` : null].filter(Boolean).join("\n"),
detail: this.sample(api, "4xx"),
};
}

failing(api, row, recent, now) {
Expand All @@ -138,6 +245,7 @@ export class Alerts {
return {
title: `${row.name} is returning ${codes.length ? codes.join(" and ") : "5xx"}`,
body: [body, lines.length ? `Latest: ${lines.join(", ")}` : null, related].filter(Boolean).join("\n"),
detail: this.sample(api, "5xx"),
};
}

Expand Down Expand Up @@ -167,15 +275,19 @@ export class Alerts {
emit(event, row, content) {
const at = Date.now();
const render = () => {
const { title, body } = typeof content === "function" ? content() : content;
const { title, body, detail = null } = typeof content === "function" ? content() : content;
/* Slack refuses a section over 3000 characters. */
const fit = (t) => (t.length > 2900 ? `${t.slice(0, 2900)}…` : t);
return {
event,
api: row.name,
title,
detail,
text: `${title}: ${body.replace(/[*`]/g, "")}`,
blocks: [
{ type: "header", text: { type: "plain_text", text: title.slice(0, 150) } },
{ type: "section", text: { type: "mrkdwn", text: body } },
{ type: "section", text: { type: "mrkdwn", text: fit(body) } },
...(detail ? [{ type: "section", text: { type: "mrkdwn", text: fit(detail) } }] : []),
{ type: "context", elements: [{ type: "mrkdwn", text: `${this.host} · ${clock(at)} · apiwatch on yeet` }] },
],
};
Expand Down
58 changes: 51 additions & 7 deletions src/lib/apis.js
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,12 @@ const MAX_RECENT = 8;
const MAX_CALLERS = 256;
const RING = 4096; // 5xx timestamps kept per API for the sliding window

/* Request and 4xx counts in 10-second buckets, about 32 minutes of them:
* enough to learn what share of an API's answers are normally 4xx, which
* a ring of timestamps cannot hold for a busy API. */
export const BUCKET_MS = 10_000;
const BUCKETS = 192;

const hostOnly = (h) => {
const s = String(h ?? "").trim().toLowerCase();
if (s.startsWith("[")) return s.slice(1, s.indexOf("]"));
Expand All @@ -44,8 +50,10 @@ const addCaller = (api, c) => {
const isIpLiteral = (h) => /^\d+\.\d+\.\d+\.\d+$/.test(h) || h.includes(":");

export class Apis {
constructor({ name = (pid) => `pid ${pid}`, info = null, isLocal = () => false, listeners = () => new Map(), callerOf = null } = {}) {
constructor({ name = (pid) => `pid ${pid}`, info = null, isLocal = () => false, listeners = () => new Map(), callerOf = null, clientCodes = null } = {}) {
this.name = name;
/* Which 4xx statuses count as client errors; null means all of them. */
this.clientCodes = clientCodes && clientCodes.size ? clientCodes : null;
this.info = info;
this.callerOf = callerOf;
this.isLocal = isLocal;
Expand Down Expand Up @@ -76,6 +84,9 @@ export class Apis {
errorTimes: [],
reqTimes: [],
recentErrors: [],
recent4xx: [],
buckets: [],
lastTx: {},
...init,
};
this.byKey.set(key, api);
Expand Down Expand Up @@ -131,13 +142,46 @@ export class Apis {
api.endpoints.set(epKey, ep);

const error = status != null && status >= 500;
if (error) {
api.errorTimes.push(Date.now());
if (api.errorTimes.length > RING) api.errorTimes.splice(0, api.errorTimes.length - RING);
api.recentErrors.push({ at: Date.now(), method: tx.method, path, status, reason: tx.reason ?? null, pid: tx.pid });
if (api.recentErrors.length > MAX_RECENT) api.recentErrors.shift();
const clientError = status != null && status >= 400 && status < 500 && (!this.clientCodes || this.clientCodes.has(status));
this.bucket(api, api.lastAt, clientError);
if (error || clientError) {
const cls = error ? "5xx" : "4xx";
if (error) {
api.errorTimes.push(api.lastAt);
if (api.errorTimes.length > RING) api.errorTimes.splice(0, api.errorTimes.length - RING);
}
const recent = error ? api.recentErrors : api.recent4xx;
recent.push({ at: api.lastAt, cls, method: tx.method, path, status, reason: tx.reason ?? null, pid: tx.pid });
if (recent.length > MAX_RECENT) recent.shift();
/* The whole exchange of the latest failure of each class, kept so an
* alert can show its bodies. One per class, so at most two bodies
* per API stay in memory. */
api.lastTx[cls] = tx;
}
return { api, error, clientError };
}

bucket(api, at, clientError) {
const t = at - (at % BUCKET_MS);
let b = api.buckets[api.buckets.length - 1];
if (!b || b.t !== t) {
api.buckets.push((b = { t, req: 0, c4: 0 }));
if (api.buckets.length > BUCKETS) api.buckets.shift();
}
b.req++;
if (clientError) b.c4++;
}

/** Requests and 4xx answers in [from, to). */
counts(api, from, to) {
let req = 0;
let c4 = 0;
for (const b of api.buckets) {
if (b.t < from || b.t >= to) continue;
req += b.req;
c4 += b.c4;
}
return { api, error };
return { req, c4 };
}

/* The pid holding the client end of a loopback connection, from the
Expand Down
Loading
Loading