Skip to content
Closed
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
31 changes: 31 additions & 0 deletions IMPLEMENTATION_DOCS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
# Implementation docs

Index of the docs that describe how this tree actually behaves. Concept
and contract pages live next to the code they cover; this file is only
the map.

## Decisions

- [docs/adr/README.md](./docs/adr/README.md) — Architecture Decision
Records (devices table, message ordering, ciphertext-only storage,
per-device envelopes, sealed-box → Signal, MLS for groups).

## Backend concepts

- [apps/backend/docs/concepts-configuration.md](./apps/backend/docs/concepts-configuration.md) —
boot-validated vs lazy env, `loadEnv()`, object-store singleton.
- [apps/backend/docs/concepts-backpressure.md](./apps/backend/docs/concepts-backpressure.md) —
socket buffer tiers; only disconnect is enforced today.
- [apps/backend/docs/concepts-stellar-listener.md](./apps/backend/docs/concepts-stellar-listener.md) —
Soroban event watch, cursors, idempotency, RPC failure.
- [apps/backend/docs/concepts-gateway-architecture.md](./apps/backend/docs/concepts-gateway-architecture.md)
- [apps/backend/docs/concepts-delivery-fanout.md](./apps/backend/docs/concepts-delivery-fanout.md)
- [apps/backend/docs/concepts-storage-push-jobs.md](./apps/backend/docs/concepts-storage-push-jobs.md)

## Operations and security

- [docs/runbook.md](./docs/runbook.md)
- [docs/observability.md](./docs/observability.md)
- [docs/threat-model.md](./docs/threat-model.md)
- [docs/security/rate-limits.md](./docs/security/rate-limits.md)
- [`.env.example`](./.env.example) — environment-variable reference
100 changes: 100 additions & 0 deletions apps/backend/docs/concepts-backpressure.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
# Backpressure and slow-consumer handling

How `services/backpressure.ts` watches per-socket send-buffer occupancy, what
the two thresholds do **today**, and how a client recovers after a hard
disconnect.

Gateway architecture mentions this module in passing
([concepts-gateway-architecture.md](./concepts-gateway-architecture.md) §8).
That overview currently describes shed as "stop sending". This page is the
accurate account of what is wired.

## Monitoring

On connect, `index.ts` calls `registerForBackpressure(socket)`. The socket is
added to an in-process set. A single `setInterval(checkBuffers, 5000)` runs
while any socket is registered and is cleared when the last one
unregisters (disconnect).

Each tick reads the engine.io transport's `bufferedAmount` (bytes waiting
in the WebSocket send buffer). If that field is missing or throws, occupancy
is treated as `0` — the socket is not shed or disconnected on a read error.

This is **per socket**, not per user. Two devices of the same account are
independent. The sets live in process memory; they are not shared across
gateway instances.

## Two tiers

| Tier | Env | Default | What the code does today |
| --- | --- | --- | --- |
| Shed | `SOCKET_SHED_THRESHOLD` | `32768` bytes | Marks the socket id in `shedSockets`, increments `clicked_backpressure_events_total{action="shed"}`, logs a warning. **Does not change emit behaviour.** |
| Disconnect | `SOCKET_BUFFER_THRESHOLD` | `65536` bytes | Same mark + metric `{action="disconnect"}`, then `socket.disconnect(true)`. **This is the only tier that affects delivery.** |

Both env values are parsed on every tick (`parseInt`, must be a positive
integer). Invalid or unset values fall back to the defaults above.

### Which tier changes emit behaviour

`isSocketShed(socketId)` exists and is exported. **Nothing in the backend
calls it.** No dispatcher, fan-out, or `socket.emit` path consults
`shedSockets` before sending.

So:

- **Disconnect** is enforced. Crossing `SOCKET_BUFFER_THRESHOLD` force-closes
the socket.
- **Shed** is telemetry and a flag only. Crossing `SOCKET_SHED_THRESHOLD`
does **not** stop new events being queued. A reader of the gateway
overview must not infer that the server sheds load today — the buffer
will keep growing until the disconnect threshold (or the client catches
up and the flag is cleared on a later tick).

When occupancy drops back below the shed threshold, the id is removed from
`shedSockets`. That only matters for the metric/flag, not for emits.

## Hard disconnect — what the client sees

