From 294616642f5cfea62a1a5ddd9656412ccc9530d7 Mon Sep 17 00:00:00 2001 From: Robert Goodall Date: Fri, 24 Jul 2026 13:11:06 -0400 Subject: [PATCH] initial public relay build repo --- .dockerignore | 7 + .env.example | 36 +++ .gitea/PULL_REQUEST_TEMPLATE.md | 11 + .gitea/workflows/ci.yml | 112 +++++++ .gitignore | 5 + CONTRIBUTING.md | 46 +++ Dockerfile | 33 ++ README.md | 211 ++++++++++++ cmd/nostermd/config.go | 74 +++++ cmd/nostermd/config_test.go | 49 +++ cmd/nostermd/federation_test.go | 135 ++++++++ cmd/nostermd/main.go | 95 ++++++ cmd/nostermd/relay.go | 177 ++++++++++ cmd/nostermd/relay_test.go | 194 +++++++++++ docker-compose.yml | 44 +++ go.mod | 39 +++ go.sum | 92 ++++++ internal/federation/federation.go | 406 +++++++++++++++++++++++ internal/federation/federation_test.go | 215 ++++++++++++ internal/group/group.go | 431 +++++++++++++++++++++++++ internal/group/group_test.go | 291 +++++++++++++++++ internal/relayinfo/relayinfo.go | 42 +++ internal/retention/retention.go | 96 ++++++ internal/retention/retention_test.go | 116 +++++++ 24 files changed, 2957 insertions(+) create mode 100644 .dockerignore create mode 100644 .env.example create mode 100644 .gitea/PULL_REQUEST_TEMPLATE.md create mode 100644 .gitea/workflows/ci.yml create mode 100644 .gitignore create mode 100644 CONTRIBUTING.md create mode 100644 Dockerfile create mode 100644 README.md create mode 100644 cmd/nostermd/config.go create mode 100644 cmd/nostermd/config_test.go create mode 100644 cmd/nostermd/federation_test.go create mode 100644 cmd/nostermd/main.go create mode 100644 cmd/nostermd/relay.go create mode 100644 cmd/nostermd/relay_test.go create mode 100644 docker-compose.yml create mode 100644 go.mod create mode 100644 go.sum create mode 100644 internal/federation/federation.go create mode 100644 internal/federation/federation_test.go create mode 100644 internal/group/group.go create mode 100644 internal/group/group_test.go create mode 100644 internal/relayinfo/relayinfo.go create mode 100644 internal/retention/retention.go create mode 100644 internal/retention/retention_test.go diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..2cbf4f7 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,7 @@ +data +relay +nostermd +.env +.env.* +!.env.example +.DS_Store diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..4763160 --- /dev/null +++ b/.env.example @@ -0,0 +1,36 @@ +# Nosterm relay deployment configuration. +# Copy to `.env` and adjust. + +# Stable 64-char hex identity that signs group metadata/rosters. +# REQUIRED when RELAY_ENV=production (the relay refuses to start without it). +# In dev it is auto-generated and persisted to the data volume. +# Generate one with: openssl rand -hex 32 +RELAY_SECRET_KEY= + +# Set to "production" to require a stable RELAY_SECRET_KEY (fail-fast if unset). +# Leave empty for development (auto-generates + persists a key). +RELAY_ENV= + +# NIP-11 relay metadata. +RELAY_NAME=Nosterm Home Relay +RELAY_DESCRIPTION= +RELAY_CONTACT=relay@nosterm.com + +# Optional message-of-the-day: shown in the client's relay "server window" on +# connect (IRC-style MOTD). Multi-line supported with \n. Empty = falls back to +# the built-in "nerdworks.io" ascii banner. +RELAY_MOTD= + +# Retention (chat messages only; membership/metadata are never pruned). +# 0 = unlimited. +RELAY_RETENTION_DAYS=0 +RELAY_RETENTION_MAX_MESSAGES=0 +RELAY_RETENTION_INTERVAL_MINUTES=60 + +# Optional federation with peer relays. Empty = disabled. +# Format: "[|:]" entries separated by ';'. +# mode = mirror (bidirectional) | ingest (read-only). Default mirror. +# channels = comma list of group ids, or '*' for all. Default '*'. +# Example: +# wss://a.example|mirror:general,dev;wss://b.example|ingest:* +RELAY_FEDERATION_PEERS= diff --git a/.gitea/PULL_REQUEST_TEMPLATE.md b/.gitea/PULL_REQUEST_TEMPLATE.md new file mode 100644 index 0000000..f7308fe --- /dev/null +++ b/.gitea/PULL_REQUEST_TEMPLATE.md @@ -0,0 +1,11 @@ + + +## Summary + +## Checklist + +- [ ] `gofmt -l .` prints nothing +- [ ] `go vet ./...` passes +- [ ] `go test ./...` passes +- [ ] `go build ./cmd/nostermd` succeeds +- [ ] No secrets committed diff --git a/.gitea/workflows/ci.yml b/.gitea/workflows/ci.yml new file mode 100644 index 0000000..a5af6a9 --- /dev/null +++ b/.gitea/workflows/ci.yml @@ -0,0 +1,112 @@ +# ci/cd for the public nostermd relay. +# +# - every push/pr: gofmt check, vet, test, build. +# - push to main / a v* tag: also build the image and push it to ecr public +# (anonymous pulls) so anyone can `docker pull` it. +# +# this is the public contribution repo — it builds and ships only the relay +# image to the public registry. it never touches private ecr or redeploys any +# environment; promoting an image to a live relay is handled out of band. +# +# gitea 1.21 actions on a self-hosted (ec2) runner. aws auth uses the runner's +# ambient iam role (no stored creds): `aws ecr-public get-login-password` mints +# a short-lived token for `docker login`. the runner role needs ecr public push +# perms (ecr-public:GetAuthorizationToken, sts:GetServiceBearerToken, +# BatchCheckLayerAvailability, Put*/Upload*/Complete*). +name: ci + +env: + # ecr public always authenticates via the us-east-1 endpoint. + AWS_REGION: us-east-1 + # public mirror for anonymous pulls. k3k1z1x5 is the account's ecr public + # registry alias; the repo name (nostermd) is provisioned out of band. + ECR_PUBLIC_REGISTRY: public.ecr.aws/k3k1z1x5 + ECR_PUBLIC_REPO: nostermd + +on: + push: + branches: [main] + tags: ['v*'] + pull_request: + +jobs: + build: + runs-on: ubuntu-latest + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Setup Go + uses: actions/setup-go@v5 + with: + go-version: '1.25' + + - name: Verify formatting (gofmt) + run: | + unformatted="$(gofmt -l .)" + if [ -n "$unformatted" ]; then + echo "These files are not gofmt-clean:" + echo "$unformatted" + exit 1 + fi + + - name: Vet + run: go vet ./... + + # no -race: go-nostr's unsafe json serializer trips the detector during + # event signing (upstream lib issue, not a relay race). see README.md#tests. + - name: Test + run: go test ./... + + - name: Build + run: go build -trimpath -o /dev/null ./cmd/nostermd + + image: + needs: build + # publish only from the default branch and version tags, never from prs. + if: gitea.event_name == 'push' && (gitea.ref == 'refs/heads/main' || startsWith(gitea.ref, 'refs/tags/v')) + runs-on: ubuntu-latest + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Install AWS CLI + # runner image ships without aws; install v2 if it's missing. + run: | + if ! command -v aws >/dev/null 2>&1; then + apt-get update -qq + apt-get install -y -qq curl unzip + curl -fsSL "https://awscli.amazonaws.com/awscli-exe-linux-x86_64.zip" -o /tmp/awscliv2.zip + unzip -q /tmp/awscliv2.zip -d /tmp + /tmp/aws/install + fi + aws --version + + - name: Log in to Amazon ECR Public + # ecr public auth is a separate endpoint (always us-east-1) with its own + # login user (AWS). this is push auth; anonymous pulls need no login. + run: | + aws ecr-public get-login-password --region us-east-1 \ + | docker login --username AWS --password-stdin public.ecr.aws + + - name: Build and push image + shell: bash + run: | + set -eo pipefail + IMAGE="$ECR_PUBLIC_REGISTRY/$ECR_PUBLIC_REPO" + # tag by commit sha for traceability. + docker build -t "$IMAGE:${{ gitea.sha }}" . + docker push "$IMAGE:${{ gitea.sha }}" + # move :latest only on main. + if [ "${{ gitea.ref }}" = "refs/heads/main" ]; then + docker tag "$IMAGE:${{ gitea.sha }}" "$IMAGE:latest" + docker push "$IMAGE:latest" + fi + # on a v* tag, also push that version tag. + case "${{ gitea.ref }}" in + refs/tags/v*) + VERSION="${{ gitea.ref_name }}" + docker tag "$IMAGE:${{ gitea.sha }}" "$IMAGE:$VERSION" + docker push "$IMAGE:$VERSION" + ;; + esac diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..03f3aa3 --- /dev/null +++ b/.gitignore @@ -0,0 +1,5 @@ +# local, untracked claude instructions +CLAUDE.local.md + +# local relay runtime state (boltdb store + generated key) +data/ diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md new file mode 100644 index 0000000..39b556e --- /dev/null +++ b/CONTRIBUTING.md @@ -0,0 +1,46 @@ +# Contributing to nostermd + +Thanks for helping improve the Nosterm home relay. This repo is the public home +for relay development — the Go source and its container image. It builds and +ships the relay image to ECR Public; it does not deploy to any live relay. + +## Workflow + +1. Fork the repo (or branch, if you have push access). +2. Create a topic branch off `main`. +3. Make your change, keeping commits focused. +4. Open a pull request against `main`. + +Every push and pull request runs CI: gofmt check, vet, tests, and a build. Get +these green before requesting review: + +```bash +gofmt -l . # must print nothing +go vet ./... +go test ./... # run without -race — see README.md#tests +go build ./cmd/nostermd +``` + +`gofmt -w .` auto-fixes formatting. + +## What happens on merge + +Merging to `main` builds the container image and pushes it to ECR Public +(`public.ecr.aws/k3k1z1x5/nostermd:latest`), so anyone can pull the latest relay +anonymously. Tagged releases (`v*`) publish a matching versioned image. This +repo does not deploy to any live environment — promotion is handled separately. + +## Guidelines + +- Match the surrounding code style; gofmt is the source of truth. +- Keep pull requests scoped to one change; smaller is easier to review. +- Add or update tests when you change behavior. Tests live beside the package + they cover. +- Never commit secrets. `RELAY_SECRET_KEY` is a per-deployment identity, not a + repo value. + +## Reporting issues + +Open an issue with steps to reproduce, what you expected, and what happened. For +security-sensitive reports, please disclose privately rather than in a public +issue. diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..45f2648 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,33 @@ +# syntax=docker/dockerfile:1 + +# --- Build stage ----------------------------------------------------------- +FROM golang:1.25-alpine AS build +WORKDIR /src + +# Cache modules independently of source. +COPY go.mod go.sum ./ +RUN go mod download + +COPY . . +# CGO disabled → a fully static binary (boltdb is pure Go), so it runs on scratch. +# The command lives in ./cmd/nostermd; internal/* holds the relay packages. +RUN CGO_ENABLED=0 go build -trimpath -ldflags="-s -w" -o /out/nostermd ./cmd/nostermd +# Stage an empty data dir we can copy in with the correct nonroot ownership. +RUN mkdir -p /out/data + +# --- Runtime stage --------------------------------------------------------- +FROM gcr.io/distroless/static-debian12:nonroot AS runtime + +# Binary + a data dir pre-owned by the distroless nonroot uid (65532), so the +# unprivileged process can write the boltdb file to the mounted volume. +COPY --from=build --chown=65532:65532 /out/nostermd /nostermd +COPY --from=build --chown=65532:65532 /out/data /data + +# Persisted event storage (mount a volume here in production). +VOLUME ["/data"] +ENV RELAY_DATA_DIR=/data \ + RELAY_ADDR=:3334 + +USER 65532:65532 +EXPOSE 3334 +ENTRYPOINT ["/nostermd"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..51b01fc --- /dev/null +++ b/README.md @@ -0,0 +1,211 @@ +# nostermd + +A minimal **NIP-29 managed-group** relay for +[Nosterm](https://www.nosterm.com), built on +[`khatru`](https://fiatjaf.com/nostr/khatru) (part of the `fiatjaf.com/nostr` monorepo), +plus the **Nosterm capability handshake** — a `features` array in its NIP-11 document so a +Nosterm client can feature-gate its UI and fall back gracefully on plain relays. + +This is the "home relay" in Nosterm's hybrid model: the client connects to this relay for +its own communities while still connecting to public relays for general Nostr content. + +## What it does + +- Serves NIP-01/11/42 via khatru, with pure-Go [BoltDB](https://github.com/etcd-io/bbolt) + storage (no cgo → a fully static binary that runs on `scratch`/distroless). +- Enforces a focused subset of **NIP-29 managed groups**: + | Kind | Event | Behavior | + | ----- | ---------------- | ------------------------------------------- | + | 9007 | create group | registers a group | + | 9021 | join request | adds the sender as a member (auto-creates) | + | 9022 | leave request | removes the sender | + | 9000 | put user (admin) | admin adds tagged members | + | 9001 | remove user | admin removes tagged members | + | 9/10 | chat / reply | accepted for members (open groups: anyone) | + | 39000 | group metadata | (re)published + signed by the relay | +- Advertises `software: "nostermd"` and `features: ["channels"]` in NIP-11 (adds + `"federation"` when any federation peer is configured). +- **Optional federation** — mirror/ingest chat with peer relays (see below). + +> **MVP scope.** Groups are auto-created as **open** on first join so the client flow works +> without an invite system. Membership is authoritative in-memory; richer moderation, +> private/closed groups, roles, and the 39001/39002 admin/member lists are future work. + +## Project layout + +Standard Go layout — the command is a thin composition root over focused, +independently-testable packages: + +``` +nostermd/ + cmd/nostermd/ # main: env config, key loading, khatru wiring, NIP-11 serving + internal/ + group/ # NIP-29 managed-group state machine (membership, admin, metadata) + federation/ # optional chat mirroring/ingest with peer relays + retention/ # background pruning of old chat messages + relayinfo/ # Nosterm capability handshake (NIP-11 features array) + Dockerfile # static binary → distroless nonroot image + docker-compose.yml +``` + +## Run it + +```bash +# Local (Go 1.25+) +go run ./cmd/nostermd # listens on :3334, ephemeral identity + ./data + +# Configurable via env +RELAY_ADDR=:3334 \ +RELAY_DATA_DIR=./data \ +RELAY_NAME="My Relay" \ +RELAY_SECRET_KEY=<64-char-hex> \ # stable identity (signs group metadata) + go run ./cmd/nostermd + +# Docker (single image) +docker build -t nostermd . +docker run -p 3334:3334 -v nostermd-data:/data nostermd + +# Or pull the prebuilt public image (no auth): +docker run -p 3334:3334 -v nostermd-data:/data public.ecr.aws/k3k1z1x5/nostermd:latest +``` + +Merges to `main` build and publish that image to ECR Public; `v*` tags publish a +matching versioned image. This repo builds the image only — it does not deploy. + +### Docker Compose (recommended) + +```bash +cp .env.example .env + +# For production, generate a stable signing identity and set RELAY_ENV: +echo "RELAY_SECRET_KEY=$(openssl rand -hex 32)" >> .env +echo "RELAY_ENV=production" >> .env + +docker compose up -d --build +``` + +The relay then listens on `ws://localhost:3334` and persists its BoltDB event store to a +named volume so rooms/messages survive restarts. + +Point a Nosterm client at it via `PUBLIC_DEFAULT_RELAYS` — see the +[nosterm-client](https://git.nerdworks.io/nerdworks/nosterm-client) repository. + +### Production TLS + +The relay speaks plain WebSocket on `:3334`. For public `wss://` access, put it behind a +TLS-terminating reverse proxy (e.g. Caddy, or the +[nosterm-client](https://git.nerdworks.io/nerdworks/nosterm-client) Caddy service, which +can proxy `/relay` to this relay on a shared Docker network). + +## Environment + +| Variable | Default | Meaning | +| --------------------------------- | -------------------- | -------------------------------------------------------------- | +| `RELAY_ADDR` | `:3334` | Listen address | +| `RELAY_DATA_DIR` | `./data` | BoltDB storage directory | +| `RELAY_NAME` | `Nosterm Home Relay` | NIP-11 name | +| `RELAY_DESCRIPTION` | _(default text)_ | NIP-11 description | +| `RELAY_CONTACT` | `relay@nosterm.com` | NIP-11 contact (NIP-05-style operator address) | +| `RELAY_MOTD` | _(nerdworks.io banner)_ | Message-of-the-day shown in the client's relay server window (defaults to the built-in ascii banner) | +| `RELAY_ENV` | _(dev)_ | Set to `production` to require a stable key (fail-fast if unset) | +| `RELAY_SECRET_KEY` | _(dev: persisted)_ | 64-char hex identity that signs group metadata/rosters | +| `RELAY_RETENTION_DAYS` | `0` | Prune chat messages older than N days (0 = unlimited) | +| `RELAY_RETENTION_MAX_MESSAGES` | `0` | Keep at most N chat messages, newest first (0 = unlimited) | +| `RELAY_RETENTION_INTERVAL_MINUTES`| `60` | How often the retention sweep runs | +| `RELAY_FEDERATION_PEERS` | _(empty: disabled)_ | Peer relays to mirror/ingest chat with (see Federation) | + +### Relay identity (`RELAY_SECRET_KEY`) + +The relay's key signs the group metadata/admins/members events (39000–39002). A **stable** +identity is required so those signatures stay valid across restarts: + +- **Production** (`RELAY_ENV=production`): `RELAY_SECRET_KEY` is **required** — the relay + refuses to start without it. Generate one with `openssl rand -hex 32`. +- **Development**: if unset, a key is generated and **persisted** to + `$RELAY_DATA_DIR/relay.key` (mode `0600`) and reused on every restart, so even local dev + keeps one identity. + +### Retention + +By default nothing is pruned. Set `RELAY_RETENTION_DAYS` and/or +`RELAY_RETENTION_MAX_MESSAGES` to cap growth. Retention **only** deletes chat messages +(kinds 9/10); group-management and metadata events (9007/9000/9001/9021/9022, 39000–39002) +are never pruned so membership always survives and can be replayed on boot. + +## Federation (optional) + +Federation lets a channel live on more than one relay, so it survives any single relay going +down and can be distributed across operators. It is **off by default** and **chat-only** by +design: only kind-9/10 messages cross the relay boundary. Membership, admins, bans, and +metadata stay **local to each relay** — every operator moderates their own instance, which +sidesteps cross-relay moderation conflicts entirely. + +Configure peers with `RELAY_FEDERATION_PEERS`. Entries are separated by `;`; each is: + +``` +[|:] +``` + +- **mode** — `mirror` (bidirectional: chat posted on either relay is forwarded to the other) + or `ingest` (read-only: pull chat _from_ the peer, never push back). Defaults to `mirror`. +- **channels** — a comma-separated list of group ids, or `*` for all. Defaults to `*`. A + leading `#` is tolerated. + +Examples: + +```bash +# Mirror two channels with peer A, ingest everything from peer B (read-only): +RELAY_FEDERATION_PEERS="wss://a.example|mirror:general,dev;wss://b.example|ingest:*" + +# Read-only announcement feed pulled from an upstream relay: +RELAY_FEDERATION_PEERS="wss://announce.example|ingest:announcements" + +# Simplest: fully mirror every channel with one peer: +RELAY_FEDERATION_PEERS="wss://peer.example" +``` + +**How it works.** The relay opens a client subscription to each peer for the federated chat +kinds and injects received events through its own add pipeline (`OnEvent` → store → +`OnEventSaved`), then broadcasts them to local subscribers. Locally-saved chat is forwarded +to `mirror` peers from the `OnEventSaved` hook. Incoming events are **re-verified** (a peer +can't inject forged messages) and scoped to the configured channels. + +**Loop prevention is structural, not tag-based.** A re-received event is a duplicate in the +eventstore, and `AddEvent` stops on `ErrDupEvent` _before_ the egress hook runs — so an event +is never forwarded twice and `A↔B↔A` cycles die on the second sight. There is no "seen" set to +tune or get wrong. Peers reconnect with capped backoff; on reconnect the peer replays stored +events so nothing missed during an outage is lost (deduped on arrival). + +**Trust model.** Federate only with relays you trust to authorize their own posters — a +federated channel behaves as open on every participating relay (federated chat bypasses +_local_ membership because the author is a member on the peer that accepted it). Federation is +strictly chat; it never imports another relay's membership, admin, or ban state. + +## Tests + +```bash +go test ./... +``` + +Tests live beside the package they cover: + +- `internal/group` — the group state machine: join → chat allowed, chat-before-join rejected, + missing `h`-tag rejected, **multi-user administrator moderation** in a closed group (creator + is admin; non-admins cannot add/remove users; admin-approved members can post; removed + members are blocked; open groups allow anyone), and **state survives a rebuild** from the + event store. +- `internal/retention` — prunes old/surplus chat by age and count but never + management/metadata events. +- `internal/federation` — config parsing (modes, channel scoping, invalid-URL/empty-list + rejection), ingress signature + scope enforcement, and the loop-prevention contract. +- `cmd/nostermd` — integration over a real in-process khatru server: **relay-key + handling** (explicit key used verbatim; dev key persisted and stable across restarts), a + **cross-user** WebSocket round-trip (one client joins and publishes, another receives via + subscription), and two end-to-end **federation** two-relay tests (a message ingested across + relays, and a mirrored post delivered exactly once despite the echo path). + +> Run without `-race`: go-nostr's `unsafe`-based JSON serializer trips the race detector's +> `checkptr` during any event signing (an upstream library issue, not a relay data race). The +> relay's own group-state locking is race-clean — verify with the state-machine tests that +> don't sign events: +> `go test ./internal/group/ -race -run 'TestAdminModerationInClosedGroup|TestOpenGroup'`. diff --git a/cmd/nostermd/config.go b/cmd/nostermd/config.go new file mode 100644 index 0000000..9c06bea --- /dev/null +++ b/cmd/nostermd/config.go @@ -0,0 +1,74 @@ +package main + +import ( + "log" + "os" + "path/filepath" + "strconv" + + "fiatjaf.com/nostr" +) + +func envOr(key, fallback string) string { + if v := os.Getenv(key); v != "" { + return v + } + return fallback +} + +// envIntOr parses a non-negative integer env var, falling back on empty/invalid. +func envIntOr(key string, fallback int) int { + if v := os.Getenv(key); v != "" { + if n, err := strconv.Atoi(v); err == nil && n >= 0 { + return n + } + log.Printf("invalid %s=%q; using %d", key, v, fallback) + } + return fallback +} + +// loadOrGenerateKey resolves the relay's signing identity, which signs the group +// metadata/admins/members events. A STABLE identity is required so those signed +// lists remain valid across restarts and re-deploys. +// +// - RELAY_SECRET_KEY set → use it (production path). +// - RELAY_ENV=production, unset → fatal. A fresh key each boot silently +// corrupts signed group rosters, so we refuse to start. +// - otherwise (dev) → persist a generated key under the data dir +// and reuse it, so even dev keeps a stable identity across restarts. +func loadOrGenerateKey(dataDir string) nostr.SecretKey { + if hex := os.Getenv("RELAY_SECRET_KEY"); hex != "" { + sk, err := nostr.SecretKeyFromHex(hex) + if err != nil { + log.Fatalf("invalid RELAY_SECRET_KEY: %v", err) + } + return sk + } + + if isProduction() { + log.Fatal("RELAY_SECRET_KEY is required when RELAY_ENV=production " + + "(a fresh key each boot would corrupt signed group rosters). " + + "Generate one with: openssl rand -hex 32") + } + + // Dev convenience: persist a generated key so restarts keep one identity. + keyPath := filepath.Join(dataDir, "relay.key") + if data, err := os.ReadFile(keyPath); err == nil { + if sk, err := nostr.SecretKeyFromHex(string(data)); err == nil { + log.Printf("using persisted dev relay identity from %s", keyPath) + return sk + } + log.Printf("ignoring unreadable %s; generating a new dev key", keyPath) + } + sk := nostr.Generate() + if err := os.WriteFile(keyPath, []byte(sk.Hex()), 0o600); err != nil { + log.Printf("could not persist dev relay key (%v); identity is ephemeral this run", err) + } else { + log.Printf("generated and persisted a dev relay identity at %s (set RELAY_SECRET_KEY for production)", keyPath) + } + return sk +} + +func isProduction() bool { + return os.Getenv("RELAY_ENV") == "production" +} diff --git a/cmd/nostermd/config_test.go b/cmd/nostermd/config_test.go new file mode 100644 index 0000000..7c390a4 --- /dev/null +++ b/cmd/nostermd/config_test.go @@ -0,0 +1,49 @@ +package main + +import ( + "os" + "path/filepath" + "testing" + + "fiatjaf.com/nostr" +) + +func TestDevKeyPersistsAcrossCalls(t *testing.T) { + t.Setenv("RELAY_SECRET_KEY", "") + t.Setenv("RELAY_ENV", "") + dir := t.TempDir() + + first := loadOrGenerateKey(dir) + // The key file should now exist and a second call must return the same key. + if _, err := os.Stat(filepath.Join(dir, "relay.key")); err != nil { + t.Fatalf("dev key not persisted: %v", err) + } + second := loadOrGenerateKey(dir) + if first != second { + t.Fatal("dev identity should be stable across restarts") + } +} + +func TestExplicitKeyIsUsed(t *testing.T) { + sk := nostr.Generate() + t.Setenv("RELAY_SECRET_KEY", sk.Hex()) + t.Setenv("RELAY_ENV", "production") // explicit key wins even in production + got := loadOrGenerateKey(t.TempDir()) + if got != sk { + t.Fatal("explicit RELAY_SECRET_KEY should be used verbatim") + } +} + +// Note: the production-without-key path calls log.Fatal (os.Exit), which can't +// be exercised in-process without a subprocess harness; its behavior is simple +// and covered by manual verification + the isProduction() guard below. +func TestIsProduction(t *testing.T) { + t.Setenv("RELAY_ENV", "production") + if !isProduction() { + t.Fatal("RELAY_ENV=production should be detected") + } + t.Setenv("RELAY_ENV", "") + if isProduction() { + t.Fatal("unset RELAY_ENV should not be production") + } +} diff --git a/cmd/nostermd/federation_test.go b/cmd/nostermd/federation_test.go new file mode 100644 index 0000000..06f997e --- /dev/null +++ b/cmd/nostermd/federation_test.go @@ -0,0 +1,135 @@ +package main + +import ( + "context" + "testing" + "time" + + "fiatjaf.com/nostr" +) + +// TestFederationIngest is an end-to-end test: relay B hosts an open channel; +// relay A ingests it. A message posted to B by a member appears on A even though +// the author is not a member of A — the chat-only, membership-bypass contract. +func TestFederationIngest(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + const groupID = "general" + + // Relay B: a normal relay hosting the channel. + bURL := startTestRelay(t, "") + + // Relay A: ingests channel "general" from B (read-only). + aURL := startTestRelay(t, bURL+"|ingest:"+groupID) + + // A subscriber on relay A watches the channel. + aReader, err := nostr.RelayConnect(ctx, aURL, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("connect reader to A: %v", err) + } + sub, err := aReader.Subscribe(ctx, nostr.Filter{ + Kinds: []nostr.Kind{nostr.KindSimpleGroupChatMessage}, + Tags: nostr.TagMap{"h": []string{groupID}}, + }, nostr.SubscriptionOptions{}) + if err != nil { + t.Fatalf("subscribe on A: %v", err) + } + defer sub.Unsub() + + // Alice joins + posts on relay B. + aliceSK := nostr.Generate() + bWriter, err := nostr.RelayConnect(ctx, bURL, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("connect alice to B: %v", err) + } + join := nostr.Event{Kind: nostr.KindSimpleGroupJoinRequest, CreatedAt: nostr.Now(), Tags: nostr.Tags{{"h", groupID}}} + mustSign(t, aliceSK, &join) + if err := bWriter.Publish(ctx, join); err != nil { + t.Fatalf("alice join on B: %v", err) + } + msg := nostr.Event{Kind: nostr.KindSimpleGroupChatMessage, CreatedAt: nostr.Now(), Content: "hello from B", Tags: nostr.Tags{{"h", groupID}}} + mustSign(t, aliceSK, &msg) + if err := bWriter.Publish(ctx, msg); err != nil { + t.Fatalf("alice post on B: %v", err) + } + + // The message should arrive on relay A via federation. + select { + case got := <-sub.Events: + if got.ID != msg.ID { + t.Fatalf("got event %s, want %s", got.ID.Hex(), msg.ID.Hex()) + } + if got.Content != "hello from B" { + t.Errorf("content = %q", got.Content) + } + case <-ctx.Done(): + t.Fatal("timed out waiting for federated message on relay A") + } +} + +// TestFederationMirrorNoDuplicate proves the loop guard end-to-end: A mirrors B, +// so a post on A is forwarded to B, which would echo it back to A. The echo is a +// store duplicate and is dropped, not re-delivered — the message arrives on A's +// subscription exactly once. +func TestFederationMirrorNoDuplicate(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + const groupID = "general" + + // B must start first (A needs B's URL); then A starts mirroring B. We assert + // the guard on the A→B→(echo)→A path. + bURL := startTestRelay(t, "") + aURL := startTestRelay(t, bURL+"|mirror:"+groupID) + + // Reader on A counts how many times the event is delivered. + aReader, err := nostr.RelayConnect(ctx, aURL, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("connect reader to A: %v", err) + } + sub, err := aReader.Subscribe(ctx, nostr.Filter{ + Kinds: []nostr.Kind{nostr.KindSimpleGroupChatMessage}, + Tags: nostr.TagMap{"h": []string{groupID}}, + }, nostr.SubscriptionOptions{}) + if err != nil { + t.Fatalf("subscribe on A: %v", err) + } + defer sub.Unsub() + + // Alice joins + posts on A. A forwards to B (mirror); B could echo back to A. + aliceSK := nostr.Generate() + aWriter, err := nostr.RelayConnect(ctx, aURL, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("connect alice to A: %v", err) + } + join := nostr.Event{Kind: nostr.KindSimpleGroupJoinRequest, CreatedAt: nostr.Now(), Tags: nostr.Tags{{"h", groupID}}} + mustSign(t, aliceSK, &join) + if err := aWriter.Publish(ctx, join); err != nil { + t.Fatalf("alice join on A: %v", err) + } + msg := nostr.Event{Kind: nostr.KindSimpleGroupChatMessage, CreatedAt: nostr.Now(), Content: "once only", Tags: nostr.Tags{{"h", groupID}}} + mustSign(t, aliceSK, &msg) + if err := aWriter.Publish(ctx, msg); err != nil { + t.Fatalf("alice post on A: %v", err) + } + + // Count deliveries of our message id for a short window; must be exactly 1. + deadline := time.After(3 * time.Second) + count := 0 + for { + select { + case got := <-sub.Events: + if got.ID == msg.ID { + count++ + } + case <-deadline: + if count != 1 { + t.Fatalf("expected message delivered exactly once, got %d", count) + } + return + case <-ctx.Done(): + t.Fatal("context canceled") + } + } +} diff --git a/cmd/nostermd/main.go b/cmd/nostermd/main.go new file mode 100644 index 0000000..e81d5cf --- /dev/null +++ b/cmd/nostermd/main.go @@ -0,0 +1,95 @@ +// Command nostermd is a minimal NIP-29 managed-group relay built on khatru, +// plus the Nosterm capability handshake (a `features` array in its NIP-11 doc). +// +// It is intentionally small: it enforces group membership for chat writes, +// auto-creates open channels on first join, and advertises itself as a +// "nostermd" so a Nosterm client can feature-gate its UI. The relay is +// assembled from focused internal packages (group, federation, retention, +// relayinfo); this file is the composition root that wires them to env config +// and the network. +package main + +import ( + "log" + "net/http" + "os" + "path/filepath" + "time" + + "fiatjaf.com/nostr/eventstore/boltdb" + + "github.com/nosterm/relay/internal/federation" + "github.com/nosterm/relay/internal/retention" +) + +// defaultMOTD is the built-in message-of-the-day: a "nerdworks.io" ascii banner +// (figlet, standard font) the client shows when RELAY_MOTD is unset. +const defaultMOTD = " _ _ _\n" + + " _ __ ___ _ __ __| |_ _____ _ __| | _____ (_) ___\n" + + "| '_ \\ / _ \\ '__/ _` \\ \\ /\\ / / _ \\| '__| |/ / __| | |/ _ \\\n" + + "| | | | __/ | | (_| |\\ V V / (_) | | | <\\__ \\_| | (_) |\n" + + "|_| |_|\\___|_| \\__,_| \\_/\\_/ \\___/|_| |_|\\_\\___(_)_|\\___/\n" + +func main() { + addr := envOr("RELAY_ADDR", ":3334") + dataDir := envOr("RELAY_DATA_DIR", "./data") + + // Embedded, pure-Go storage (no cgo → clean static Docker builds). + if err := os.MkdirAll(dataDir, 0o755); err != nil { + log.Fatalf("cannot create data dir: %v", err) + } + + // Relay identity: signs the 39000-series group metadata/roster events. A + // stable key is required so those signatures survive restarts (see below). + sk := loadOrGenerateKey(dataDir) + + store := &boltdb.BoltBackend{Path: filepath.Join(dataDir, "events.bolt")} + if err := store.Init(); err != nil { + log.Fatalf("cannot open eventstore: %v", err) + } + + // Parse optional federation config up front so warnings surface at boot. + fedCfg, fedWarnings := federation.ParseConfig(os.Getenv("RELAY_FEDERATION_PEERS")) + for _, w := range fedWarnings { + log.Printf("federation config: %s", w) + } + + app := buildRelay(relayConfig{ + secretKey: sk, + store: store, + name: envOr("RELAY_NAME", "Nosterm Home Relay"), + description: envOr("RELAY_DESCRIPTION", "A Nosterm NIP-29 managed-group relay."), + // NIP-11 contact: a NIP-05-style identifier / address for the relay operator. + contact: envOr("RELAY_CONTACT", "relay@nosterm.com"), + // Optional message-of-the-day shown in the client's relay server window; + // defaults to the nerdworks.io banner when unset. + motd: envOr("RELAY_MOTD", defaultMOTD), + fedCfg: fedCfg, + }) + + // Reconstruct membership/metadata from persisted events so groups survive + // restarts (without this, every restart silently drops all group state). + if replayed := app.groups.Rebuild(store); replayed > 0 { + log.Printf("replayed %d group management event(s) from storage", replayed) + } + + // Retention: prune old chat messages so the store doesn't grow forever. + stopRetention := retention.Start(store, retention.Config{ + MaxAge: time.Duration(envIntOr("RELAY_RETENTION_DAYS", 0)) * 24 * time.Hour, + MaxMessages: envIntOr("RELAY_RETENTION_MAX_MESSAGES", 0), + Interval: time.Duration(envIntOr("RELAY_RETENTION_INTERVAL_MINUTES", 60)) * time.Minute, + }) + defer stopRetention() + + // Start federation last, once the relay pipeline is fully wired. + if app.federator != nil { + app.federator.Start() + defer app.federator.Stop() + } + + log.Printf("nostermd pubkey: %s", app.pubKey.Hex()) + log.Printf("listening on %s (%s; %s)", addr, app.groups.Summary(), fedCfg.Summary()) + if err := http.ListenAndServe(addr, app.handler); err != nil { + log.Fatal(err) + } +} diff --git a/cmd/nostermd/relay.go b/cmd/nostermd/relay.go new file mode 100644 index 0000000..c7da43f --- /dev/null +++ b/cmd/nostermd/relay.go @@ -0,0 +1,177 @@ +package main + +import ( + "context" + "encoding/json" + "net/http" + + "fiatjaf.com/nostr" + "fiatjaf.com/nostr/eventstore" + "fiatjaf.com/nostr/khatru" + + "github.com/nosterm/relay/internal/federation" + "github.com/nosterm/relay/internal/group" + "github.com/nosterm/relay/internal/relayinfo" +) + +// relayConfig holds everything buildRelay needs to assemble a relay. It is the +// seam between env/flag parsing (main) and the wiring (buildRelay), so tests can +// construct a relay with an in-memory store and explicit config. +type relayConfig struct { + secretKey nostr.SecretKey + store eventstore.Store + name string + description string + contact string + motd string + fedCfg federation.Config +} + +// relayApp is a fully-wired relay: the khatru relay, its HTTP handler, the group +// state authority, and (when configured) the federator. main starts/stops these; +// tests assert against them. +type relayApp struct { + relay *khatru.Relay + handler http.Handler + groups *group.State + federator *federation.Federator + pubKey nostr.PubKey +} + +// buildRelay assembles the relay from its internal packages and wires khatru's +// hooks: NIP-29 enforcement on ingress, addressable-list (re)publishing on +// group mutations, federation ingress/egress, and NIP-11 serving with the +// Nosterm capability handshake. It does NOT rebuild state, start retention, +// or start federation — the caller owns lifecycle so tests can drive it. +func buildRelay(cfg relayConfig) *relayApp { + pk := nostr.GetPublicKey(cfg.secretKey) + + groups := group.NewState(cfg.secretKey) + if cfg.fedCfg.Enabled() { + // Chat on a federated channel is accepted regardless of local membership. + groups.IsFederated = cfg.fedCfg.FederatesChannel + } + + relay := khatru.NewRelay() + relay.Info.Name = cfg.name + relay.Info.Description = cfg.description + relay.Info.Contact = cfg.contact + relay.Info.PubKey = &pk + relay.Info.SupportedNIPs = []any{1, 11, 29, 42} + relay.UseEventstore(cfg.store, 500) + + // Enforce NIP-29 group rules on every write. + relay.OnEvent = func(ctx context.Context, evt nostr.Event) (reject bool, msg string) { + return groups.Evaluate(evt) + } + + // Federator: connects to peers, pulls their chat (injecting it through the + // relay's own add pipeline — whose eventstore dedup breaks forwarding loops), + // and forwards locally-saved chat to mirror peers via the OnEventSaved hook. + var fed *federation.Federator + if cfg.fedCfg.Enabled() { + fed = federation.NewFederator(cfg.fedCfg, func(ctx context.Context, evt nostr.Event) bool { + // Ensure the channel exists locally so federated chat has a home, then + // run it through the normal pipeline (OnEvent → store → OnEventSaved). + // AddEvent returns skipBroadcast=true on a duplicate: that is our loop + // guard — a re-seen event never reaches the egress hook again. + groups.EnsureFederatedGroup(group.GroupIDFromEvent(evt)) + skip, err := relay.AddEvent(ctx, evt) + if err != nil || skip { + return false + } + // AddEvent stores but does not broadcast; deliver to local subscribers. + relay.BroadcastEvent(evt) + return true + }) + } + + // After a group-mutating event is stored, (re)publish the affected NIP-29 + // addressable lists (39000 metadata, 39001 admins, 39002 members) so clients + // can discover the channel and track who may moderate it. + publish := func(evt nostr.Event, ok bool) { + if ok { + relay.BroadcastEvent(evt) + _ = cfg.store.SaveEvent(evt) + } + } + relay.OnEventSaved = func(ctx context.Context, evt nostr.Event) { + id := group.GroupIDFromEvent(evt) + if id == "" { + return + } + switch evt.Kind { + case nostr.KindSimpleGroupCreateGroup: + publish(groups.MetadataEvent(ctx, id)) + publish(groups.AdminsEvent(id)) + publish(groups.MembersEvent(id)) + case nostr.KindSimpleGroupEditMetadata: + publish(groups.MetadataEvent(ctx, id)) + case nostr.KindSimpleGroupJoinRequest, + nostr.KindSimpleGroupLeaveRequest, + nostr.KindSimpleGroupPutUser, + nostr.KindSimpleGroupRemoveUser: + // Membership (and possibly admin) changed → refresh both lists. + publish(groups.AdminsEvent(id)) + publish(groups.MembersEvent(id)) + case nostr.KindSimpleGroupChatMessage, nostr.KindSimpleGroupThreadedReply: + // Egress: forward newly-saved chat to mirror peers. A chat event + // re-injected from a peer is a store duplicate and never reaches here + // (AddEvent stops on dedup before OnEventSaved), so this can't loop. + if fed != nil { + fed.Forward(evt) + } + } + } + + // Serve NIP-11 ourselves so we can add the nostermd `features` array + // (the upstream document type has no field for it), then dispatch the real + // protocol traffic to khatru's specific handlers. We must NOT call + // relay.ServeHTTP here: khatru's ServeHTTP dispatches through this same mux, + // so routing "/" back into it recurses infinitely (stack overflow on the + // first non-NIP-11 request). Call the concrete handlers directly instead. + mux := relay.Router() + federated := cfg.fedCfg.Enabled() + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + switch { + case r.Header.Get("Accept") == "application/nostr+json": + serveNIP11(w, cfg.name, cfg.description, cfg.contact, cfg.motd, pk, federated) + case r.Header.Get("Upgrade") == "websocket": + relay.HandleWebsocket(w, r) + case r.Header.Get("Content-Type") == "application/nostr+json+rpc": + relay.HandleNIP86(w, r) + default: + http.Error(w, "expected a nostr relay connection (WebSocket) or NIP-11 request", http.StatusUpgradeRequired) + } + }) + + return &relayApp{ + relay: relay, + handler: mux, + groups: groups, + federator: fed, + pubKey: pk, + } +} + +// serveNIP11 emits a standard NIP-11 document merged with the Nosterm handshake +// keys (software/version/features), plus an optional `motd` the client shows in +// the relay's server window. +func serveNIP11(w http.ResponseWriter, name, description, contact, motd string, pk nostr.PubKey, federated bool) { + doc := map[string]any{ + "name": name, + "description": description, + "contact": contact, + "pubkey": pk.Hex(), + "supported_nips": []int{1, 11, 29, 42}, + } + if motd != "" { + doc["motd"] = motd + } + for k, v := range relayinfo.Extras(federated) { + doc[k] = v + } + w.Header().Set("Content-Type", "application/nostr+json") + w.Header().Set("Access-Control-Allow-Origin", "*") + _ = json.NewEncoder(w).Encode(doc) +} diff --git a/cmd/nostermd/relay_test.go b/cmd/nostermd/relay_test.go new file mode 100644 index 0000000..7d2d9c7 --- /dev/null +++ b/cmd/nostermd/relay_test.go @@ -0,0 +1,194 @@ +package main + +import ( + "context" + "net/http/httptest" + "testing" + "time" + + "fiatjaf.com/nostr" + "fiatjaf.com/nostr/eventstore/slicestore" + + "github.com/nosterm/relay/internal/federation" +) + +// startTestRelay spins up a fully-wired relay backed by an in-memory store, +// federating per peersSpec (empty = no federation). It returns the relay's ws:// +// URL. The federator, if any, is started and cleaned up automatically. This is +// the integration seam: it exercises exactly the wiring buildRelay produces for +// main(), over a real in-process khatru websocket server. +func startTestRelay(t *testing.T, peersSpec string) string { + t.Helper() + + store := &slicestore.SliceStore{} + if err := store.Init(); err != nil { + t.Fatalf("store init: %v", err) + } + fedCfg, _ := federation.ParseConfig(peersSpec) + + app := buildRelay(relayConfig{ + secretKey: nostr.Generate(), + store: store, + name: "test-relay", + fedCfg: fedCfg, + }) + + srv := httptest.NewServer(app.handler) + t.Cleanup(srv.Close) + if app.federator != nil { + app.federator.Start() + t.Cleanup(app.federator.Stop) + } + return "ws" + srv.URL[len("http"):] +} + +func mustSign(t *testing.T, sk nostr.SecretKey, evt *nostr.Event) { + t.Helper() + if err := evt.Sign(sk); err != nil { + t.Fatalf("sign: %v", err) + } +} + +func hpk(sk nostr.SecretKey) nostr.PubKey { return nostr.GetPublicKey(sk) } + +func TestGroupJoinSendReceive(t *testing.T) { + url := startTestRelay(t, "") + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + // Two independent identities: alice writes, bob reads. + aliceSK := nostr.Generate() + bobSK := nostr.Generate() + const groupID = "general" + + // Bob connects and subscribes to the group's chat messages first. + bob, err := nostr.RelayConnect(ctx, url, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("bob connect: %v", err) + } + sub, err := bob.Subscribe(ctx, nostr.Filter{ + Kinds: []nostr.Kind{nostr.KindSimpleGroupChatMessage}, + Tags: nostr.TagMap{"h": []string{groupID}}, + }, nostr.SubscriptionOptions{}) + if err != nil { + t.Fatalf("bob subscribe: %v", err) + } + + // Alice connects, joins the group, then sends a chat message. + alice, err := nostr.RelayConnect(ctx, url, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("alice connect: %v", err) + } + + join := nostr.Event{ + Kind: nostr.KindSimpleGroupJoinRequest, + CreatedAt: nostr.Now(), + Tags: nostr.Tags{{"h", groupID}}, + } + mustSign(t, aliceSK, &join) + if err := alice.Publish(ctx, join); err != nil { + t.Fatalf("alice join publish: %v", err) + } + + msg := nostr.Event{ + Kind: nostr.KindSimpleGroupChatMessage, + CreatedAt: nostr.Now(), + Content: "hello nosterm relay", + Tags: nostr.Tags{{"h", groupID}}, + } + mustSign(t, aliceSK, &msg) + if err := alice.Publish(ctx, msg); err != nil { + t.Fatalf("alice message publish: %v", err) + } + + // Bob must receive alice's message via the subscription (cross-user). + select { + case got := <-sub.Events: + if got.Content != "hello nosterm relay" { + t.Fatalf("unexpected content: %q", got.Content) + } + if got.PubKey != nostr.GetPublicKey(aliceSK) { + t.Fatalf("unexpected author") + } + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for the group message") + } + + _ = bobSK // bob reads without needing to sign in this open-group MVP +} + +// End-to-end over a real websocket: an admin and two members interact, and the +// relay rejects a non-admin's moderation attempt on the wire. +// +// NOTE: run the suite without -race. go-nostr's unsafe-based JSON serializer +// trips the race detector's checkptr under this test's concurrent signing; it +// is an upstream library issue, not a data race in the relay. The relay's own +// group-state locking is race-clean — verify with: +// +// go test ./internal/group/ -race +func TestMultiUserAdminOverWebsocket(t *testing.T) { + url := startTestRelay(t, "") + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + admin := nostr.Generate() + alice := nostr.Generate() + const g = "ops" + + adminConn, err := nostr.RelayConnect(ctx, url, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("admin connect: %v", err) + } + aliceConn, err := nostr.RelayConnect(ctx, url, nostr.RelayOptions{}) + if err != nil { + t.Fatalf("alice connect: %v", err) + } + + publish := func(conn *nostr.Relay, sk nostr.SecretKey, kind nostr.Kind, tags nostr.Tags) error { + evt := nostr.Event{Kind: kind, CreatedAt: nostr.Now(), Tags: tags} + mustSign(t, sk, &evt) + return conn.Publish(ctx, evt) + } + + // Admin creates a closed group. + if err := publish(adminConn, admin, nostr.KindSimpleGroupCreateGroup, nostr.Tags{{"h", g}, {"closed"}}); err != nil { + t.Fatalf("create closed group: %v", err) + } + + // Alice (not yet a member) cannot post — the relay rejects it on the wire. + if err := publish(aliceConn, alice, nostr.KindSimpleGroupChatMessage, nostr.Tags{{"h", g}}); err == nil { + t.Fatal("expected the relay to reject a non-member's post to a closed group") + } + + // Alice cannot moderate (add herself) — she is not an admin. + if err := publish(aliceConn, alice, nostr.KindSimpleGroupPutUser, nostr.Tags{{"h", g}, {"p", hpk(alice).Hex()}}); err == nil { + t.Fatal("expected the relay to reject a non-admin moderation event") + } + + // The admin admits alice. + if err := publish(adminConn, admin, nostr.KindSimpleGroupPutUser, nostr.Tags{{"h", g}, {"p", hpk(alice).Hex()}}); err != nil { + t.Fatalf("admin add user: %v", err) + } + + // Alice can now post, and the admin — subscribed — receives it. + sub, err := adminConn.Subscribe(ctx, nostr.Filter{ + Kinds: []nostr.Kind{nostr.KindSimpleGroupChatMessage}, + Tags: nostr.TagMap{"h": []string{g}}, + }, nostr.SubscriptionOptions{}) + if err != nil { + t.Fatalf("admin subscribe: %v", err) + } + + if err := publish(aliceConn, alice, nostr.KindSimpleGroupChatMessage, nostr.Tags{{"h", g}}); err != nil { + t.Fatalf("member post after admission should succeed: %v", err) + } + + select { + case got := <-sub.Events: + if got.PubKey != hpk(alice) { + t.Fatal("admin received a message from an unexpected author") + } + case <-time.After(5 * time.Second): + t.Fatal("admin did not receive the admitted member's message") + } +} diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..1af2737 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,44 @@ +# Standalone Nosterm home relay deployment. +# +# Runs ONLY the relay: a minimal NIP-29 managed-group Nostr relay with the +# Nosterm capability handshake. It listens on port 3334 and persists its BoltDB +# event store to a named volume so rooms/messages survive restarts. +# +# cp .env.example .env # then set RELAY_SECRET_KEY (see .env.example) +# docker compose up -d --build + +services: + relay: + build: + context: . + image: nostermd:latest + environment: + # Stable 64-char hex identity that signs group metadata/rosters. REQUIRED + # when RELAY_ENV=production (the relay refuses to start without it, since a + # fresh key each boot corrupts signed rosters). In dev it is persisted to + # the data volume automatically. Generate with: openssl rand -hex 32 + RELAY_SECRET_KEY: ${RELAY_SECRET_KEY:-} + RELAY_ENV: ${RELAY_ENV:-} + RELAY_NAME: ${RELAY_NAME:-Nosterm Home Relay} + RELAY_DESCRIPTION: ${RELAY_DESCRIPTION:-} + RELAY_CONTACT: ${RELAY_CONTACT:-relay@nosterm.com} + RELAY_MOTD: ${RELAY_MOTD:-} + # Retention: prune old CHAT messages so the store doesn't grow unbounded. + # 0 = unlimited. Group membership/metadata is never pruned. + RELAY_RETENTION_DAYS: ${RELAY_RETENTION_DAYS:-0} + RELAY_RETENTION_MAX_MESSAGES: ${RELAY_RETENTION_MAX_MESSAGES:-0} + RELAY_RETENTION_INTERVAL_MINUTES: ${RELAY_RETENTION_INTERVAL_MINUTES:-60} + # Optional federation: mirror/ingest NIP-29 CHAT with peer relays. Empty = + # disabled. Format: "url|mode:channels" entries separated by ';', where + # mode is mirror (bidirectional) or ingest (read-only) and channels is a + # comma list or '*'. e.g. + # wss://peer-a.example|mirror:general,dev;wss://peer-b.example|ingest:* + RELAY_FEDERATION_PEERS: ${RELAY_FEDERATION_PEERS:-} + volumes: + - relay-data:/data + ports: + - '3334:3334' + restart: unless-stopped + +volumes: + relay-data: diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..aa31137 --- /dev/null +++ b/go.mod @@ -0,0 +1,39 @@ +module github.com/nosterm/relay + +go 1.25.11 + +require fiatjaf.com/nostr v0.0.0-20260716191248-c205ed45b97e + +require ( + fiatjaf.com/lib v0.3.7 // indirect + github.com/ImVexed/fasturl v0.0.0-20230304231329-4e41488060f3 // indirect + github.com/andybalholm/brotli v1.1.1 // indirect + github.com/bep/debounce v1.2.1 // indirect + github.com/btcsuite/btcd/btcec/v2 v2.3.4 // indirect + github.com/btcsuite/btcd/chaincfg/chainhash v1.1.0 // indirect + github.com/coder/websocket v1.8.13 // indirect + github.com/decred/dcrd/crypto/blake256 v1.1.0 // indirect + github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.0 // indirect + github.com/fasthttp/websocket v1.5.12 // indirect + github.com/josharian/intern v1.0.0 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/klauspost/compress v1.18.0 // indirect + github.com/mailru/easyjson v0.9.0 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect + github.com/puzpuzpuz/xsync/v3 v3.5.1 // indirect + github.com/rs/cors v1.11.1 // indirect + github.com/savsgio/gotils v0.0.0-20240704082632-aef3928b8a38 // indirect + github.com/templexxx/cpu v0.0.1 // indirect + github.com/templexxx/xhex v0.0.0-20200614015412-aed53437177b // indirect + github.com/tidwall/gjson v1.18.0 // indirect + github.com/tidwall/match v1.1.1 // indirect + github.com/tidwall/pretty v1.2.1 // indirect + github.com/valyala/bytebufferpool v1.0.0 // indirect + github.com/valyala/fasthttp v1.59.0 // indirect + go.etcd.io/bbolt v1.4.2 // indirect + golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 // indirect + golang.org/x/net v0.41.0 // indirect + golang.org/x/sync v0.15.0 // indirect + golang.org/x/sys v0.35.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..fbd30a6 --- /dev/null +++ b/go.sum @@ -0,0 +1,92 @@ +fiatjaf.com/lib v0.3.7 h1:mXZOn7NrUcjSdy4oNvwQyAmes7Ueb+Zr5hjqMIe2dxI= +fiatjaf.com/lib v0.3.7/go.mod h1:UlHaZvPHj25PtKLh9GjZkUHRmQ2xZ8Jkoa4VRaLeeQ8= +fiatjaf.com/nostr v0.0.0-20260716191248-c205ed45b97e h1:/vgoytiH4qQ28O/oX1hj656X69CnTy1pGj/17HFjH+A= +fiatjaf.com/nostr v0.0.0-20260716191248-c205ed45b97e/go.mod h1:b1EIUDnd133Ie8Pg8O/biaKdFyCMz28aD4n64g1GqvM= +github.com/ImVexed/fasturl v0.0.0-20230304231329-4e41488060f3 h1:ClzzXMDDuUbWfNNZqGeYq4PnYOlwlOVIvSyNaIy0ykg= +github.com/ImVexed/fasturl v0.0.0-20230304231329-4e41488060f3/go.mod h1:we0YA5CsBbH5+/NUzC/AlMmxaDtWlXeNsqrwXjTzmzA= +github.com/PowerDNS/lmdb-go v1.9.3 h1:AUMY2pZT8WRpkEv39I9Id3MuoHd+NZbTVpNhruVkPTg= +github.com/PowerDNS/lmdb-go v1.9.3/go.mod h1:TE0l+EZK8Z1B4dx070ZxkWTlp8RG1mjN0/+FkFRQMtU= +github.com/andybalholm/brotli v1.1.1 h1:PR2pgnyFznKEugtsUo0xLdDop5SKXd5Qf5ysW+7XdTA= +github.com/andybalholm/brotli v1.1.1/go.mod h1:05ib4cKhjx3OQYUY22hTVd34Bc8upXjOLL2rKwwZBoA= +github.com/bep/debounce v1.2.1 h1:v67fRdBA9UQu2NhLFXrSg0Brw7CexQekrBwDMM8bzeY= +github.com/bep/debounce v1.2.1/go.mod h1:H8yggRPQKLUhUoqrJC1bO2xNya7vanpDl7xR3ISbCJ0= +github.com/btcsuite/btcd/btcec/v2 v2.3.4 h1:3EJjcN70HCu/mwqlUsGK8GcNVyLVxFDlWurTXGPFfiQ= +github.com/btcsuite/btcd/btcec/v2 v2.3.4/go.mod h1:zYzJ8etWJQIv1Ogk7OzpWjowwOdXY1W/17j2MW85J04= +github.com/btcsuite/btcd/chaincfg/chainhash v1.1.0 h1:59Kx4K6lzOW5w6nFlA0v5+lk/6sjybR934QNHSJZPTQ= +github.com/btcsuite/btcd/chaincfg/chainhash v1.1.0/go.mod h1:7SFka0XMvUgj3hfZtydOrQY2mwhPclbT2snogU7SQQc= +github.com/coder/websocket v1.8.13 h1:f3QZdXy7uGVz+4uCJy2nTZyM0yTBj8yANEHhqlXZ9FE= +github.com/coder/websocket v1.8.13/go.mod h1:LNVeNrXQZfe5qhS9ALED3uA+l5pPqvwXg3CKoDBB2gs= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/decred/dcrd/crypto/blake256 v1.1.0 h1:zPMNGQCm0g4QTY27fOCorQW7EryeQ/U0x++OzVrdms8= +github.com/decred/dcrd/crypto/blake256 v1.1.0/go.mod h1:2OfgNZ5wDpcsFmHmCK5gZTPcCXqlm2ArzUIkw9czNJo= +github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.0 h1:NMZiJj8QnKe1LgsbDayM4UoHwbvwDRwnI3hwNaAHRnc= +github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.0/go.mod h1:ZXNYxsqcloTdSy/rNShjYzMhyjf0LaoftYK0p+A3h40= +github.com/dvyukov/go-fuzz v0.0.0-20200318091601-be3528f3a813/go.mod h1:11Gm+ccJnvAhCNLlf5+cS9KjtbaD5I5zaZpFMsTHWTw= +github.com/fasthttp/websocket v1.5.12 h1:e4RGPpWW2HTbL3zV0Y/t7g0ub294LkiuXXUuTOUInlE= +github.com/fasthttp/websocket v1.5.12/go.mod h1:I+liyL7/4moHojiOgUOIKEWm9EIxHqxZChS+aMFltyg= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/josharian/intern v1.0.0 h1:vlS4z54oSdjm0bgjRigI+G1HpF+tI+9rE5LLzOg8HmY= +github.com/josharian/intern v1.0.0/go.mod h1:5DoeVV0s6jJacbCEi61lwdGj/aVlrQvzHFFd8Hwg//Y= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/mailru/easyjson v0.9.0 h1:PrnmzHw7262yW8sTBwxi1PdJA3Iw/EKBa8psRf7d9a4= +github.com/mailru/easyjson v0.9.0/go.mod h1:1+xMtQp2MRNVL/V1bOzuP3aP8VNwRW55fQUto+XFtTU= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/puzpuzpuz/xsync/v3 v3.5.1 h1:GJYJZwO6IdxN/IKbneznS6yPkVC+c3zyY/j19c++5Fg= +github.com/puzpuzpuz/xsync/v3 v3.5.1/go.mod h1:VjzYrABPabuM4KyBh1Ftq6u8nhwY5tBPKP9jpmh0nnA= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/rs/cors v1.11.1 h1:eU3gRzXLRK57F5rKMGMZURNdIG4EoAmX8k94r9wXWHA= +github.com/rs/cors v1.11.1/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU= +github.com/savsgio/gotils v0.0.0-20240704082632-aef3928b8a38 h1:D0vL7YNisV2yqE55+q0lFuGse6U8lxlg7fYTctlT5Gc= +github.com/savsgio/gotils v0.0.0-20240704082632-aef3928b8a38/go.mod h1:sM7Mt7uEoCeFSCBM+qBrqvEo+/9vdmj19wzp3yzUhmg= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/templexxx/cpu v0.0.1 h1:hY4WdLOgKdc8y13EYklu9OUTXik80BkxHoWvTO6MQQY= +github.com/templexxx/cpu v0.0.1/go.mod h1:w7Tb+7qgcAlIyX4NhLuDKt78AHA5SzPmq0Wj6HiEnnk= +github.com/templexxx/xhex v0.0.0-20200614015412-aed53437177b h1:XeDLE6c9mzHpdv3Wb1+pWBaWv/BlHK0ZYIu/KaL6eHg= +github.com/templexxx/xhex v0.0.0-20200614015412-aed53437177b/go.mod h1:7rwmCH0wC2fQvNEvPZ3sKXukhyCTyiaZ5VTZMQYpZKQ= +github.com/tidwall/gjson v1.18.0 h1:FIDeeyB800efLX89e5a8Y0BNH+LOngJyGrIWxG2FKQY= +github.com/tidwall/gjson v1.18.0/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= +github.com/tidwall/match v1.1.1 h1:+Ho715JplO36QYgwN9PGYNhgZvoUSc9X2c80KVTi+GA= +github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= +github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= +github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= +github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= +github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= +github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= +github.com/valyala/fasthttp v1.59.0 h1:Qu0qYHfXvPk1mSLNqcFtEk6DpxgA26hy6bmydotDpRI= +github.com/valyala/fasthttp v1.59.0/go.mod h1:GTxNb9Bc6r2a9D0TWNSPwDz78UxnTGBViY3xZNEqyYU= +github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU= +github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= +go.etcd.io/bbolt v1.4.2 h1:IrUHp260R8c+zYx/Tm8QZr04CX+qWS5PGfPdevhdm1I= +go.etcd.io/bbolt v1.4.2/go.mod h1:Is8rSHO/b4f3XigBC0lL0+4FwAQv3HXEEIgFMuKHceM= +golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 h1:nDVHiLt8aIbd/VzvPWN6kSOPE7+F/fNFDSXLVYkE/Iw= +golang.org/x/exp v0.0.0-20250305212735-054e65f0b394/go.mod h1:sIifuuw/Yco/y6yb6+bDNfyeQ/MdPUy/hKEMYQV17cM= +golang.org/x/net v0.41.0 h1:vBTly1HeNPEn3wtREYfy4GZ/NECgw2Cnl+nK6Nz3uvw= +golang.org/x/net v0.41.0/go.mod h1:B/K4NNqkfmg07DQYrbwvSluqCJOOXwUjeb/5lOisjbA= +golang.org/x/sync v0.15.0 h1:KWH3jNZsfyT6xfAfKiz6MRNmd46ByHDYaZ7KSkCtdW8= +golang.org/x/sync v0.15.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI= +golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/federation/federation.go b/internal/federation/federation.go new file mode 100644 index 0000000..eacd45d --- /dev/null +++ b/internal/federation/federation.go @@ -0,0 +1,406 @@ +// Package federation lets the relay mirror NIP-29 chat with peer relays, so a +// channel can live on more than one relay and survive any single one going down. +// It is OPTIONAL and chat-only by design: only kind-9/10 messages cross the +// boundary. Membership, admins, bans and metadata stay local to each relay (each +// operator moderates their own instance) — this sidesteps cross-relay moderation +// conflicts entirely. +// +// Two per-channel modes: +// +// - mirror bidirectional. Chat posted on either relay is forwarded to the +// other, so the channel is fully replicated. +// - ingest read-only. This relay pulls chat FROM the peer but never pushes +// back (announcement feeds, syndication, public-room mirrors). +// +// Loop prevention is structural rather than tag-based: an incoming federated +// event is injected through the relay's normal add pipeline, whose eventstore +// rejects duplicates (ErrDupEvent) BEFORE the egress hook runs. So an event that +// has already been seen is never re-forwarded, and A↔B↔A cycles die on the +// second sight. No custom "seen" bookkeeping to get wrong. +package federation + +import ( + "context" + "log" + "strings" + "sync" + "time" + + "fiatjaf.com/nostr" + + "github.com/nosterm/relay/internal/group" +) + +// Mode is a peer's federation direction. +type Mode int + +const ( + // ModeMirror is bidirectional federation. + ModeMirror Mode = iota + // ModeIngest is read-only (pull-only) federation. + ModeIngest +) + +func (m Mode) String() string { + if m == ModeIngest { + return "ingest" + } + return "mirror" +} + +// Peer is one configured peer relay and the channels federated with it. +type Peer struct { + URL string + Mode Mode + Channels map[string]bool // group ids; empty map with AllChans=true means "all" + AllChans bool +} + +// wants reports whether events for the given group id federate with this peer. +func (p *Peer) wants(groupID string) bool { + if p.AllChans { + return true + } + return p.Channels[groupID] +} + +// Config is the parsed set of peers. +type Config struct { + Peers []Peer +} + +// Enabled reports whether any peer is configured. +func (c Config) Enabled() bool { + return len(c.Peers) > 0 +} + +// FederatesChannel reports whether any peer federates the given group id. Used +// by group evaluation to accept federated chat regardless of local membership +// (a federated channel is effectively open on every participating relay). +func (c Config) FederatesChannel(id string) bool { + for i := range c.Peers { + if c.Peers[i].wants(id) { + return true + } + } + return false +} + +// ParseConfig parses the RELAY_FEDERATION_PEERS specification. +// +// Grammar (peers separated by ';', ignoring surrounding whitespace): +// +// |: +// +// where is "mirror" or "ingest" and is a comma-separated list +// of group ids, or "*" for every channel. The "|:" suffix is +// optional and defaults to "mirror:*". Blank entries are skipped. +// +// Examples: +// +// wss://a.example +// wss://a.example|mirror:general,dev +// wss://a.example|mirror:general;wss://b.example|ingest:announcements +// wss://a.example|ingest:* +func ParseConfig(spec string) (Config, []string) { + var cfg Config + var warnings []string + + for _, raw := range strings.Split(spec, ";") { + entry := strings.TrimSpace(raw) + if entry == "" { + continue + } + + urlPart := entry + modePart := "" + if i := strings.Index(entry, "|"); i >= 0 { + urlPart = strings.TrimSpace(entry[:i]) + modePart = strings.TrimSpace(entry[i+1:]) + } + + // Require an explicit ws/wss scheme on the RAW input: NormalizeURL would + // silently rewrite http→ws and bare hosts→wss, masking config mistakes. + lower := strings.ToLower(urlPart) + if !(strings.HasPrefix(lower, "ws://") || strings.HasPrefix(lower, "wss://")) { + warnings = append(warnings, "ignoring peer with non-ws(s) URL: "+urlPart) + continue + } + url := nostr.NormalizeURL(urlPart) + if url == "" { + warnings = append(warnings, "ignoring peer with invalid URL: "+urlPart) + continue + } + + peer := Peer{ + URL: url, + Mode: ModeMirror, + Channels: map[string]bool{}, + AllChans: true, // default: all channels + } + + if modePart != "" { + modeName := modePart + chanList := "*" + if j := strings.Index(modePart, ":"); j >= 0 { + modeName = strings.TrimSpace(modePart[:j]) + chanList = strings.TrimSpace(modePart[j+1:]) + } + switch strings.ToLower(modeName) { + case "mirror", "": + peer.Mode = ModeMirror + case "ingest", "pull", "read", "readonly": + peer.Mode = ModeIngest + default: + warnings = append(warnings, "unknown federation mode "+modeName+" for "+url+"; using mirror") + peer.Mode = ModeMirror + } + if chanList != "" && chanList != "*" { + peer.AllChans = false + for _, c := range strings.Split(chanList, ",") { + id := strings.TrimSpace(strings.TrimPrefix(c, "#")) + if id != "" { + peer.Channels[id] = true + } + } + // A mode with an explicit but empty channel list federates nothing. + if len(peer.Channels) == 0 { + warnings = append(warnings, "peer "+url+" lists no valid channels; skipping") + continue + } + } + } + + cfg.Peers = append(cfg.Peers, peer) + } + + return cfg, warnings +} + +// EventInjector accepts a federated event into the local relay pipeline. It +// returns whether the event was newly stored (false for duplicates/rejects), so +// tests can assert loop prevention without a live relay. In production this is +// backed by khatru's Relay.AddEvent. +type EventInjector func(ctx context.Context, evt nostr.Event) (stored bool) + +// Federator owns the live peer connections and the forwarding logic. +type Federator struct { + cfg Config + inject EventInjector + dialCtx context.Context + cancel context.CancelFunc + + mu sync.Mutex + connected map[string]*nostr.Relay // url -> live client connection + + wg sync.WaitGroup +} + +// federationKinds are the only kinds we ever subscribe to / forward: chat only. +var federationKinds = []nostr.Kind{ + nostr.KindSimpleGroupChatMessage, // 9 + nostr.KindSimpleGroupThreadedReply, // 10 +} + +// NewFederator creates a federator for the given config. inject feeds received +// peer events into the local relay pipeline. +func NewFederator(cfg Config, inject EventInjector) *Federator { + ctx, cancel := context.WithCancel(context.Background()) + return &Federator{ + cfg: cfg, + inject: inject, + dialCtx: ctx, + cancel: cancel, + connected: map[string]*nostr.Relay{}, + } +} + +// Start dials every peer and begins pulling their chat. Non-blocking: each peer +// runs in its own goroutine that reconnects with backoff until Stop is called. +func (f *Federator) Start() { + for i := range f.cfg.Peers { + peer := f.cfg.Peers[i] + f.wg.Add(1) + go func() { + defer f.wg.Done() + f.runPeer(peer) + }() + } + log.Printf("federation: started with %d peer(s)", len(f.cfg.Peers)) +} + +// Stop tears down all peer connections and waits for their goroutines to exit. +func (f *Federator) Stop() { + f.cancel() + f.mu.Lock() + for _, r := range f.connected { + _ = r.Close() + } + f.mu.Unlock() + f.wg.Wait() +} + +// runPeer maintains one peer connection, reconnecting with capped backoff. Each +// (re)connection opens a fresh subscription for the federated chat kinds and +// injects received events into the local pipeline (dedup handles loops). +func (f *Federator) runPeer(peer Peer) { + backoff := time.Second + const maxBackoff = 30 * time.Second + + for f.dialCtx.Err() == nil { + if err := f.connectAndPull(peer); err != nil { + if f.dialCtx.Err() != nil { + return + } + log.Printf("federation: peer %s (%s) error: %v; retrying in %s", + peer.URL, peer.Mode, err, backoff) + select { + case <-f.dialCtx.Done(): + return + case <-time.After(backoff): + } + backoff *= 2 + if backoff > maxBackoff { + backoff = maxBackoff + } + continue + } + // Clean subscription end (peer closed / context canceled) — reset backoff + // and try to re-establish unless we're shutting down. + backoff = time.Second + } +} + +// connectAndPull opens a connection + subscription and blocks pumping events +// until the subscription ends or the connection drops. Returns nil on a clean +// end (so the caller re-establishes), or an error to trigger backoff. +func (f *Federator) connectAndPull(peer Peer) error { + relay, err := nostr.RelayConnect(f.dialCtx, peer.URL, nostr.RelayOptions{}) + if err != nil { + return err + } + f.mu.Lock() + f.connected[peer.URL] = relay + f.mu.Unlock() + defer func() { + f.mu.Lock() + delete(f.connected, peer.URL) + f.mu.Unlock() + _ = relay.Close() + }() + + filter := nostr.Filter{Kinds: federationKinds} + // Scope the pull to the configured channels when not "all", so we don't drag + // down chat for channels this peer isn't federating. + if !peer.AllChans { + ids := make([]string, 0, len(peer.Channels)) + for id := range peer.Channels { + ids = append(ids, id) + } + filter.Tags = nostr.TagMap{"h": ids} + } + + sub, err := relay.Subscribe(f.dialCtx, filter, nostr.SubscriptionOptions{Label: "federation"}) + if err != nil { + return err + } + defer sub.Unsub() + + log.Printf("federation: pulling %s from %s", peer.Mode, peer.URL) + + for { + select { + case <-f.dialCtx.Done(): + return nil + case <-sub.EndOfStoredEvents: + // Historical backfill done; keep receiving live events. + case reason := <-sub.ClosedReason: + log.Printf("federation: peer %s closed subscription: %s", peer.URL, reason) + return nil + case evt, ok := <-sub.Events: + if !ok { + return nil + } + f.receive(peer, evt) + } + } +} + +// receive validates and injects an event pulled from a peer. Signature is +// re-verified defensively (never trust a peer to have done so), the channel +// scope is re-checked, and the event is fed through the local pipeline where +// the eventstore dedups it — that dedup is what prevents forwarding loops. +func (f *Federator) receive(peer Peer, evt nostr.Event) { + id := group.GroupIDFromEvent(evt) + if id == "" || !peer.wants(id) { + return + } + if !evt.VerifySignature() { + log.Printf("federation: dropped event %s from %s (bad signature)", evt.ID.Hex(), peer.URL) + return + } + f.inject(f.dialCtx, evt) +} + +// Forward pushes a locally-saved chat event out to every peer that federates its +// channel in a writable (mirror) mode. Ingest peers are pull-only and are +// skipped. Called from the relay's OnEventSaved hook for chat kinds; a duplicate +// re-injected from a peer never reaches here because AddEvent stops on dedup +// before OnEventSaved fires. +func (f *Federator) Forward(evt nostr.Event) { + id := group.GroupIDFromEvent(evt) + if id == "" { + return + } + f.mu.Lock() + defer f.mu.Unlock() + for i := range f.cfg.Peers { + peer := f.cfg.Peers[i] + if peer.Mode != ModeMirror || !peer.wants(id) { + continue + } + relay := f.connected[peer.URL] + if relay == nil || !relay.IsConnected() { + continue // reconnect logic will backfill via subscription + } + // Publish in the background so a slow peer can't block the save path. + f.wg.Add(1) + go func(r *nostr.Relay, url string) { + defer f.wg.Done() + ctx, cancel := context.WithTimeout(f.dialCtx, 10*time.Second) + defer cancel() + if err := r.Publish(ctx, evt); err != nil { + log.Printf("federation: publish to %s failed: %v", url, err) + } + }(relay, peer.URL) + } +} + +// Summary is a short human-readable description for logging/inspection. +func (c Config) Summary() string { + if !c.Enabled() { + return "federation disabled" + } + var b strings.Builder + b.WriteString("federation: ") + for i, p := range c.Peers { + if i > 0 { + b.WriteString(", ") + } + scope := "*" + if !p.AllChans { + ids := make([]string, 0, len(p.Channels)) + for id := range p.Channels { + ids = append(ids, id) + } + scope = strings.Join(ids, "/") + } + b.WriteString(p.URL) + b.WriteString("(") + b.WriteString(p.Mode.String()) + b.WriteString(":") + b.WriteString(scope) + b.WriteString(")") + } + return b.String() +} diff --git a/internal/federation/federation_test.go b/internal/federation/federation_test.go new file mode 100644 index 0000000..d03f269 --- /dev/null +++ b/internal/federation/federation_test.go @@ -0,0 +1,215 @@ +package federation + +import ( + "context" + "sync" + "testing" + + "fiatjaf.com/nostr" +) + +func TestParseConfig(t *testing.T) { + t.Run("empty spec disables federation", func(t *testing.T) { + cfg, warnings := ParseConfig("") + if cfg.Enabled() { + t.Fatal("expected federation disabled for empty spec") + } + if len(warnings) != 0 { + t.Fatalf("unexpected warnings: %v", warnings) + } + }) + + t.Run("bare url defaults to mirror all channels", func(t *testing.T) { + cfg, _ := ParseConfig("wss://a.example") + if len(cfg.Peers) != 1 { + t.Fatalf("expected 1 peer, got %d", len(cfg.Peers)) + } + p := cfg.Peers[0] + if p.Mode != ModeMirror { + t.Errorf("expected mirror mode, got %s", p.Mode) + } + if !p.AllChans { + t.Error("expected AllChans=true for bare url") + } + if !p.wants("anything") { + t.Error("bare url should federate every channel") + } + }) + + t.Run("explicit mode and channels", func(t *testing.T) { + cfg, _ := ParseConfig("wss://a.example|ingest:general,dev") + p := cfg.Peers[0] + if p.Mode != ModeIngest { + t.Errorf("expected ingest, got %s", p.Mode) + } + if p.AllChans { + t.Error("expected scoped channels, got AllChans") + } + if !p.wants("general") || !p.wants("dev") { + t.Error("should federate listed channels") + } + if p.wants("secret") { + t.Error("should NOT federate unlisted channel") + } + }) + + t.Run("multiple peers separated by semicolons", func(t *testing.T) { + cfg, _ := ParseConfig("wss://a.example|mirror:general; wss://b.example|ingest:*") + if len(cfg.Peers) != 2 { + t.Fatalf("expected 2 peers, got %d", len(cfg.Peers)) + } + if cfg.Peers[0].Mode != ModeMirror || cfg.Peers[1].Mode != ModeIngest { + t.Error("modes parsed incorrectly") + } + if !cfg.Peers[1].AllChans { + t.Error("ingest:* should mean all channels") + } + }) + + t.Run("channel names tolerate a leading #", func(t *testing.T) { + cfg, _ := ParseConfig("wss://a.example|mirror:#general") + if !cfg.Peers[0].wants("general") { + t.Error("should strip leading # from channel names") + } + }) + + t.Run("invalid url is skipped with a warning", func(t *testing.T) { + cfg, warnings := ParseConfig("http://not-websocket.example") + if cfg.Enabled() { + t.Error("non-ws url should be skipped") + } + if len(warnings) == 0 { + t.Error("expected a warning for invalid url") + } + }) + + t.Run("explicit but empty channel list is skipped", func(t *testing.T) { + cfg, warnings := ParseConfig("wss://a.example|mirror:,") + if cfg.Enabled() { + t.Error("empty channel list should skip the peer") + } + if len(warnings) == 0 { + t.Error("expected a warning for empty channel list") + } + }) + + t.Run("unknown mode falls back to mirror with a warning", func(t *testing.T) { + cfg, warnings := ParseConfig("wss://a.example|bogus:general") + if cfg.Peers[0].Mode != ModeMirror { + t.Error("unknown mode should fall back to mirror") + } + if len(warnings) == 0 { + t.Error("expected a warning for unknown mode") + } + }) +} + +func TestFederatesChannel(t *testing.T) { + cfg, _ := ParseConfig("wss://a.example|mirror:general;wss://b.example|ingest:dev") + if !cfg.FederatesChannel("general") { + t.Error("general federates via peer a") + } + if !cfg.FederatesChannel("dev") { + t.Error("dev federates via peer b") + } + if cfg.FederatesChannel("private") { + t.Error("private federates with nobody") + } +} + +// TestReceiveLoopPrevention proves the loop guard: an event injected twice is +// only stored (and thus only forwardable) once, because the injector reports the +// second sight as not-newly-stored. This mirrors AddEvent's ErrDupEvent path. +func TestReceiveLoopPrevention(t *testing.T) { + seen := map[nostr.ID]bool{} + var mu sync.Mutex + injectCount := 0 + storedCount := 0 + + inject := func(_ context.Context, evt nostr.Event) bool { + mu.Lock() + defer mu.Unlock() + injectCount++ + if seen[evt.ID] { + return false // duplicate: not newly stored (the loop breaker) + } + seen[evt.ID] = true + storedCount++ + return true + } + + cfg, _ := ParseConfig("wss://peer.example|mirror:general") + f := NewFederator(cfg, inject) + peer := cfg.Peers[0] + + sk := nostr.Generate() + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupChatMessage, + CreatedAt: nostr.Now(), + Content: "hello federation", + Tags: nostr.Tags{{"h", "general"}}, + } + if err := evt.Sign(sk); err != nil { + t.Fatalf("sign: %v", err) + } + + // Same event received twice (e.g. A→B then echoed B→A→B). + f.receive(peer, evt) + f.receive(peer, evt) + + if injectCount != 2 { + t.Errorf("expected 2 inject attempts, got %d", injectCount) + } + if storedCount != 1 { + t.Errorf("loop guard failed: event stored %d times, want 1", storedCount) + } +} + +// TestReceiveDropsBadSignature ensures a peer can't inject forged events. +func TestReceiveDropsBadSignature(t *testing.T) { + injected := false + inject := func(_ context.Context, _ nostr.Event) bool { + injected = true + return true + } + cfg, _ := ParseConfig("wss://peer.example|mirror:general") + f := NewFederator(cfg, inject) + + // Unsigned (invalid) event tagged for a federated channel. + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupChatMessage, + CreatedAt: nostr.Now(), + Content: "forged", + Tags: nostr.Tags{{"h", "general"}}, + } + f.receive(cfg.Peers[0], evt) + if injected { + t.Error("event with an invalid signature must not be injected") + } +} + +// TestReceiveIgnoresUnfederatedChannel ensures scope is enforced on ingress. +func TestReceiveIgnoresUnfederatedChannel(t *testing.T) { + injected := false + inject := func(_ context.Context, _ nostr.Event) bool { + injected = true + return true + } + cfg, _ := ParseConfig("wss://peer.example|mirror:general") + f := NewFederator(cfg, inject) + + sk := nostr.Generate() + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupChatMessage, + CreatedAt: nostr.Now(), + Content: "off-topic", + Tags: nostr.Tags{{"h", "other-channel"}}, + } + if err := evt.Sign(sk); err != nil { + t.Fatalf("sign: %v", err) + } + f.receive(cfg.Peers[0], evt) + if injected { + t.Error("chat for an unfederated channel must not be injected") + } +} diff --git a/internal/group/group.go b/internal/group/group.go new file mode 100644 index 0000000..4cb2e0e --- /dev/null +++ b/internal/group/group.go @@ -0,0 +1,431 @@ +// Package group implements the relay's in-memory NIP-29 managed-group state: the +// authority for channel membership, admin rights, and metadata. +// +// It is a deliberately small subset of NIP-29 (the managed-group MVP): +// +// kind 9007 create group +// kind 9021 join request -> adds a member (open groups) +// kind 9022 leave request -> removes a member +// kind 9000 put user (admin) -> adds a member +// kind 9001 remove user -> removes a member +// kind 9002 edit metadata -> admin sets name/about +// kind 9/10 chat / reply -> allowed for members (or anyone if open) +// +// Metadata (39000), admins (39001) and members (39002) are (re)published as +// signed addressable events whenever the group changes. State is rebuilt from +// the event store on boot so it survives restarts. +package group + +import ( + "cmp" + "context" + "fmt" + "slices" + "sync" + + "fiatjaf.com/nostr" + "fiatjaf.com/nostr/eventstore" +) + +// State is the relay's in-memory NIP-29 managed-group state. It is the authority +// for membership and metadata and is safe for concurrent use. +type State struct { + mu sync.RWMutex + groups map[string]*group // keyed by group id (the `h` tag value) + sk nostr.SecretKey // relay identity, signs 39000-39002 events + pk nostr.PubKey + + // IsFederated, when set, reports whether a channel is federated with a peer. + // Chat on a federated channel is accepted regardless of local membership, + // because the author is authorized on the peer relay that federates it. Nil + // (federation disabled) means "no channel is federated". + IsFederated func(groupID string) bool +} + +type group struct { + id string + name string + about string + creator nostr.PubKey + members map[nostr.PubKey]bool + admins map[nostr.PubKey]bool + // open groups auto-accept join requests and allow anyone to post (public + // channels). Closed groups are members-only: a join request does NOT grant + // membership — an admin must add the user (kind 9000) — and only members may + // post. A group is created closed when its 9007 event carries a `closed` tag. + open bool +} + +// NewState returns an empty group state that signs its addressable list events +// (39000-39002) with the given relay identity. +func NewState(sk nostr.SecretKey) *State { + return &State{ + groups: make(map[string]*group), + sk: sk, + pk: nostr.GetPublicKey(sk), + } +} + +// hasTag reports whether the event has a tag with the given name. Unlike +// Tags.Find (which ignores value-less tags), this matches flag tags such as +// ["closed"] that carry no value. +func hasTag(evt nostr.Event, name string) bool { + for _, t := range evt.Tags { + if len(t) >= 1 && t[0] == name { + return true + } + } + return false +} + +// managementKinds are the NIP-29 events that mutate group state and therefore +// must be replayed on boot to reconstruct membership. +var managementKinds = []nostr.Kind{ + nostr.KindSimpleGroupCreateGroup, + nostr.KindSimpleGroupJoinRequest, + nostr.KindSimpleGroupLeaveRequest, + nostr.KindSimpleGroupPutUser, + nostr.KindSimpleGroupRemoveUser, +} + +// Rebuild reconstructs in-memory group state from the persisted event store so +// membership and metadata survive relay restarts. Events are replayed in +// chronological (oldest-first) order so the final state is correct; Evaluate is +// reused so replay and live handling can never diverge. Returns the number of +// management events replayed. +func (gs *State) Rebuild(store eventstore.Store) int { + var events []nostr.Event + for evt := range store.QueryEvents(nostr.Filter{Kinds: managementKinds}, 100_000) { + events = append(events, evt) + } + // Apply oldest-first. On a created_at tie, a group-create (9007) must come + // first so later same-second events (joins/puts) find the group — otherwise + // membership can be lost across a restart when events share a timestamp. + slices.SortFunc(events, func(a, b nostr.Event) int { + if a.CreatedAt != b.CreatedAt { + return cmp.Compare(a.CreatedAt, b.CreatedAt) + } + aCreate := a.Kind == nostr.KindSimpleGroupCreateGroup + bCreate := b.Kind == nostr.KindSimpleGroupCreateGroup + if aCreate != bCreate { + if aCreate { + return -1 + } + return 1 + } + return 0 + }) + for _, evt := range events { + gs.Evaluate(evt) + } + return len(events) +} + +// GroupIDFromEvent returns the `h` tag value, or "" if absent. +func GroupIDFromEvent(evt nostr.Event) string { + t := evt.Tags.Find("h") + if len(t) >= 2 { + return t[1] + } + return "" +} + +// firstTagValue returns the value of the first tag with the given name, or "". +func firstTagValue(evt nostr.Event, name string) string { + t := evt.Tags.Find(name) + if len(t) >= 2 { + return t[1] + } + return "" +} + +// ensureGroup returns the group, creating one on first reference (channel +// auto-creation). New groups default to open unless created closed. `open` is +// only applied at creation time — it never downgrades an existing group. +func (gs *State) ensureGroup(id string, creator nostr.PubKey, open bool) *group { + g, ok := gs.groups[id] + if !ok { + g = &group{ + id: id, + name: id, + creator: creator, + members: map[nostr.PubKey]bool{}, + admins: map[nostr.PubKey]bool{creator: true}, + open: open, + } + g.members[creator] = true + gs.groups[id] = g + } + return g +} + +// Evaluate applies a management/chat event to group state and reports whether +// the relay should reject it. Called from khatru's OnEvent hook. +func (gs *State) Evaluate(evt nostr.Event) (reject bool, msg string) { + id := GroupIDFromEvent(evt) + + switch evt.Kind { + case nostr.KindSimpleGroupCreateGroup: // 9007 + if id == "" { + return true, "missing group id (h tag)" + } + // A `closed` tag makes the group members-only (admin-approved joins). + open := !hasTag(evt, "closed") + gs.mu.Lock() + gs.ensureGroup(id, evt.PubKey, open) + gs.mu.Unlock() + return false, "" + + case nostr.KindSimpleGroupJoinRequest: // 9021 + if id == "" { + return true, "missing group id (h tag)" + } + gs.mu.Lock() + g := gs.ensureGroup(id, evt.PubKey, true) + // Open groups auto-join. Closed groups accept the request event but do + // NOT grant membership — an admin must add the user (kind 9000). + if g.open { + g.members[evt.PubKey] = true + } + gs.mu.Unlock() + return false, "" + + case nostr.KindSimpleGroupLeaveRequest: // 9022 + gs.mu.Lock() + if g := gs.groups[id]; g != nil { + delete(g.members, evt.PubKey) + } + gs.mu.Unlock() + return false, "" + + case nostr.KindSimpleGroupPutUser: // 9000 (admin adds user) + if !gs.IsAdmin(id, evt.PubKey) { + return true, "only admins can add users" + } + for _, pk := range taggedPubkeys(evt) { + gs.setMember(id, pk, true) + } + return false, "" + + case nostr.KindSimpleGroupRemoveUser: // 9001 (admin removes user) + if !gs.IsAdmin(id, evt.PubKey) { + return true, "only admins can remove users" + } + for _, pk := range taggedPubkeys(evt) { + gs.setMember(id, pk, false) + } + return false, "" + + case nostr.KindSimpleGroupEditMetadata: // 9002 (admin sets name/about) + if id == "" { + return true, "missing group id (h tag)" + } + if !gs.IsAdmin(id, evt.PubKey) { + return true, "only admins can edit group metadata" + } + gs.mu.Lock() + if g := gs.groups[id]; g != nil { + if name := firstTagValue(evt, "name"); name != "" { + g.name = name + } + // `about` is the room topic; allow clearing it with an empty value. + if t := evt.Tags.Find("about"); len(t) >= 2 { + g.about = t[1] + } + } + gs.mu.Unlock() + return false, "" + + case nostr.KindSimpleGroupChatMessage, nostr.KindSimpleGroupThreadedReply: // 9 / 10 + if id == "" { + return true, "chat message missing group id (h tag)" + } + gs.mu.RLock() + g := gs.groups[id] + gs.mu.RUnlock() + if g == nil { + // A channel we don't know locally may still be one we federate — accept + // so peer chat can land (the group is created lazily on federated write). + if gs.IsFederated != nil && gs.IsFederated(id) { + return false, "" + } + return true, "unknown group; send a join request (kind 9021) first" + } + // Federated channels behave as open across all participating relays: the + // author is a member on the peer that authorized the post, so don't gate + // federated chat on local membership. + if !g.open && !gs.IsMember(id, evt.PubKey) { + if gs.IsFederated != nil && gs.IsFederated(id) { + return false, "" + } + return true, "you are not a member of this group" + } + return false, "" + } + + // Non-group events are handled by the relay's normal rules. + return false, "" +} + +// IsMember reports whether pk is a member of the group. +func (gs *State) IsMember(id string, pk nostr.PubKey) bool { + gs.mu.RLock() + defer gs.mu.RUnlock() + g := gs.groups[id] + return g != nil && g.members[pk] +} + +// IsAdmin reports whether pk is an admin of the group. +func (gs *State) IsAdmin(id string, pk nostr.PubKey) bool { + gs.mu.RLock() + defer gs.mu.RUnlock() + g := gs.groups[id] + return g != nil && g.admins[pk] +} + +func (gs *State) setMember(id string, pk nostr.PubKey, member bool) { + gs.mu.Lock() + defer gs.mu.Unlock() + g := gs.groups[id] + if g == nil { + return + } + if member { + g.members[pk] = true + } else { + delete(g.members, pk) + } +} + +// taggedPubkeys extracts valid `p` tag pubkeys from an event. +func taggedPubkeys(evt nostr.Event) []nostr.PubKey { + var out []nostr.PubKey + for t := range evt.Tags.FindAll("p") { + if len(t) >= 2 { + if pk, err := nostr.PubKeyFromHex(t[1]); err == nil { + out = append(out, pk) + } + } + } + return out +} + +// MetadataEvent builds and signs the relay's kind-39000 metadata event for a +// group, so clients can discover the group's name/about via NIP-29. +func (gs *State) MetadataEvent(ctx context.Context, id string) (nostr.Event, bool) { + gs.mu.RLock() + g := gs.groups[id] + gs.mu.RUnlock() + if g == nil { + return nostr.Event{}, false + } + + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupMetadata, + CreatedAt: nostr.Now(), + Content: "", + Tags: nostr.Tags{ + {"d", id}, + {"name", g.name}, + {"about", g.about}, + }, + } + if g.open { + evt.Tags = append(evt.Tags, nostr.Tag{"open"}) + } else { + evt.Tags = append(evt.Tags, nostr.Tag{"closed"}) + } + if err := evt.Sign(gs.sk); err != nil { + return nostr.Event{}, false + } + return evt, true +} + +// AdminsEvent builds and signs the relay's kind-39001 admins list for a group. +// Each admin is a `p` tag with an "admin" role, per NIP-29. +func (gs *State) AdminsEvent(id string) (nostr.Event, bool) { + gs.mu.RLock() + g := gs.groups[id] + var admins []nostr.PubKey + if g != nil { + for pk := range g.admins { + admins = append(admins, pk) + } + } + gs.mu.RUnlock() + if g == nil { + return nostr.Event{}, false + } + + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupAdmins, + CreatedAt: nostr.Now(), + Tags: nostr.Tags{{"d", id}}, + } + for _, pk := range admins { + evt.Tags = append(evt.Tags, nostr.Tag{"p", pk.Hex(), "admin"}) + } + if err := evt.Sign(gs.sk); err != nil { + return nostr.Event{}, false + } + return evt, true +} + +// MembersEvent builds and signs the relay's kind-39002 members list for a group. +func (gs *State) MembersEvent(id string) (nostr.Event, bool) { + gs.mu.RLock() + g := gs.groups[id] + var members []nostr.PubKey + if g != nil { + for pk := range g.members { + members = append(members, pk) + } + } + gs.mu.RUnlock() + if g == nil { + return nostr.Event{}, false + } + + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupMembers, + CreatedAt: nostr.Now(), + Tags: nostr.Tags{{"d", id}}, + } + for _, pk := range members { + evt.Tags = append(evt.Tags, nostr.Tag{"p", pk.Hex()}) + } + if err := evt.Sign(gs.sk); err != nil { + return nostr.Event{}, false + } + return evt, true +} + +// EnsureFederatedGroup makes sure a group exists locally so that chat federated +// from a peer relay has somewhere to land. Federated chat is chat-only: the +// author is a member on the PEER (which already authorized the post), not +// necessarily here, so we must not gate it on local membership. We create the +// channel open on first sight of federated chat if it doesn't exist yet; we +// never alter membership or admin state from federated events. +func (gs *State) EnsureFederatedGroup(id string) { + if id == "" { + return + } + gs.mu.Lock() + if _, ok := gs.groups[id]; !ok { + gs.groups[id] = &group{ + id: id, + name: id, + members: map[nostr.PubKey]bool{}, + admins: map[nostr.PubKey]bool{}, + open: true, + } + } + gs.mu.Unlock() +} + +// Summary is a short human-readable description for logging/inspection. +func (gs *State) Summary() string { + gs.mu.RLock() + defer gs.mu.RUnlock() + return fmt.Sprintf("%d group(s) tracked", len(gs.groups)) +} diff --git a/internal/group/group_test.go b/internal/group/group_test.go new file mode 100644 index 0000000..4b7a331 --- /dev/null +++ b/internal/group/group_test.go @@ -0,0 +1,291 @@ +package group + +import ( + "context" + "testing" + + "fiatjaf.com/nostr" + "fiatjaf.com/nostr/eventstore/slicestore" +) + +func mustSign(t *testing.T, sk nostr.SecretKey, evt *nostr.Event) { + t.Helper() + if err := evt.Sign(sk); err != nil { + t.Fatalf("sign: %v", err) + } +} + +func hpk(sk nostr.SecretKey) nostr.PubKey { return nostr.GetPublicKey(sk) } + +// evalOK applies an event and fails if the relay rejects it. +func evalOK(t *testing.T, gs *State, evt nostr.Event) { + t.Helper() + if reject, msg := gs.Evaluate(evt); reject { + t.Fatalf("kind %d unexpectedly rejected: %s", evt.Kind, msg) + } +} + +// evalRejected applies an event and fails if it is allowed. +func evalRejected(t *testing.T, gs *State, evt nostr.Event) string { + t.Helper() + reject, msg := gs.Evaluate(evt) + if !reject { + t.Fatalf("kind %d unexpectedly allowed", evt.Kind) + } + return msg +} + +func TestChatRejectedWithoutGroup(t *testing.T) { + gs := NewState(nostr.Generate()) + // A chat message to a group that was never created/joined is rejected. + evt := nostr.Event{ + Kind: nostr.KindSimpleGroupChatMessage, + Tags: nostr.Tags{{"h", "nonexistent"}}, + } + reject, msg := gs.Evaluate(evt) + if !reject { + t.Fatalf("expected rejection for chat to unknown group, got allow (%q)", msg) + } +} + +func TestChatMissingHTagRejected(t *testing.T) { + gs := NewState(nostr.Generate()) + evt := nostr.Event{Kind: nostr.KindSimpleGroupChatMessage} + reject, _ := gs.Evaluate(evt) + if !reject { + t.Fatal("expected rejection for chat with no h tag") + } +} + +func TestJoinThenChatAllowed(t *testing.T) { + gs := NewState(nostr.Generate()) + alice := nostr.Generate() + alicePK := nostr.GetPublicKey(alice) + + join := nostr.Event{PubKey: alicePK, Kind: nostr.KindSimpleGroupJoinRequest, Tags: nostr.Tags{{"h", "g1"}}} + if reject, msg := gs.Evaluate(join); reject { + t.Fatalf("join rejected: %s", msg) + } + chat := nostr.Event{PubKey: alicePK, Kind: nostr.KindSimpleGroupChatMessage, Tags: nostr.Tags{{"h", "g1"}}} + if reject, msg := gs.Evaluate(chat); reject { + t.Fatalf("chat after join rejected: %s", msg) + } + if !gs.IsMember("g1", alicePK) { + t.Fatal("alice should be a member after joining") + } +} + +// The creator of a group is its administrator; a plain joiner is not, and only +// the admin may add or remove users. +func TestAdminModerationInClosedGroup(t *testing.T) { + gs := NewState(nostr.Generate()) + + admin := nostr.Generate() // creates the group → admin + member + alice := nostr.Generate() // will be admitted by the admin + bob := nostr.Generate() // will try to moderate without permission + mallory := nostr.Generate() + const g = "staff" + + adminPK, alicePK, bobPK, malloryPK := hpk(admin), hpk(alice), hpk(bob), hpk(mallory) + + // Admin creates a CLOSED group. + evalOK(t, gs, nostr.Event{ + PubKey: adminPK, Kind: nostr.KindSimpleGroupCreateGroup, + Tags: nostr.Tags{{"h", g}, {"closed"}}, + }) + if !gs.IsAdmin(g, adminPK) { + t.Fatal("creator should be an admin") + } + if !gs.IsMember(g, adminPK) { + t.Fatal("creator should be a member") + } + + // A join request to a closed group is accepted as an event but does NOT + // grant membership — it awaits admin approval. + evalOK(t, gs, nostr.Event{ + PubKey: alicePK, Kind: nostr.KindSimpleGroupJoinRequest, Tags: nostr.Tags{{"h", g}}, + }) + if gs.IsMember(g, alicePK) { + t.Fatal("closed group: join request alone must not grant membership") + } + + // A non-member/non-admin cannot post to the closed group. + evalRejected(t, gs, nostr.Event{ + PubKey: alicePK, Kind: nostr.KindSimpleGroupChatMessage, Tags: nostr.Tags{{"h", g}}, + }) + + // A non-admin (bob) cannot add users. + msg := evalRejected(t, gs, nostr.Event{ + PubKey: bobPK, Kind: nostr.KindSimpleGroupPutUser, + Tags: nostr.Tags{{"h", g}, {"p", alicePK.Hex()}}, + }) + if msg == "" { + t.Fatal("expected a rejection reason for non-admin add") + } + + // The admin admits alice (kind 9000). + evalOK(t, gs, nostr.Event{ + PubKey: adminPK, Kind: nostr.KindSimpleGroupPutUser, + Tags: nostr.Tags{{"h", g}, {"p", alicePK.Hex()}}, + }) + if !gs.IsMember(g, alicePK) { + t.Fatal("admin add (9000) should grant alice membership") + } + + // Now alice, a member, can post. + evalOK(t, gs, nostr.Event{ + PubKey: alicePK, Kind: nostr.KindSimpleGroupChatMessage, Tags: nostr.Tags{{"h", g}}, + }) + + // The admin can add multiple users in one event. + evalOK(t, gs, nostr.Event{ + PubKey: adminPK, Kind: nostr.KindSimpleGroupPutUser, + Tags: nostr.Tags{{"h", g}, {"p", bobPK.Hex()}, {"p", malloryPK.Hex()}}, + }) + if !gs.IsMember(g, bobPK) || !gs.IsMember(g, malloryPK) { + t.Fatal("admin should add multiple tagged users") + } + + // The admin removes mallory (kind 9001); she can no longer post. + evalOK(t, gs, nostr.Event{ + PubKey: adminPK, Kind: nostr.KindSimpleGroupRemoveUser, + Tags: nostr.Tags{{"h", g}, {"p", malloryPK.Hex()}}, + }) + if gs.IsMember(g, malloryPK) { + t.Fatal("admin remove (9001) should revoke mallory's membership") + } + evalRejected(t, gs, nostr.Event{ + PubKey: malloryPK, Kind: nostr.KindSimpleGroupChatMessage, Tags: nostr.Tags{{"h", g}}, + }) + + // A non-admin (bob) cannot remove users either. + evalRejected(t, gs, nostr.Event{ + PubKey: bobPK, Kind: nostr.KindSimpleGroupRemoveUser, + Tags: nostr.Tags{{"h", g}, {"p", alicePK.Hex()}}, + }) + if !gs.IsMember(g, alicePK) { + t.Fatal("a rejected removal must not change membership") + } +} + +// Open groups let anyone post without admin approval (public channels), which +// contrasts with the closed-group moderation above. +func TestOpenGroupAllowsNonMemberPosts(t *testing.T) { + gs := NewState(nostr.Generate()) + admin := nostr.Generate() + stranger := nostr.Generate() + const g = "lobby" + + evalOK(t, gs, nostr.Event{ + PubKey: hpk(admin), Kind: nostr.KindSimpleGroupCreateGroup, Tags: nostr.Tags{{"h", g}}, + }) + // A stranger who never joined can still post to an open group. + evalOK(t, gs, nostr.Event{ + PubKey: hpk(stranger), Kind: nostr.KindSimpleGroupChatMessage, Tags: nostr.Tags{{"h", g}}, + }) +} + +// Group membership and metadata must survive a relay restart: state is rebuilt +// from the persisted event store rather than lost. This reproduces the bug +// where an in-memory-only State silently dropped all groups on restart. +func TestStateSurvivesRebuild(t *testing.T) { + store := &slicestore.SliceStore{} + if err := store.Init(); err != nil { + t.Fatalf("store init: %v", err) + } + + admin := nostr.Generate() + alice := nostr.Generate() + const g = "persisted" + + // Simulate the live path: Evaluate accepts the event, then it is stored + // (khatru stores accepted events; we do it explicitly here). + apply := func(gs *State, evt nostr.Event) { + if reject, msg := gs.Evaluate(evt); reject { + t.Fatalf("kind %d rejected: %s", evt.Kind, msg) + } + if err := store.SaveEvent(evt); err != nil { + t.Fatalf("save: %v", err) + } + } + + signed := func(sk nostr.SecretKey, kind nostr.Kind, tags nostr.Tags) nostr.Event { + e := nostr.Event{Kind: kind, CreatedAt: nostr.Now(), Tags: tags} + mustSign(t, sk, &e) + return e + } + + // Build state in the "first boot": admin creates a closed group and admits alice. + first := NewState(nostr.Generate()) + apply(first, signed(admin, nostr.KindSimpleGroupCreateGroup, nostr.Tags{{"h", g}, {"closed"}})) + apply(first, signed(admin, nostr.KindSimpleGroupPutUser, nostr.Tags{{"h", g}, {"p", hpk(alice).Hex()}})) + + if !first.IsMember(g, hpk(alice)) || !first.IsAdmin(g, hpk(admin)) { + t.Fatal("precondition: alice member + admin admin should hold before restart") + } + + // "Restart": a fresh State that knows nothing until it replays the store. + second := NewState(nostr.Generate()) + if second.IsMember(g, hpk(alice)) { + t.Fatal("fresh state should be empty before rebuild") + } + replayed := second.Rebuild(store) + if replayed != 2 { + t.Fatalf("expected 2 management events replayed, got %d", replayed) + } + + // Membership, admin rights, and closed-ness must be restored. + if !second.IsMember(g, hpk(admin)) { + t.Fatal("admin membership lost across rebuild") + } + if !second.IsMember(g, hpk(alice)) { + t.Fatal("alice membership lost across rebuild") + } + if !second.IsAdmin(g, hpk(admin)) { + t.Fatal("admin rights lost across rebuild") + } + // Closed-group enforcement must still apply after rebuild: a non-member is rejected. + stranger := nostr.Generate() + if reject, _ := second.Evaluate(signed(stranger, nostr.KindSimpleGroupChatMessage, nostr.Tags{{"h", g}})); !reject { + t.Fatal("closed-group enforcement lost across rebuild: stranger post should be rejected") + } + // A restored member can still post. + if reject, msg := second.Evaluate(signed(alice, nostr.KindSimpleGroupChatMessage, nostr.Tags{{"h", g}})); reject { + t.Fatalf("restored member alice should be able to post: %s", msg) + } +} + +// Only admins can edit group metadata (topic), and the edit is applied to state. +func TestEditMetadataAdminOnly(t *testing.T) { + gs := NewState(nostr.Generate()) + admin := nostr.Generate() + stranger := nostr.Generate() + const g = "topictest" + + // Admin creates the group. + evalOK(t, gs, nostr.Event{ + PubKey: hpk(admin), Kind: nostr.KindSimpleGroupCreateGroup, Tags: nostr.Tags{{"h", g}}, + }) + + // A non-admin cannot set the topic. + evalRejected(t, gs, nostr.Event{ + PubKey: hpk(stranger), Kind: nostr.KindSimpleGroupEditMetadata, + Tags: nostr.Tags{{"h", g}, {"about", "hostile takeover"}}, + }) + + // The admin sets name + topic, and state reflects it. + evalOK(t, gs, nostr.Event{ + PubKey: hpk(admin), Kind: nostr.KindSimpleGroupEditMetadata, + Tags: nostr.Tags{{"h", g}, {"name", "Topic Test"}, {"about", "the topic"}}, + }) + meta, ok := gs.MetadataEvent(context.Background(), g) + if !ok { + t.Fatal("expected metadata event") + } + if firstTagValue(meta, "name") != "Topic Test" { + t.Fatalf("name not applied: %q", firstTagValue(meta, "name")) + } + if firstTagValue(meta, "about") != "the topic" { + t.Fatalf("about (topic) not applied: %q", firstTagValue(meta, "about")) + } +} diff --git a/internal/relayinfo/relayinfo.go b/internal/relayinfo/relayinfo.go new file mode 100644 index 0000000..a0fee97 --- /dev/null +++ b/internal/relayinfo/relayinfo.go @@ -0,0 +1,42 @@ +// Package relayinfo implements the Nosterm capability handshake. +// +// A standard NIP-11 relay information document describes supported NIPs. A +// "nostermd" additionally advertises a `features` array so that a Nosterm +// client can enable/disable UI capabilities and gracefully fall back when +// connected to a plain Nostr relay that omits it. +// +// The upstream nip11.RelayInformationDocument has no `features` field and its +// custom JSON marshaller drops unknown keys, so the relay serves NIP-11 itself +// and merges these fields into the JSON object. +package relayinfo + +const ( + // SoftwareName is advertised as the NIP-11 `software` value. + SoftwareName = "nostermd" + // SoftwareVersion is advertised as the NIP-11 `version` value. + SoftwareVersion = "0.1.0" +) + +// baseFeatures advertised by every relay build. Presence/typing/moderation/etc. +// are declared here as they are implemented; the baseline relay supports managed +// NIP-29 channels only. Runtime-dependent features (e.g. federation) are added +// in Extras based on configuration. +var baseFeatures = []string{ + "channels", // NIP-29 managed group channels +} + +// Extras returns the extra top-level keys merged into the NIP-11 document JSON. +// Keeping this in one place makes the handshake easy to extend. `federated` +// reflects whether this relay has any federation peer configured, so a client +// can surface that a channel here may be mirrored elsewhere. +func Extras(federated bool) map[string]any { + features := append([]string{}, baseFeatures...) + if federated { + features = append(features, "federation") // chat-only relay↔relay mirroring + } + return map[string]any{ + "software": SoftwareName, + "version": SoftwareVersion, + "features": features, + } +} diff --git a/internal/retention/retention.go b/internal/retention/retention.go new file mode 100644 index 0000000..acaf5fc --- /dev/null +++ b/internal/retention/retention.go @@ -0,0 +1,96 @@ +// Package retention prunes old CHAT messages (kinds 9/10) so the event store +// does not grow without bound. Group-management and addressable metadata events +// (9007/9000/9001/9021/9022, 39000-39002) are NEVER pruned — they define +// membership and must survive to be replayed on boot. +package retention + +import ( + "log" + "slices" + "time" + + "fiatjaf.com/nostr" + "fiatjaf.com/nostr/eventstore" +) + +// Config controls what the retention sweeper deletes. +type Config struct { + // MaxAge: chat messages older than this are deleted. Zero disables age pruning. + MaxAge time.Duration + // MaxMessages: keep at most this many chat messages (newest kept). Zero + // disables count pruning. + MaxMessages int + // Interval between sweeps. + Interval time.Duration +} + +// chatKinds are the only kinds eligible for retention pruning. +var chatKinds = []nostr.Kind{ + nostr.KindSimpleGroupChatMessage, + nostr.KindSimpleGroupThreadedReply, +} + +// Sweep runs a single retention pass and returns the number of events deleted. +// `now` is injected so the logic is deterministically testable. +func Sweep(store eventstore.Store, cfg Config, now time.Time) int { + // Collect all chat events (newest-first from the store). + var chats []nostr.Event + for evt := range store.QueryEvents(nostr.Filter{Kinds: chatKinds}, 1_000_000) { + chats = append(chats, evt) + } + // Sort newest-first so index >= MaxMessages are the surplus oldest ones. + slices.SortFunc(chats, nostr.CompareEventReverse) + + cutoff := nostr.Timestamp(0) + if cfg.MaxAge > 0 { + cutoff = nostr.Timestamp(now.Add(-cfg.MaxAge).Unix()) + } + + deleted := 0 + for i, evt := range chats { + tooOld := cfg.MaxAge > 0 && evt.CreatedAt < cutoff + surplus := cfg.MaxMessages > 0 && i >= cfg.MaxMessages + if tooOld || surplus { + if err := store.DeleteEvent(evt.ID); err == nil { + deleted++ + } + } + } + return deleted +} + +// Start launches a background sweeper. It is a no-op (returns a no-op stop) when +// neither limit is configured. The returned stop function halts it. +func Start(store eventstore.Store, cfg Config) (stop func()) { + if cfg.MaxAge == 0 && cfg.MaxMessages == 0 { + log.Println("retention disabled (no RELAY_RETENTION_* limits set)") + return func() {} + } + if cfg.Interval == 0 { + cfg.Interval = time.Hour + } + + done := make(chan struct{}) + go func() { + // Sweep once at startup, then on the interval. + if n := Sweep(store, cfg, time.Now()); n > 0 { + log.Printf("retention: pruned %d old chat message(s)", n) + } + ticker := time.NewTicker(cfg.Interval) + defer ticker.Stop() + for { + select { + case <-done: + return + case <-ticker.C: + if n := Sweep(store, cfg, time.Now()); n > 0 { + log.Printf("retention: pruned %d old chat message(s)", n) + } + } + } + }() + + log.Printf("retention enabled (maxAge=%s maxMessages=%d interval=%s)", + cfg.MaxAge, cfg.MaxMessages, cfg.Interval) + return func() { close(done) } +} diff --git a/internal/retention/retention_test.go b/internal/retention/retention_test.go new file mode 100644 index 0000000..4d947e9 --- /dev/null +++ b/internal/retention/retention_test.go @@ -0,0 +1,116 @@ +package retention + +import ( + "testing" + "time" + + "fiatjaf.com/nostr" + "fiatjaf.com/nostr/eventstore/slicestore" +) + +func mkStore(t *testing.T) *slicestore.SliceStore { + t.Helper() + s := &slicestore.SliceStore{} + if err := s.Init(); err != nil { + t.Fatalf("init: %v", err) + } + return s +} + +func save(t *testing.T, s *slicestore.SliceStore, sk nostr.SecretKey, kind nostr.Kind, createdAt nostr.Timestamp) nostr.ID { + t.Helper() + e := nostr.Event{Kind: kind, CreatedAt: createdAt, Tags: nostr.Tags{{"h", "g"}}} + if kind == nostr.KindSimpleGroupChatMessage { + e.Content = "msg" + } + if err := e.Sign(sk); err != nil { + t.Fatalf("sign: %v", err) + } + if err := s.SaveEvent(e); err != nil { + t.Fatalf("save: %v", err) + } + return e.ID +} + +func count(s *slicestore.SliceStore, kinds ...nostr.Kind) int { + n := 0 + for range s.QueryEvents(nostr.Filter{Kinds: kinds}, 1_000_000) { + n++ + } + return n +} + +func TestRetentionPrunesOldChatByAge(t *testing.T) { + s := mkStore(t) + sk := nostr.Generate() + now := time.Unix(1_700_000_000, 0) + old := nostr.Timestamp(now.Add(-48 * time.Hour).Unix()) + fresh := nostr.Timestamp(now.Add(-1 * time.Hour).Unix()) + + save(t, s, sk, nostr.KindSimpleGroupChatMessage, old) + save(t, s, sk, nostr.KindSimpleGroupChatMessage, fresh) + + deleted := Sweep(s, Config{MaxAge: 24 * time.Hour}, now) + if deleted != 1 { + t.Fatalf("expected 1 pruned, got %d", deleted) + } + if c := count(s, nostr.KindSimpleGroupChatMessage); c != 1 { + t.Fatalf("expected 1 chat remaining, got %d", c) + } +} + +func TestRetentionPrunesSurplusByCount(t *testing.T) { + s := mkStore(t) + sk := nostr.Generate() + now := time.Unix(1_700_000_000, 0) + for i := 0; i < 5; i++ { + save(t, s, sk, nostr.KindSimpleGroupChatMessage, nostr.Timestamp(1_700_000_000+int64(i))) + } + deleted := Sweep(s, Config{MaxMessages: 2}, now) + if deleted != 3 { + t.Fatalf("expected 3 pruned, got %d", deleted) + } + if c := count(s, nostr.KindSimpleGroupChatMessage); c != 2 { + t.Fatalf("expected 2 chats remaining, got %d", c) + } +} + +func TestRetentionNeverPrunesManagementEvents(t *testing.T) { + s := mkStore(t) + sk := nostr.Generate() + now := time.Unix(1_700_000_000, 0) + ancient := nostr.Timestamp(now.Add(-1000 * time.Hour).Unix()) + + // Very old management + metadata events must be preserved. + save(t, s, sk, nostr.KindSimpleGroupCreateGroup, ancient) + save(t, s, sk, nostr.KindSimpleGroupPutUser, ancient) + save(t, s, sk, nostr.KindSimpleGroupJoinRequest, ancient) + save(t, s, sk, nostr.KindSimpleGroupMetadata, ancient) + save(t, s, sk, nostr.KindSimpleGroupAdmins, ancient) + save(t, s, sk, nostr.KindSimpleGroupMembers, ancient) + // An old chat that SHOULD be pruned. + save(t, s, sk, nostr.KindSimpleGroupChatMessage, ancient) + + deleted := Sweep(s, Config{MaxAge: time.Hour, MaxMessages: 1}, now) + if deleted != 1 { + t.Fatalf("expected only the 1 chat pruned, got %d", deleted) + } + // All 6 management/metadata events survive. + mgmt := count(s, + nostr.KindSimpleGroupCreateGroup, nostr.KindSimpleGroupPutUser, + nostr.KindSimpleGroupJoinRequest, nostr.KindSimpleGroupMetadata, + nostr.KindSimpleGroupAdmins, nostr.KindSimpleGroupMembers) + if mgmt != 6 { + t.Fatalf("management events must never be pruned; have %d/6", mgmt) + } +} + +func TestRetentionDisabledIsNoop(t *testing.T) { + s := mkStore(t) + sk := nostr.Generate() + now := time.Unix(1_700_000_000, 0) + save(t, s, sk, nostr.KindSimpleGroupChatMessage, nostr.Timestamp(1)) + if d := Sweep(s, Config{}, now); d != 0 { + t.Fatalf("no limits should prune nothing, got %d", d) + } +}