diff --git a/README.md b/README.md index 58b7941..deb2400 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ -**`apiwatch` is an eBPF API watchdog for Linux: it lists every HTTP API a machine serves and calls, from live traffic, and messages Slack when one starts failing.** +**`apiwatch` is an eBPF API watchdog for Linux: it lists every HTTP API a machine serves and calls, from live traffic, and messages Slack with the failing request when one breaks.** ## Quick start @@ -66,6 +66,10 @@ The thresholds `--watch` uses, and why each default is what it is: | `--down-after` | `5` | Seconds a served port must be gone before "stopped listening" fires, so a restart doesn't page you. | | `--ignore` | none | APIs never to alert on, by name, host or port: `--ignore httpbin.org,8082`. | | `--ports` | all | Capture only these ports, filtered in the kernel. Also watches them for going down from the first second, before any traffic is seen. | +| `--client-errors` | `baseline` | 4xx alerts. `baseline` learns each API's normal 4xx share over its first five minutes and alerts when the last minute is at least three times that and ten points higher, because a 404 for a missing record is normal traffic for many APIs. `all` alerts on any 4xx like a 5xx; `off` never. | +| `--client-codes` | every 4xx | Count only these statuses as 4xx, e.g. `--client-codes 401,403,429` to watch auth failures and rate limiting and ignore 404s. | +| `--min-client-errors` | `5` | 4xx answers within the window before a jump can fire, so three 404s on a quiet API are not a spike. | +| `--bodies` | `redacted` | The failing request and its response in each 5xx and 4xx alert, about 1000 characters each. `redacted` blanks values under secret-looking keys (password, token, api_key, authorization, card, cvv, ssn and similar) in JSON, forms and query strings, plus bearer tokens, JWTs and card numbers anywhere. `raw` sends them as captured; `off` sends method, path and status only. | ### Slack, once @@ -76,6 +80,8 @@ Messages go out through yeet's own Slack connection rather than a webhook you ma Then `--test-alert` sends one message through `yeet.alert`. If the host isn't signed in, or the post is refused, it prints why and exits non-zero. +Alerts carry the failing request and response, and like every alert they travel through yeet's servers to Slack. Secret fields are blanked by default; other personal data, an email address for instance, is not. Use `--bodies off` if request bodies must not leave the host. + ### Leave it running as a service `--watch` belongs in a [yeet service](https://yeet.cx/docs/cli/services?utm_source=github&utm_medium=readme&utm_campaign=apiwatch), which the daemon keeps running after your shell closes, restarts if it dies, and starts again at boot: @@ -85,10 +91,15 @@ git clone --depth 1 https://github.com/yeet-src/apiwatch ~/.local/share/apiwatch make -C ~/.local/share/apiwatch yeet service new apiwatch -C ~/.local/share/apiwatch -R always yeet service unit add apiwatch/watch -I ~/.local/share/apiwatch/src/main.js -- --watch --slack "#api-alerts" --name "$(hostname)" +yeet service unit add apiwatch/web -W http://127.0.0.1:9470 +yeet service mount apiwatch/web -L /log -t watch -p console yeet service enable apiwatch yeet service start apiwatch +curl -sN http://127.0.0.1:9470/log # the watcher's JSON lines, streamed as they happen ``` +The `web` unit and the `/log` route serve the watcher's console over plain HTTP: a `GET` gets a chunked `text/plain` body, one JSON line per event, open until you disconnect. Routes on a service's web server answer only while the host is signed in, so on a signed-out host `/log` returns `403` with a "Pair this host" page. `127.0.0.1` keeps it on this machine; bind another address only if you mean to share the log. + Build it yourself and point the unit at the built checkout. A unit added straight from `gh:yeet-src/apiwatch` is cloned but never built, so it restarts forever on a missing `bin/socket.bpf.o`. Give `-I` an absolute path, since a relative one resolves against wherever you ran the command. The service runs its own copy of the directory, so a `git pull`, a `make` and a restart still run the old code. To update, stop the service, `yeet service unit remove apiwatch/watch`, add the unit again with the same arguments, and start it. Or create the service with `--dev`, which runs the checkout in place so a restart picks up a rebuild (and breaks if you delete the checkout). `yeet service stop apiwatch` pauses it, `yeet service remove apiwatch` removes it. @@ -142,15 +153,18 @@ Each API is latched: one message when it breaks, one when it recovers (no 5xx fo **What's the lightest way to add 5xx alerting to one Linux box without standing up Prometheus, an exporter and Alertmanager?** `apiwatch` is one process and one yeet service with no config file; the thresholds are flags. On a test box handling about seven HTTP exchanges a second, the script's JavaScript used about 0.3% of one core over five minutes with its heap between 4 and 9 MiB, and the kernel probes ran for 16 ms in two minutes (about 2 µs per send, under 1 µs per receive), not counting the kernel's own cost of firing a kprobe. That cost grows with how much TCP traffic the box carries, not just HTTP, so on a busy machine give it `--ports`. See [What it can't see](#what-it-cant-see). +**An API started answering 401s or 429s and nobody noticed until customers complained. How do I get told when an API's client errors jump, without paging on every 404?** +Leave `--client-errors baseline` on. Each API learns its own normal share of 4xx answers over its first five minutes, and the alert fires when the last minute is at least three times that share and ten points higher, with the request that got the 4xx and the response it got. An API that normally answers 15% 404s alerts at 45%, not at the first 404. `--client-codes 401,403,429` narrows it to auth failures and rate limiting. + +**When an API starts failing, how do I see the actual request that failed and the error it returned, without turning on request logging?** +Every 5xx and 4xx alert carries the latest failing exchange: the method and path, the request body, the status, and the response body, about 1000 characters each, read off the socket. Passwords, tokens, API keys and card numbers are blanked before they leave the host; `--bodies off` drops the bodies. For HTTPS this needs a readable TLS library (see the question above). + **Is this a replacement for Datadog, New Relic, or Prometheus with Alertmanager?** -No. It watches one host, keeps nothing once it restarts, has no dashboards, no history, no latency or 4xx alerts, no on-call routing or escalation, and posts to Slack only. It is for knowing which APIs a machine has and hearing about it when one of them starts returning 5xx, on hosts where a full observability stack isn't there or isn't watching these APIs. +No. It watches one host, keeps nothing once it restarts, has no dashboards, no history, no latency alerts, no on-call routing or escalation, and posts to Slack only. It is for knowing which APIs a machine has and hearing about it when one of them starts returning 5xx, on hosts where a full observability stack isn't there or isn't watching these APIs. **When should I use this instead of an uptime checker like UptimeRobot or Pingdom, Datadog Synthetics, or a deeper tool like httpscope?** Use an uptime checker when the question is "can the outside world reach my site", since `apiwatch` sees only traffic that reached the box. Use synthetics when you need a scripted multi-step check, like a login followed by a checkout. Use [`httpscope`](https://github.com/yeet-src/httpscope) when you want the full shape of each API (request and response schemas, drift between deploys, a GraphQL interface an agent can query), and [`container-traffic`](https://github.com/yeet-src/container-traffic) for a live per-container rate, error and latency dashboard. Use `apiwatch` for the inventory plus a 5xx alert from real traffic, with nothing else to run. -**Can I run this on a server where I can't install a proxy, a sidecar, or an agent into the application?** -Yes. Nothing attaches to your processes' configuration: the kprobes are in the kernel and the uprobes are on the TLS library file, which the daemon attaches from outside. Your apps keep running unchanged, and removing `apiwatch` leaves nothing behind. - ## What you're looking at `--discover` on a test box running nginx in front of two Python services, with a load generator as the clients, a sync job calling a partner API that fails one call in seven, and an order service fetching a weather forecast: @@ -197,42 +211,73 @@ The first line is the window. **Served by this machine** has one entry per liste | `called by` | Who sent the requests: local processes by name, or `remote clients` for anything off the box. A process that exits within milliseconds (a `curl` in a loop) is counted, not named. | | `from shop-partner-sync` | For a called API, the process that made the calls. | -`--watch` writes one JSON line per event. Running as `--watch --dry-run --name shop-box --ignore httpbin.org` with the payments upstream stopped for 26 seconds: +`--watch` writes one JSON line per event. Running as `--watch --dry-run --name shop-box --ignore httpbin.org` while the payments upstream was stopped and a client kept posting orders with a password and a card number in them: ```console -{"t":"2026-10-07T20:36:49.675Z","event":"failing","api":"shop-orders","kind":"served","port":8081,"title":"shop-orders is returning 502"} -{"t":"2026-10-07T20:36:49.675Z","event":"failing","api":"nginx","kind":"served","port":80,"title":"nginx is returning 502"} -{"t":"2026-10-07T20:36:52.673Z","event":"down","api":"shop-payments","kind":"served","port":8082,"title":"shop-payments stopped listening"} -{"t":"2026-10-07T20:37:13.672Z","event":"up","api":"shop-payments","kind":"served","port":8082,"title":"shop-payments is listening again"} -{"t":"2026-10-07T20:39:11.677Z","event":"recovered","api":"shop-orders","kind":"served","port":8081,"title":"shop-orders recovered"} -{"t":"2026-10-07T20:39:11.677Z","event":"recovered","api":"nginx","kind":"served","port":80,"title":"nginx recovered"} +{"t":"2026-10-08T17:26:42.577Z","event":"failing","api":"shop-orders","kind":"served","port":8081,"title":"shop-orders is returning 502"} +{"t":"2026-10-08T17:26:42.577Z","event":"failing","api":"nginx","kind":"served","port":80,"title":"nginx is returning 502"} +{"t":"2026-10-08T17:26:47.575Z","event":"down","api":"shop-payments","kind":"served","port":8082,"title":"shop-payments stopped listening"} ``` -The first three landed within six seconds, so they go out as one message. This is the message as Slack would show it, from the same run: +Those three landed within six seconds, so they go out as one message. This is it as Slack would show it, from the same run: -```text +````text 3 APIs broke on shop-box shop-orders is returning 502 - shop-orders (port 8081 on shop-box) answered 5 of 109 requests with a 5xx in the last 60s. - Latest: POST /orders → 502 ×5 - On this host, shop-payments (port 8082) stopped listening at 20:36:47 UTC. + shop-orders (port 8081 on shop-box) answered 25 of 82 requests with a 5xx in the last 60s. + Latest: POST /orders → 502 ×8 + On this host, shop-payments (port 8082) stopped listening at 17:26:42 UTC. + Latest failing request (secrets redacted) + ```POST /orders + content-type: application/json + + {"amount":1999,"email":"ana@example.com","card_number":"[redacted]","password":"[redacted]"}``` + Response + ```502 Bad Gateway + content-type: application/json + + {"error":"payments unreachable: "}``` nginx is returning 502 - nginx (port 80 on shop-box) answered 7 of 116 requests with a 5xx in the last 60s. - Latest: POST /api/orders → 502 ×5, GET /payments/health → 502 ×2 - On this host, shop-payments (port 8082) stopped listening at 20:36:47 UTC. + nginx (port 80 on shop-box) answered 45 of 105 requests with a 5xx in the last 60s. + Latest: GET /payments/health → 502 ×4, POST /api/orders → 502 ×4 + On this host, shop-payments (port 8082) stopped listening at 17:26:42 UTC. + Latest failing request (secrets redacted) + ```GET /payments/health?region=us&access_token=[redacted]``` + Response + ```502 Bad Gateway + content-type: text/html + + 502 Bad Gateway …``` shop-payments stopped listening Nothing on shop-box is listening on port 8082 any more (it was shop-payments). Every request to it fails until it is back. -shop-box · 20:36:49 UTC · apiwatch on yeet -``` +shop-box · 17:26:42 UTC · apiwatch on yeet +```` + +A 4xx jump, from a Node API that normally answers 15% 404s for users that don't exist, when a client started asking for one that never would: + +````text +users-api 4xx jumped to 52% (404) + users-api (port 3000 on users-box) answered 39 of 75 requests with a 4xx in the last 60s (52%), against 15% normally. + Latest (sampled): GET /users/{n} → 404 ×8 + Latest failing request (secrets redacted) + ```GET /users/99?api_key=[redacted]``` + Response + ```404 Not Found + content-type: application/json + + {"error":"not found"}``` +users-box · 17:28:41 UTC · apiwatch on yeet +```` -Twenty-six seconds later comes "shop-payments is listening again, after 26s down", and two minutes after the last 502 a single "2 APIs recovered on shop-box": "No 5xx from nginx for 2 min. It returned 22 5xx responses over 23s, starting 20:36:48 UTC." +Four minutes later the same API reported "users-api 4xx back to normal: 13% of the last 60s, against 15% normally." | event | when | | --- | --- | | `start` | The watcher started: host label, channel, whether the host is signed in. | -| `status` | 20 s and 60 s after start, then every ten minutes: every API being watched with its request and 5xx counts, whether the host is signed in, which TLS libraries are tapped. | +| `status` | 20 s after start, then every minute: every API being watched with its request and 5xx counts, whether the host is signed in, the `--bodies` and `--client-errors` modes, which TLS libraries are tapped. | | `failing` | An API crossed `--min-errors` 5xx within `--window`. | +| `failing_4xx` / `recovered_4xx` / `reminder_4xx` | An API's 4xx share jumped past its baseline (or, with `--client-errors all`, any 4xx), came back to normal, or is still high after `--remind` seconds. | | `down` / `up` | A served port stopped listening for `--down-after` seconds, and came back. | | `reminder` | Still failing after `--remind` seconds. | | `recovered` | No 5xx for `--recover` seconds. | @@ -244,7 +289,7 @@ Twenty-six seconds later comes "shop-payments is listening again, after 26s down `apiwatch` never draws a screen, so it is safe to pipe, redirect, and run from an agent or a CI job. - `--discover` prints plain text and exits after `--seconds`. `--discover --json` prints one JSON object with `served`, `called`, `unreadable`, `quiet` and `tls` arrays and the same fields as the text. -- `--watch` prints JSON lines on stdout until stopped. As a service, read them with `yeet attach -c `, the id from `yeet service tree apiwatch`. Attaching shows only lines printed after you attach, which is why `status` repeats: attach within a minute of starting it and you see one. +- `--watch` prints JSON lines on stdout until stopped. As a service with the `/log` route, `curl -sN http://127.0.0.1:9470/log` streams them; without the route, `yeet attach -c ` does, the id from `yeet service tree apiwatch`. Either way you see only lines printed after you connect, which is why `status` repeats every minute: connect at any time and one arrives within a minute. - `--test-alert` exits non-zero and prints the reason when the host isn't signed in or Slack refuses the post, so it works as a check in a script. To verify an install, run `--discover --seconds 30` with something generating HTTP, as in [Try it without real traffic](#try-it-without-real-traffic), and look for that port in the served list. @@ -265,7 +310,8 @@ src/ main.js flags, the three modes, output, the Slack queue lib/capture.js loads the taps, keeps the socket inventory, names TLS calls by SNI lib/apis.js transactions → served and called APIs, endpoints, status counts - lib/alerts.js the 5xx window, port-down, latching, recovery and reminders + lib/alerts.js 5xx, the 4xx baseline, port-down, latching, recovery and reminders + lib/bodies.js the request and response an alert shows, redacted lib/procs.js pid → systemd unit, container, script or command lib/sni.js the hostname in a TLS ClientHello lib/path.js /users/42 → /users/{n} @@ -306,8 +352,9 @@ So both spellings go through local flavors, `struct iov_iter___old { iov }` and | --- | --- | | `lib/http/decoder.js` | One entry per connection. Decides from the first bytes whether it carries HTTP/1, the HTTP/2 preface, a TLS record or something else, holds records 20 ms to put them back in kernel-timestamp order, pairs requests with responses, and accounts for bytes the copy missed so a gap costs a body, not the framing. | | `lib/http/h1.js`, `h2.js`, `hpack.js` | HTTP/1.x (content-length, chunked, pipelining) and HTTP/2 frames with HPACK header decoding. | -| `lib/apis.js` | Keys each transaction to an API: served by local port, called by `Host`. Drops the client side of a loopback hop into the served API's caller list. | -| `lib/alerts.js` | Per-API state: `ok`, `failing`, `down`. Only transitions and reminders produce messages. | +| `lib/apis.js` | Keys each transaction to an API: served by local port, called by `Host`. Drops the client side of a loopback hop into the served API's caller list. Keeps request and 4xx counts in 10-second buckets for about half an hour, the baseline's memory, and the latest failing exchange of each class. | +| `lib/alerts.js` | Per-API state: `ok`, `failing`, `down` for 5xx and ports, and a separate `ok`/`failing` for 4xx with the baseline frozen while it fires. Only transitions and reminders produce messages. | +| `lib/bodies.js` | The exchange an alert shows: inflates gzip, deflate and brotli bodies through `yeet:compression`, cuts each to about 1000 characters, and redacts by key name in JSON, forms and query strings and by shape (bearer tokens, JWTs, Luhn-valid card numbers) everywhere else. Headers other than the content type are never shown. | | `lib/procs.js` | Names a pid from its cgroup and command line, looked up the moment its first byte is captured, before a short-lived process can exit. | | `lib/capture.js` | Polls the socket inventory every 2 s for listeners, rescans for TLS libraries every 60 s, and records the SNI of each ClientHello so an unreadable call still has a name. | @@ -366,7 +413,10 @@ yeet run . -- --discover --seconds 25 - **HTTP/3.** QUIC runs over UDP, which these hooks never see. - **gRPC failures that arrive as HTTP 200.** gRPC reports its own errors in a `grpc-status` trailer on a 200 response, which `apiwatch` does not treat as an error. [`grpcsnoop`](https://github.com/yeet-src/grpcsnoop) decodes gRPC calls and their messages. - **Failures with no HTTP response.** A called API that refuses the connection, times out, or never answers produces no status code, so it is not alerted on as that API. If your app turns it into a 5xx, that is what you'll hear about. -- **Slowness, 4xx, and silence.** Alerts are 5xx and port-down only. A latency spike, a run of 404s, or an API whose traffic stops (while its port stays up) does not alert. +- **Slowness and silence.** A latency spike, or an API whose traffic stops while its port stays up, does not alert. +- **A 4xx jump in the first five minutes.** The baseline needs five minutes and 50 requests of an API before it can tell a jump from normal, and it restarts from nothing when the watcher does. `--client-errors all` alerts from the first second, on every 4xx. +- **Every secret.** Redaction goes by key names and recognisable shapes. A secret under an innocent key (`{"note":"my password is …"}`), a free-text body, or personal data such as names and email addresses goes through as captured. Use `--bodies off` where that matters. +- **Bodies of HTTPS it cannot read**, which is Go and other unhookable TLS stacks, and bodies over about 32 KiB, which arrive cut. Alerts for those APIs still fire, with whatever was captured. - **Anything before it started.** State lives in memory: discovery sees only its window, and a restarted watcher starts from zero, so a port that was already down when it started is unknown until you name it with `--ports`. - **The cost of capturing everything.** Without `--ports` the socket tap copies every TCP call on the box, not just HTTP, and the decoder discards what isn't. That is cheap on an API server and expensive on a database or a file server pushing gigabytes. Use `--ports` there. - **Every read under heavy concurrency.** A kretprobe has a fixed pool of in-flight instances, so when more threads sit in `tcp_recvmsg` at once than the pool holds, some reads are skipped. A skipped read is never recorded as a wrong one. diff --git a/package.json b/package.json new file mode 100644 index 0000000..f4465da --- /dev/null +++ b/package.json @@ -0,0 +1,8 @@ +{ + "name": "apiwatch", + "private": true, + "type": "module", + "scripts": { + "test": "node --test test/bodies.test.js test/alerts.test.js" + } +} diff --git a/src/lib/alerts.js b/src/lib/alerts.js index 976609c..59a03fb 100644 --- a/src/lib/alerts.js +++ b/src/lib/alerts.js @@ -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. @@ -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())); @@ -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) { @@ -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"), }; } @@ -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` }] }, ], }; diff --git a/src/lib/apis.js b/src/lib/apis.js index afca964..7b18b82 100644 --- a/src/lib/apis.js +++ b/src/lib/apis.js @@ -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("]")); @@ -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; @@ -76,6 +84,9 @@ export class Apis { errorTimes: [], reqTimes: [], recentErrors: [], + recent4xx: [], + buckets: [], + lastTx: {}, ...init, }; this.byKey.set(key, api); @@ -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 diff --git a/src/lib/bodies.js b/src/lib/bodies.js new file mode 100644 index 0000000..cb9cea4 --- /dev/null +++ b/src/lib/bodies.js @@ -0,0 +1,196 @@ +/* The request and response an alert shows, with secrets taken out. + * + * An alert that says "POST /orders → 502" tells you something broke; the + * body the client sent and the error the server gave back usually tell + * you why. Both are already in the captured transaction. What this module + * adds is the care they need before they leave the box: request bodies + * carry passwords, tokens and card numbers, and an alert goes through + * yeet's servers into a Slack channel. + * + * redacted (the default) values under a sensitive-looking key are + * replaced in JSON, form and query strings, and bearer + * tokens, JWTs and card numbers are replaced anywhere + * raw bodies as captured, still cut to `limit` + * off method, path and status only + * + * Headers are never shown except the content type. Pure: `inflate` + * (yeet:compression's decodeContentEncoding) is injected. + */ + +import { header } from "./http/h1.js"; +import { utf8 } from "./http/bytes.js"; + +export const MODES = ["redacted", "raw", "off"]; +const MASK = "[redacted]"; + +/* A key is sensitive when its words, split on separators and camelCase, + * name a secret. Whole words, so `shipping` is not a `pin` and `author` + * is not `auth`. */ +const SECRET_WORDS = new Set([ + "password", "passwd", "pass", "pwd", "passphrase", "secret", "token", "auth", "authorization", + "credential", "credentials", "cookie", "session", "sessionid", "sid", "otp", "pin", "cvv", "cvc", + "ssn", "card", "iban", "signature", "sig", "apikey", "key", "jwt", "bearer", +]); +/* A run-together word (`accesstoken`, `clientsecret`, `userpassword`) + * counts when it ends in a secret, or starts with `password`; that keeps + * `tokenize` and `secretary` out. */ +const SECRET_ENDINGS = /(password|passwd|passphrase|secret|token|apikey|accesskey|privatekey|authorization|credentials?|cardnumber|ccnumber|cvv|cvc|ssn|iban)$/; +const SECRET_START = /^(password|passwd)/; + +export function isSecretKey(key) { + const words = String(key) + .replace(/([a-z0-9])([A-Z])/g, "$1 $2") + .toLowerCase() + .split(/[^a-z0-9]+/) + .filter(Boolean); + return words.some((w) => SECRET_WORDS.has(w) || SECRET_ENDINGS.test(w) || SECRET_START.test(w)) || SECRET_ENDINGS.test(words.join("")); +} + +/* 13 to 19 digits, spaced or dashed, passing Luhn: a card number. */ +const luhn = (digits) => { + let sum = 0; + for (let i = 0; i < digits.length; i++) { + let d = digits.charCodeAt(digits.length - 1 - i) - 48; + if (i % 2) { + d *= 2; + if (d > 9) d -= 9; + } + sum += d; + } + return sum % 10 === 0; +}; + +/** Secrets that are recognisable by their shape, wherever they sit. */ +export function scrubText(s) { + return String(s) + .replace(/\b(Bearer|Basic|Token)\s+[A-Za-z0-9._~+/=-]{8,}/gi, `$1 ${MASK}`) + .replace(/\beyJ[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]*/g, MASK) + .replace(/\b\d(?:[ -]?\d){12,18}\b/g, (m) => { + const digits = m.replace(/\D/g, ""); + return digits.length >= 13 && digits.length <= 19 && luhn(digits) ? MASK : m; + }); +} + +export function redactJson(value) { + if (Array.isArray(value)) return value.map(redactJson); + if (value && typeof value === "object") { + const out = {}; + for (const [k, v] of Object.entries(value)) out[k] = isSecretKey(k) ? MASK : redactJson(v); + return out; + } + return typeof value === "string" ? scrubText(value) : value; +} + +/** `a=1&password=x` with the sensitive values replaced. */ +export function redactForm(s) { + return String(s) + .split("&") + .map((pair) => { + const eq = pair.indexOf("="); + if (eq < 0) return pair; + let key = pair.slice(0, eq); + try { + key = decodeURIComponent(key.replace(/\+/g, " ")); + } catch { + /* keep it as written */ + } + return isSecretKey(key) ? `${pair.slice(0, eq)}=${MASK}` : scrubText(pair); + }) + .join("&"); +} + +/** A request target with its query string redacted. */ +export function redactTarget(target) { + const t = String(target ?? ""); + const q = t.indexOf("?"); + return q < 0 ? scrubText(t) : `${scrubText(t.slice(0, q))}?${redactForm(t.slice(q + 1))}`; +} + +const cut = (s, limit) => (s.length > limit ? `${s.slice(0, limit)}… (${s.length} characters, cut)` : s); + +function isText(bytes) { + const n = Math.min(bytes.length, 512); + let bad = 0; + for (let i = 0; i < n; i++) { + const c = bytes[i]; + if (c === 0 || (c < 32 && c !== 9 && c !== 10 && c !== 13)) bad++; + } + return bad * 20 < n; +} + +/** One body as text for an alert, or a note saying why there is none. */ +export function bodyText(body, headers, { mode = "redacted", inflate = null, limit = 1000 } = {}) { + if (mode === "off" || !body || !body.len) return null; + let data = body.data ?? new Uint8Array(0); + const partial = body.truncated || body.holes || !body.complete; + const enc = (header(headers, "content-encoding") ?? "").toLowerCase(); + if (enc && enc !== "identity") { + if (partial || !inflate) return `[${enc} body, ${body.len} bytes, not decoded]`; + try { + data = inflate(enc, data); + } catch { + return `[${enc} body, ${body.len} bytes, did not decode]`; + } + } + if (!isText(data)) return `[binary body, ${body.len} bytes]`; + let text = utf8(data); + const ct = (header(headers, "content-type") ?? "").toLowerCase(); + if (mode === "redacted") { + const first = text.trimStart()[0]; + let done = false; + if (!partial && (ct.includes("json") || first === "{" || first === "[")) { + try { + text = JSON.stringify(redactJson(JSON.parse(text))); + done = true; + } catch { + /* not JSON after all */ + } + } + if (!done) text = ct.includes("x-www-form-urlencoded") ? redactForm(text) : scrubText(text); + } + text = cut(text, limit); + return partial ? `${text} (captured in part)` : text; +} + +/** + * The exchange an alert shows: `{ request, response }`, each with a first + * line and a body (string or null). Never throws. + */ +export function sampleOf(tx, opts = {}) { + const mode = opts.mode ?? "redacted"; + try { + const target = mode === "raw" ? String(tx.target ?? "") : redactTarget(tx.target); + const reqType = header(tx.reqHeaders ?? [], "content-type"); + const resType = header(tx.resHeaders ?? [], "content-type"); + return { + request: { + line: `${tx.method ?? "?"} ${target}`, + type: reqType, + body: bodyText(tx.reqBody, tx.reqHeaders ?? [], opts), + }, + response: { + line: `${tx.status ?? "no response"}${tx.reason ? ` ${tx.reason}` : ""}`, + type: resType, + body: bodyText(tx.resBody, tx.resHeaders ?? [], opts), + }, + }; + } catch (error) { + return { request: { line: `${tx.method ?? "?"} ?`, body: null }, response: { line: String(tx.status ?? "?"), body: null }, error: String(error) }; + } +} + +/* Slack mrkdwn needs &, < and > escaped, and a code block cannot hold + * three backticks in a row. */ +const slackSafe = (s) => String(s).replace(/&/g, "&").replace(//g, ">").replace(/```/g, "``​`"); + +/** A sample as Slack mrkdwn: the request and the response, each in a code block. */ +export function sampleMrkdwn(sample, label = "Latest failing request") { + if (!sample) return null; + const block = (part) => { + const lines = [part.line]; + if (part.type) lines.push(`content-type: ${part.type}`); + if (part.body) lines.push("", part.body); + return "```" + slackSafe(lines.join("\n")) + "```"; + }; + return `*${label}*\n${block(sample.request)}\n*Response*\n${block(sample.response)}`; +} diff --git a/src/main.js b/src/main.js index b79ecb8..764126b 100644 --- a/src/main.js +++ b/src/main.js @@ -14,6 +14,8 @@ import { startCapture, SELF_COMMS } from "./lib/capture.js"; import { Apis } from "./lib/apis.js"; import { Alerts } from "./lib/alerts.js"; import { describe, known } from "./lib/procs.js"; +import { MODES as BODY_MODES, sampleMrkdwn, sampleOf } from "./lib/bodies.js"; +import { decodeContentEncoding } from "yeet:compression"; const HELP = `apiwatch: the HTTP APIs this machine serves and calls, and a Slack alert when one breaks. @@ -24,6 +26,15 @@ Alert when one breaks (run as a yeet service so it outlives your shell): yeet run github:yeet-src/apiwatch -- --watch --slack "#channel" [--name web-1] --window 60 seconds of history each check looks at --min-errors 1 5xx responses within the window that count as broken + --client-errors baseline + 4xx alerts: baseline (when an API's 4xx share jumps far above + its own normal, learned over its first 5 min), all, or off + --client-codes 401,403,429 + count only these 4xx statuses (default: every 4xx) + --min-client-errors 5 + 4xx answers within the window before a jump can fire + --bodies redacted the failing request and response in each alert: redacted + (secret fields, tokens and card numbers blanked), raw, or off --recover 120 seconds without a 5xx before an API counts as recovered --remind 1800 seconds between "still broken" reminders --down-after 5 seconds a served port must be gone before it counts as down @@ -52,6 +63,14 @@ const portsArg = String(arg("ports", "") || "") const host = String(arg("name", "") || "this host"); const channel = arg("slack", null); const json = Boolean(arg("json", false)); +const bodies = String(arg("bodies", "redacted")); +const clientErrors = String(arg("client-errors", "baseline")); +const clientCodes = new Set( + String(arg("client-codes", "") || "") + .split(",") + .map((c) => Number(c.trim())) + .filter((c) => c >= 400 && c < 500), +); const out = (line) => console.log(line); const jlog = (obj) => console.log(JSON.stringify({ t: new Date().toISOString(), ...obj })); @@ -94,6 +113,9 @@ async function boot() { const capture = await startCapture({ base: import.meta.dirname, ports: portsArg, + /* Enough of each body that a compressed one can be decoded whole; an + * alert shows at most about 1000 characters of it. */ + bodyLimit: 16_384, onPid: (pid) => { if (!known(pid)) describe(pid); }, @@ -107,7 +129,7 @@ async function boot() { if (row && !known(row.pid)) describe(row.pid); return row?.pid ?? null; }; - apis = new Apis({ name, info: known, isLocal: capture.isLocal, listeners: capture.listeners, callerOf }); + apis = new Apis({ name, info: known, isLocal: capture.isLocal, listeners: capture.listeners, callerOf, clientCodes }); return { capture, apis, errors }; } @@ -200,6 +222,8 @@ async function discover() { async function watch() { const dryRun = Boolean(arg("dry-run", false)); if (!channel && !dryRun) throw new Error('--watch needs --slack "#channel" (or --dry-run)'); + if (!BODY_MODES.includes(bodies)) throw new Error(`--bodies must be one of ${BODY_MODES.join(", ")}, not "${bodies}"`); + if (!["baseline", "all", "off"].includes(clientErrors)) throw new Error(`--client-errors must be baseline, all or off, not "${clientErrors}"`); const window = num("window", 60); const { capture, apis, errors } = await boot(); @@ -235,14 +259,19 @@ async function watch() { .filter(Boolean) .join("\n"); const combine = (batch) => { - const broken = batch.filter((m) => m.event === "failing" || m.event === "down" || m.event === "reminder").length; + const broken = batch.filter((m) => /^(failing|down|reminder)/.test(m.event)).length; const title = broken === batch.length ? `${batch.length} APIs broke on ${host}` : broken === 0 ? `${batch.length} APIs recovered on ${host}` : `${batch.length} API changes on ${host}`; const blocks = [{ type: "header", text: { type: "plain_text", text: title } }]; - for (const m of batch) blocks.push({ type: "section", text: { type: "mrkdwn", text: `*${m.title}*\n${m.blocks[1].text.text}` } }, { type: "divider" }); + for (const m of batch.slice(0, 12)) { + blocks.push({ type: "section", text: { type: "mrkdwn", text: `*${m.title}*\n${m.blocks[1].text.text}` } }); + if (m.detail) blocks.push(m.blocks[2]); + blocks.push({ type: "divider" }); + } blocks.pop(); - blocks.push(batch[0].blocks[2]); + /* The footer (host, time) is each message's last block. */ + blocks.push(batch[0].blocks[batch[0].blocks.length - 1]); return { title, text: `${title}: ${batch.map((m) => m.title).join("; ")}`, blocks }; }; @@ -257,6 +286,14 @@ async function watch() { downAfter: num("down-after", 5), ports: portsArg, ignore: String(arg("ignore", "") || "").split(",").map((x) => x.trim()).filter(Boolean), + clientErrors, + minClient: Math.max(1, num("min-client-errors", 5)), + sample: (api, cls) => { + const tx = api.lastTx?.[cls]; + if (!tx) return null; + const label = bodies === "off" ? "Latest failing request (bodies off)" : `Latest failing request${bodies === "redacted" ? " (secrets redacted)" : ""}`; + return sampleMrkdwn(sampleOf(tx, { mode: bodies, inflate: decodeContentEncoding }), label); + }, log: jlog, send: (m) => { if (queue.length < 50) queue.push(m); @@ -273,10 +310,11 @@ async function watch() { alerts.check(); }, 1000); - /* What is being watched, 20 s and 60 s after start and then every ten - * minutes, so whoever attaches to the service log (which shows only - * what is printed after they attach) sees that capture is alive and - * whether alerts can be delivered. */ + /* What is being watched, 20 s after start and then every minute, so + * whoever connects to the log (which shows only what is printed after + * they connect) sees within a minute that capture is alive and whether + * alerts can be delivered. An agent that takes a while between starting + * the service and reading its log must not land in a silent gap. */ const status = async () => { await Promise.all([...apis.pids()].map((p) => describe(p))); const list = apis.list().filter((r) => r.requests > 0); @@ -287,6 +325,8 @@ async function watch() { channel, dryRun, signedIn, + bodies, + clientErrors, transactions: apis.transactions, apis: list.map((r) => ({ kind: r.kind, name: r.name, port: r.port, requests: r.requests, errors5xx: r.errors5xx })), tls: capture.tlsStatus().filter((t) => t.state === "attached").map((t) => t.path), @@ -294,8 +334,7 @@ async function watch() { }); }; setTimeout(status, 20_000); - setTimeout(status, 60_000); - setInterval(status, 600_000); + setInterval(status, 60_000); } /* ---------------------------------------------------------------- */ diff --git a/test/alerts.test.js b/test/alerts.test.js new file mode 100644 index 0000000..48169f7 --- /dev/null +++ b/test/alerts.test.js @@ -0,0 +1,81 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { Apis } from "../src/lib/apis.js"; +import { Alerts } from "../src/lib/alerts.js"; + +/* A served API on port 3000, fed synthetic transactions at chosen times. */ +function rig(opts = {}) { + let now = 1_000_000_000_000; + const realNow = Date.now; + Date.now = () => now; + const apis = new Apis({ name: () => "users-api", listeners: () => new Map([[3000, { pid: 1 }]]) }); + const sent = []; + const alerts = new Alerts({ apis, listeners: () => new Map([[3000, { pid: 1 }]]), send: (m) => sent.push(m.render()), host: "box", sample: () => "*sample*", ...opts }); + const tx = (status) => apis.observe({ role: "server", pid: 1, method: "GET", target: "/users/7", status, flow: { sport: 3000, daddr: "10.0.0.9", dport: 5555 } }); + const advance = (ms, perSec, share404) => { + for (let t = 0; t < ms; t += 1000) { + now += 1000; + for (let i = 0; i < perSec; i++) tx(i < perSec * share404 ? 404 : 200); + alerts.check(now); + } + }; + return { apis, alerts, sent, advance, restore: () => (Date.now = realNow), now: () => now }; +} + +test("a normal 404 rate never alerts", () => { + const r = rig(); + r.advance(15 * 60_000, 10, 0.1); + assert.equal(r.sent.filter((m) => m.event === "failing_4xx").length, 0); + r.restore(); +}); + +test("a jump far above the learned baseline alerts once, with the sample, then recovers", () => { + const r = rig(); + r.advance(10 * 60_000, 10, 0.1); // learn: 10% normal + r.advance(90_000, 10, 0.6); // spike to 60% + const fired = r.sent.filter((m) => m.event === "failing_4xx"); + assert.equal(fired.length, 1); + assert.match(fired[0].title, /4xx jumped to/); + assert.match(fired[0].text, /against 10% normally/); + assert.equal(fired[0].detail, "*sample*"); + r.advance(3 * 60_000, 10, 0.1); // back to normal + assert.equal(r.sent.filter((m) => m.event === "recovered_4xx").length, 1); + r.restore(); +}); + +test("nothing fires before the baseline is learned", () => { + const r = rig(); + r.advance(3 * 60_000, 10, 0.9); + assert.equal(r.sent.filter((m) => m.event === "failing_4xx").length, 0); + r.restore(); +}); + +test("a long spike does not become the new normal while it fires", () => { + const r = rig({ remind: 100_000 }); + r.advance(10 * 60_000, 10, 0.05); + r.advance(20 * 60_000, 10, 0.5); + assert.equal(r.sent.filter((m) => m.event === "recovered_4xx").length, 0); + r.restore(); +}); + +test("--client-errors all fires on the first 4xx; off never does", () => { + const a = rig({ clientErrors: "all" }); + a.advance(5_000, 10, 0.1); + assert.equal(a.sent.filter((m) => m.event === "failing_4xx").length, 1); + a.restore(); + const o = rig({ clientErrors: "off" }); + o.advance(15 * 60_000, 10, 0.9); + assert.equal(o.sent.filter((m) => /4xx/.test(m.event)).length, 0); + o.restore(); +}); + +test("5xx still alerts on the first error and carries the sample", () => { + const r = rig(); + r.advance(5_000, 10, 0); + r.apis.observe({ role: "server", pid: 1, method: "POST", target: "/users", status: 502, flow: { sport: 3000, daddr: "10.0.0.9", dport: 5555 } }); + r.alerts.check(r.now()); + const f = r.sent.filter((m) => m.event === "failing"); + assert.equal(f.length, 1); + assert.equal(f[0].detail, "*sample*"); + r.restore(); +}); diff --git a/test/bodies.test.js b/test/bodies.test.js new file mode 100644 index 0000000..1c2335c --- /dev/null +++ b/test/bodies.test.js @@ -0,0 +1,78 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { gzipSync } from "node:zlib"; +import { isSecretKey, redactForm, redactTarget, scrubText, bodyText, sampleOf, sampleMrkdwn } from "../src/lib/bodies.js"; + +const enc = (s) => new TextEncoder().encode(s); +const body = (s, extra = {}) => ({ len: enc(s).length, data: enc(s), complete: true, holes: 0, truncated: false, ...extra }); + +test("secret keys are whole words, not substrings", () => { + for (const k of ["password", "user_password", "apiKey", "api_key", "X-Api-Key", "accessToken", "client_secret", "cardNumber", "cvv", "ssn", "pin", "Authorization", "session_id", "accesstoken", "clientSecret", "passwordHash"]) { + assert.ok(isSecretKey(k), k); + } + for (const k of ["shipping", "author", "passenger", "keyboard_layout", "amount", "email", "spinner", "tokenize_count", "secretary"]) { + assert.ok(!isSecretKey(k), k); + } +}); + +test("JSON bodies have secret values replaced, everything else kept", () => { + const t = bodyText(body('{"user":"ana","password":"hunter2","card":{"number":"4111111111111111"},"amount":2599,"note":"Bearer abcdefghijklmnop"}'), [["content-type", "application/json"]]); + const o = JSON.parse(t); + assert.equal(o.user, "ana"); + assert.equal(o.password, "[redacted]"); + assert.equal(o.card, "[redacted]"); + assert.equal(o.amount, 2599); + assert.equal(o.note, "Bearer [redacted]"); +}); + +test("card numbers are caught by Luhn, other long numbers are not", () => { + assert.equal(scrubText("pay 4111 1111 1111 1111 now"), "pay [redacted] now"); + assert.equal(scrubText("order 1234567890123"), "order 1234567890123"); +}); + +test("JWTs are replaced anywhere", () => { + assert.equal(scrubText("t=eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiIxMjM0NTY3ODkwIn0.abcdefghij"), "t=[redacted]"); +}); + +test("forms and query strings", () => { + assert.equal(redactForm("user=ana&password=x&amount=3"), "user=ana&password=[redacted]&amount=3"); + assert.equal(redactTarget("/login?next=%2Fhome&token=abc123"), "/login?next=%2Fhome&token=[redacted]"); + assert.equal(redactTarget("/orders/42"), "/orders/42"); +}); + +test("raw mode leaves bodies alone, off mode drops them", () => { + const b = body('{"password":"hunter2"}'); + assert.equal(bodyText(b, [], { mode: "raw" }), '{"password":"hunter2"}'); + assert.equal(bodyText(b, [], { mode: "off" }), null); +}); + +test("compressed bodies are inflated when an inflater is given", () => { + const z = gzipSync(Buffer.from('{"error":"payments unreachable","token":"zzz"}')); + const b = { len: z.length, data: new Uint8Array(z), complete: true, holes: 0, truncated: false }; + const inflate = (_enc, data) => new Uint8Array(require_gunzip(data)); + assert.equal(JSON.parse(bodyText(b, [["content-encoding", "gzip"]], { inflate })).token, "[redacted]"); + assert.match(bodyText(b, [["content-encoding", "gzip"]]), /gzip body, \d+ bytes, not decoded/); +}); +import { gunzipSync } from "node:zlib"; +function require_gunzip(d) { return gunzipSync(Buffer.from(d)); } + +test("binary bodies are shown as a size", () => { + const b = { len: 4, data: new Uint8Array([0, 1, 2, 0]), complete: true, holes: 0, truncated: false }; + assert.equal(bodyText(b, []), "[binary body, 4 bytes]"); +}); + +test("long bodies are cut, partial ones are marked", () => { + assert.match(bodyText(body("x".repeat(3000)), [], { limit: 100 }), /3000 characters, cut/); + assert.match(bodyText(body("abc", { complete: false }), []), /captured in part/); +}); + +test("a sample renders as Slack mrkdwn with escaping", () => { + const tx = { method: "POST", target: "/orders?api_key=k1", status: 502, reason: "Bad Gateway", + reqHeaders: [["content-type", "application/json"]], reqBody: body('{"amount":1,"password":"p"}'), + resHeaders: [["content-type", "text/html"]], resBody: body("

502 Bad Gateway

") }; + const md = sampleMrkdwn(sampleOf(tx), "Latest failing request"); + assert.match(md, /POST \/orders\?api_key=\[redacted\]/); + assert.match(md, /"password":"\[redacted\]"/); + assert.match(md, /<h1>502 Bad Gateway<\/h1>/); + assert.match(md, /^\*Latest failing request\*/); +});