`socket.disconnect(true)` is a server-initiated close (`close` / disconnect
on the client, reason typically `io server disconnect`). The socket is
removed from rooms, presence, the device-delivery subscriber, and
backpressure monitoring. In-flight emits for that connection are gone.

The client is not told "you were too slow". There is no dedicated
backpressure event. From the client's point of view this is the same shape
as any other unexpected drop: the connection is dead and it must open a
new one and authenticate again.

## How resume recovers

A new connection is a new socket. After auth it re-joins conversation /
user / device rooms and is registered for backpressure from scratch
(`shedSockets` does not survive the old socket id).

Missed **durable** chat is not in the WebSocket buffer and is not in the
resume stream. The client recovers it with envelope sync
(`GET /sync`, `syncRequired: true` on `resume_complete`). See
[api-messages-sync.md](./api-messages-sync.md) and
[concepts-delivery-fanout.md](./concepts-delivery-fanout.md).

Missed **ephemeral** events (read/delivery receipts, presence, system
notices) are recovered by the resume/replay path:

1. Client emits `resume { lastEventId }` with the last Redis stream id it
persisted.
2. The gateway reads `resume:events:${userId}` after that id (exclusive)
and emits `ephemeral_replay` for each entry still within the 300s TTL /
500-entry cap.
3. It finishes with `resume_complete { lastEventId, syncRequired: true }`.

Implementation: `services/resumeStream.ts`, handler in
`socket/messaging.ts`. Protocol notes:
[concepts-gateway-architecture.md](./concepts-gateway-architecture.md) §5
and [api-websocket-events.md](./api-websocket-events.md) (`resume` /
`resume_complete` / `ephemeral_replay`).

Events that were only sitting in the killed socket's send buffer and were
never recorded to Postgres or the resume stream are not replayed. That is
the loss window a slow consumer accepts when it is hard-disconnected.

Cross-referenced from [`IMPLEMENTATION_DOCS.md`](../../../IMPLEMENTATION_DOCS.md).
119 changes: 119 additions & 0 deletions apps/backend/docs/concepts-configuration.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
# Backend configuration loading and validation

How the gateway reads environment variables: which ones fail boot, which ones
are read later, and how that interacts with the object-store singleton.

This page is about **when** a variable is consumed. The catalog of names,
defaults, and meanings lives in the environment-variable reference
([`.env.example`](../../../.env.example)). Rate-limit buckets and override
syntax are in [`docs/security/rate-limits.md`](../../../docs/security/rate-limits.md).
Do not treat the tables below as a second copy of those lists.

## Boot-validated set (`src/config.ts`)

`index.ts` calls `loadEnv()` immediately after `dotenv.config()`. `loadEnv()`
parses `process.env` against `EnvSchema` (Zod). On success it returns the
typed object and prints nothing. On failure it logs the offending names and
**exits the process with code 1** — the HTTP server is never bound.

Required (missing or empty → hard startup failure):

| Variable | Constraint |
| --- | --- |
| `DATABASE_URL` | non-empty string |
| `REDIS_URL` | non-empty string |
| `JWT_SECRET` | non-empty string |
| `PORT` | positive integer (string is coerced) |
| `TOKEN_TRANSFER_CONTRACT_ID` | non-empty string |
| `OBJECT_STORE_ENDPOINT` | non-empty string |
| `OBJECT_STORE_BUCKET` | non-empty string |
| `OBJECT_STORE_ACCESS_KEY` | non-empty string |
| `OBJECT_STORE_SECRET_KEY` | non-empty string |
| `OBJECT_STORE_REGION` | non-empty string |
| `OBJECT_STORE_FORCE_PATH_STYLE` | `true` / `false` / `1` / `0` |

Optional in the same schema (validated only when present; absence does not
fail boot): `VAPID_PUBLIC_KEY`, `VAPID_PRIVATE_KEY`, `VAPID_SUBJECT`, the
legacy `S3_*` aliases, and `IDEMPOTENCY_TTL_SECONDS`.

`assertTransportSecurityConfig()` runs next. That is a separate check
(`ENFORCE_TLS`, `ALLOWED_ORIGINS`) and is not part of `EnvSchema`.

### Failure output

If `DATABASE_URL` is missing, stderr looks like this and the process exits 1:

