diff --git a/README.md b/README.md index 73c75da..2fb7129 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ Automations for the AT Protocol — listen to events, filter them, and trigger actions like webhook deliveries or PDS record creation. -Airglow connects to [Jetstream](https://atproto.com/guides/streaming-data#jetstream) (AT Protocol's event streaming service), matches incoming records against user-defined automations, and executes actions automatically. Think IFTTT or Zapier, but the trigger side is always "something happened on the AT Protocol." +Airglow connects to [Jetstream](https://bsky.network/docs/jetstream) (AT Protocol's event streaming service), matches incoming records against user-defined automations, and executes actions automatically. Think IFTTT or Zapier, but the trigger side is always "something happened on the AT Protocol." Website: [airglow.run](https://airglow.run) @@ -133,14 +133,15 @@ goat lex new record run.airglow. # create a new lexicon Airglow is designed to be easy to self-host. Configuration is done via environment variables (see `.env.example`): -| Variable | Purpose | -| ---------------- | --------------------------------------------------- | -| `PUBLIC_URL` | Public-facing base URL of the instance | -| `DATABASE_PATH` | Path to the SQLite database file | -| `JETSTREAM_URL` | Jetstream WebSocket endpoint | -| `COOKIE_SECRET` | Secret for session cookies (min 32 chars) | -| `NSID_ALLOWLIST` | Comma-separated NSIDs to allow (empty = allow all) | -| `NSID_BLOCKLIST` | Comma-separated NSIDs to block (empty = block none) | +| Variable | Purpose | +| --------------------------- | ------------------------------------------------------------- | +| `PUBLIC_URL` | Public-facing base URL of the instance | +| `DATABASE_PATH` | Path to the SQLite database file | +| `JETSTREAM_URL` | Jetstream WebSocket endpoint (defaults to a v2 instance) | +| `JETSTREAM_MAX_LOOKBACK_US` | Max age of a stored cursor still used to resume (default 36h) | +| `COOKIE_SECRET` | Secret for session cookies (min 32 chars) | +| `NSID_ALLOWLIST` | Comma-separated NSIDs to allow (empty = allow all) | +| `NSID_BLOCKLIST` | Comma-separated NSIDs to block (empty = block none) | Instance operators can configure `NSID_ALLOWLIST` and `NSID_BLOCKLIST` to control which lexicons their instance handles. For example, a typical instance may want to block `app.bsky.*` or `app.bsky.feed.*` since those collections are very active and could overwhelm a small instance. @@ -148,5 +149,5 @@ Instance operators can configure `NSID_ALLOWLIST` and `NSID_BLOCKLIST` to contro - [AT Protocol docs](https://atproto.com/docs) - [Lexicon spec](https://atproto.com/specs/lexicon) - - [Jetstream](https://atproto.com/guides/streaming-data#jetstream) + - [Jetstream](https://bsky.network/docs/jetstream) - [OAuth](https://atproto.com/specs/oauth) diff --git a/docs/architecture.md b/docs/architecture.md index f2e41b0..5d44046 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -72,7 +72,9 @@ The `inFlight` set lives on `JetstreamManager` (not per-subscription) because su ### Jetstream cursor persistence -Each subscription tracks the latest `time_us` it has observed and persists it to `jetstream_cursors` (one row per canonical DID-set key — empty = global). On (re)connect, if the persisted cursor is within `JETSTREAM_MAX_LOOKBACK_US` (default 24h), it is appended as `?cursor=` so Jetstream replays from that point. `lastSeenTimeUs` is advanced **before** dispatch — at-least-once delivery is the action layer's responsibility, via the composite PK on `automation_executions(automation_uri, source_uri)`. +Each subscription tracks the latest `time_us` it has observed and persists it to `jetstream_cursors` (one row per canonical DID-set key — empty = global). On (re)connect, if the persisted cursor is within `JETSTREAM_MAX_LOOKBACK_US` (default 36h), it is appended as `?cursor=` so Jetstream replays from that point. `lastSeenTimeUs` is advanced **before** dispatch — at-least-once delivery is the action layer's responsibility, via the composite PK on `automation_executions(automation_uri, source_uri)`. + +Jetstream v2 also emits its own monotonic sequence number as a top-level `cursor` field on each event. Airglow deliberately does **not** persist it: `time_us` is accepted as a resume token by both v1 and v2, so keeping the timestamp is what makes a rollback to v1 (which emits no sequence number) work without a cursor reset. Writes happen in three places: @@ -80,7 +82,30 @@ Writes happen in three places: - WS `close` handler before the reconnect timer fires, bounding loss-on-disconnect to ≤ one batch of in-memory events. Suppressed during manager shutdown so step 4 above is authoritative. - `refreshAutomations()` pre-flushes any sub it is about to tear down, so a same-key recreate in the same process reads the fresh row instead of racing the close handler's upsert. -Stale rows are swept hourly by [cleanup.ts](../lib/db/cleanup.ts) once they're older than `2 × maxLookbackUs` — past that age they're already beyond Jetstream's ~72h server-side replay buffer. +Stale rows are swept hourly by [cleanup.ts](../lib/db/cleanup.ts) once they're older than `2 × maxLookbackUs` — past that age they're already beyond the server-side replay buffer, so the row can no longer be used. + +### Lookback window and the server's replay buffer + +`JETSTREAM_MAX_LOOKBACK_US` bounds how old a stored cursor may be and still be used. Beyond it, the subscription starts live and the gap is lost — every automation that would have fired in that window never fires. + +The bound tracks Jetstream v2's live-tail buffer, measured at **~37.4h on 2026-08-17** (a 36h cursor was served in full; 38h already clamped). That buffer is a rolling window with no published guarantee, so 36h is a best-effort tuning with a little margin, not a contract — re-measure before changing it. + +Two properties worth knowing before touching this value: + +- **It is a client-side volume bound, not a correctness boundary.** An over-deep cursor is clamped by the server to its oldest retained event; it does not error and does not silently jump to live. So a larger value only risks a bigger replay burst, and anything past ~37h buys no extra depth. +- **It is only safe to raise because catch-up breaches don't disable.** Per-automation rate limits are 10/s, 100/min and 500/hr, and a breach normally auto-disables. A multi-hour replay would blow the hourly window within seconds of reconnecting, so `JETSTREAM_CATCH_UP_THRESHOLD_MS` (default 5 min) marks sufficiently-old events as catch-up: a breach there skips the event and logs once, but leaves the automation enabled. Events skipped that way are **not** retried, so gap recovery is best-effort. + +Closing gaps deeper than the live-tail buffer needs v2's archive replay (`planSnapshot` / `listSegments` / `getSegment`), which requires an API token and a separate HTTP ingestion path. Not implemented; 36h is a chosen stopping point, not v2's limit. + +### Rolling back to Jetstream v1 + +Set `JETSTREAM_URL` to a v1 host and restart: + +``` +JETSTREAM_URL=wss://jetstream2.us-east.bsky.network/subscribe +``` + +No cursor reset is needed in either direction, because the persisted cursor is a `time_us` timestamp that both versions accept. Note the host naming: `jetstream2` was instance 2 of the **v1** fleet, not version 2 — the v2 hosts are `jetstream.us-east` and `jetstream.us-west`. ### Known limitation: SIGKILL window