```text
Missing or invalid environment variables: DATABASE_URL
- DATABASE_URL: DATABASE_URL is required
```

An empty environment lists every required name, one `- name: message` line
each. A non-numeric `PORT` is reported the same way (`PORT must be an integer`).

## Lazy vs eager

**Eager (boot).** `EnvSchema` variables above. The process does not start
without them. `index.ts` then immediately calls `getObjectStore()`, so the
object-store singleton is constructed on the boot path too — but the
construction itself is still lazy at the *module* level (see below).

**Lazy (first use, no boot failure).** Everything else is read from
`process.env` at the call site, with a hardcoded default when unset or
malformed:

- Rate-limit buckets in `src/config/rateLimits.ts` via `getRateLimitRule()`.
Each call re-reads `RATE_LIMIT_<BUCKET>` (and the legacy
`SOCKET_RATE_LIMIT_PER_SEC` for `socket_default`). A malformed override is
warned and ignored; the default stands. `RATE_LIMIT_DISABLED=true` is a
kill switch, also read on demand.
- Backpressure thresholds (`SOCKET_BUFFER_THRESHOLD`, `SOCKET_SHED_THRESHOLD`)
— see [concepts-backpressure.md](./concepts-backpressure.md).
- Stellar RPC wiring (`STELLAR_RPC_URL`, `GROUP_TREASURY_CONTRACT_ID`) —
missing values disable the listener rather than crashing; see
[concepts-stellar-listener.md](./concepts-stellar-listener.md).
- Sync/envelope knobs such as `ENVELOPE_TTL_SECONDS` and `SYNC_PAGE_SIZE`.

Rate-limit rules are deliberately not cached at import time so tests and the
runbook can change a limit without restarting the module graph.

## `loadEnv()` and the object-store singleton

`lib/objectStore.ts` keeps a process-wide singleton behind `getObjectStore()`.
The client is **not** built at import time. The first call constructs it:

- `NODE_ENV === 'production'` → `createObjectStore(loadEnv())` (real S3 /
MinIO / R2 client, credentials from the boot-validated `OBJECT_STORE_*`
fields).
- otherwise → `getLocalObjectStore()` (fs-backed store; no live S3 needed).

Building the singleton in `objectStore.ts` (instead of exporting a client
from `index.ts`) avoids a circular import:

```text
index.ts → app.ts → routes → lib/storage.ts → index.ts
```

`storage.ts` and `services/fileCleanup.ts` both call `getObjectStore()`, so
there is exactly one client/bucket pair per process. `loadEnv()` is invoked
again inside that first production call; by then boot validation has already
run, so it is a typed re-parse, not a second chance to start with a broken
env.

`resetObjectStoreForTests()` clears the memo so tests can change env.

## Production vs development in `lib/storage.ts`

`generatePresignedPut` / `generatePresignedGet` branch on
`NODE_ENV === 'production'`:

| Branch | Implementation |
| --- | --- |
| production | `getObjectStore()` — the S3-compatible client from `OBJECT_STORE_*` |
| any other `NODE_ENV` | `getLocalObjectStore()` — local-disk store, URLs served by `routes/localStorage.ts` |

The same `NODE_ENV` split exists inside `getObjectStore()` itself. The
storage helpers short-circuit to the local store in non-production so
presigned-URL callers never touch the S3 constructor on a laptop or in CI.

Cross-referenced from [`IMPLEMENTATION_DOCS.md`](../../../IMPLEMENTATION_DOCS.md).
124 changes: 124 additions & 0 deletions apps/backend/docs/concepts-stellar-listener.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
# Stellar chain listener and on-chain reconciliation

How `services/stellarListener.ts` watches Soroban contract events, writes
`token_transfers` / `treasury_proposals`, tracks its place in the ledger, and
behaves after downtime or an unreachable RPC.

The listener is started from `index.ts` only when both `STELLAR_RPC_URL` and
`TOKEN_TRANSFER_CONTRACT_ID` are set. Missing either logs

```text
[stellar-listener] STELLAR_RPC_URL or TOKEN_TRANSFER_CONTRACT_ID unset; listener disabled.
```

and leaves the API up. `GROUP_TREASURY_CONTRACT_ID` is optional: when set, a
second fetcher is attached for treasury events.

`runForever` never rethrows. RPC and DB errors are logged inside the loop so
a dead chain cannot take the HTTP/WebSocket process down.

## Watched events and database writes

Two independent pollers share the same loop (default interval 5s).

### `token_transfer` — topic `transfer`

`buildRpcFetcher` calls Soroban RPC `getEvents` filtered to
`TOKEN_TRANSFER_CONTRACT_ID` and topic `transfer`. Each event is mapped to:

| RPC field | Persisted column (`token_transfers`) |
| --- | --- |
| `txHash` | `tx_hash` (unique) |
| `value.to` | `recipient_address` |
| `value.amount` | `amount` (decimal string) |
| `value.memo` | `memo` (hex, optional) |

`conversation_id` / `sender_id` are filled by decoding `memo` as a message
UUID and looking up `messages`. If that misses, the writer falls back to an
arbitrary existing conversation and user; if the database has neither, the
event is dropped.

Write path: `INSERT … ON CONFLICT (tx_hash) DO UPDATE SET created_at = now()`.
A replayed ledger page refreshes `created_at` on the same row; it does not
insert a second transfer.

### `group_treasury` — proposal lifecycle

`buildTreasuryRpcFetcher` watches `GROUP_TREASURY_CONTRACT_ID` for:

| Contract event | Row status |
| --- | --- |
| `proposal_created` | `active` |
| `proposal_approved` | `approved` |
| `proposal_rejected` | `rejected` |
| `proposal_executed` | `executed` |
| `proposal_expired` | `expired` |

Persistence is an **update** on `(contract_id, proposal_id)` (unique index
`treasury_proposals_contract_proposal_idx`), copying `approvals` /
`rejections` when present. If no matching row exists, the event is ignored
(`if (!row) return`) — the listener does not insert treasury proposals. After
a successful update it emits `treasury_proposal_updated` to the linked
conversation Socket.IO room.

## Cursor / position tracking

Both pollers keep an **in-process** cursor (`pagingToken` from the last
successfully persisted event). There is no cursor table and nothing is
written to Redis or Postgres.

- First poll after process start: `cursor` is `null`. `getEvents` is called
with `startLedger` and `cursor` both unset (the fetcher comment: "resume
on cursor only").
- After each successful persist, the cursor advances to that event's
`pagingToken`. A persist failure leaves the cursor unmoved so the next
poll retries the same page.
- Token-transfer and treasury streams have separate cursors.

Cursors die with the process. A restart is a cold start: both cursors are
`null` again.

## Catch-up after downtime

Two different "down" cases:

1. **Process was down (restart / deploy).** Cursors reset. The next
`getEvents` page is whatever the RPC returns without a cursor. Events
still inside the RPC's retention window are re-delivered and upserted.
Events that have aged out of that window are not replayed by this
listener — those mirrored rows stay missing until something else writes
them.
2. **Process stayed up but a poll failed.** `consecutiveFailures`
increments. The loop waits `min(1000 * 2^(n-1), 30000)` ms and retries
**from the last good in-memory cursor**, so it continues where it left
off rather than rewinding.

## Idempotency

A replayed ledger event must not double-write:

- Token transfers: unique `tx_hash` + `ON CONFLICT DO UPDATE`. The same
hash can be persisted any number of times; there is still one row.
- Treasury proposals: unique `(contract_id, proposal_id)`. A repeated
event updates status / counts on the existing row.

Because the cursor only advances after a successful persist, a crash
mid-page can re-offer the same events; the unique keys absorb that.

## RPC unreachable — drift from on-chain state

When `getEvents` throws (timeout, 5xx, network partition):

- The error is logged as `fetch failed; reconnecting after backoff`.
- The API and WebSocket server keep running.
- No new rows are written for the duration of the outage.

**Yes, on-chain state can drift from the mirrored tables.** The listener is
a best-effort projection, not a consensus participant. During an RPC outage
the chain continues; `token_transfers` and `treasury_proposals` do not.
After reconnect, catch-up depends on RPC retention and on the treasury
update-only rule (a `proposal_*` event for a row that was never inserted
is silently skipped). There is no compensating backfill job in this
service.

Cross-referenced from [`IMPLEMENTATION_DOCS.md`](../../../IMPLEMENTATION_DOCS.md).
Loading
Loading