diff --git a/.changeset/haproxy-blue-green.md b/.changeset/haproxy-blue-green.md new file mode 100644 index 00000000..dd428e7b --- /dev/null +++ b/.changeset/haproxy-blue-green.md @@ -0,0 +1,7 @@ +--- +"nostream": minor +--- + +feat(deploy): add HAProxy blue/green compose stack + +Two relays behind HAProxy with `/readyz` health checks and `option redispatch`, plus a rolling recreate script that replaces one relay at a time for zero-downtime image updates. Optional Redis stream fan-out (`RELAY_BROADCAST_FANOUT`) lets both relays share live WebSocket broadcasts while workers keep using cluster `process.send`. diff --git a/CONFIGURATION.md b/CONFIGURATION.md index ec492738..fafec7f5 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -54,6 +54,10 @@ The following environment variables can be set: | REDIS_PORT | Redis Port | 6379 | | REDIS_USER | Redis User | default | | REDIS_PASSWORD | Redis Password | nostr_ts_relay | +| RELAY_BROADCAST_FANOUT | Publish accepted events to a Redis stream so multiple relay containers share live fan-out (`true`/`false`) | `false` | +| RELAY_BROADCAST_STREAM_KEY | Redis stream key used when `RELAY_BROADCAST_FANOUT` is enabled | `nostream:relay:broadcast` | +| RELAY_BROADCAST_STREAM_MAXLEN | Approximate max entries retained on the broadcast Redis stream (`XADD` trim) | `50000` | +| RELAY_INSTANCE_ID | Optional stable id for this relay instance (defaults to hostname and pid) | | | PROMETHEUS_URL | Prometheus base URL for admin metrics queries | http://127.0.0.1:9090 | | PROMETHEUS_QUERY_TIMEOUT_MS | Timeout for each Prometheus admin metrics query (ms) | 5000 | | ADMIN_METRICS_SSE_INTERVAL_MS | Interval between `/admin/metrics` SSE snapshots (ms) | 5000 | diff --git a/deploy/README.md b/deploy/README.md index d6b4428d..77e369c8 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -156,3 +156,64 @@ the new checkout (or copy the updated files). Automated sync is planned separate ``` Existing `.env` and `.nostr/settings.yaml` are preserved. + +## Zero-downtime updates (HAProxy blue/green) + +`deploy/docker-compose.haproxy.yml` replaces the single-relay stack with two +relays (`nostream-blue`, `nostream-green`) behind HAProxy on `127.0.0.1:8008`. +Postgres, Redis, and migrations are unchanged. + +The HAProxy compose file **always sets `RELAY_BROADCAST_FANOUT=true` on relay +services** (even if bootstrap `.env` leaves it `false` for the single-relay +stack). Both relays publish accepted events to a shared Redis stream; each +cluster primary subscribes and fans out to its workers so live WebSocket clients +stay in sync when HAProxy balances across blue and green. + +Do **not** run the single-relay `docker-compose.yml` stack and the HAProxy stack +at the same time: both bind `127.0.0.1:8008`. Stop the old stack before starting +blue/green: + +```bash +cd /opt/nostream +docker compose down # single-relay stack, if it was running +``` + +Install alongside `.env` and `postgresql.conf`, then start: + +```bash +cp deploy/docker-compose.haproxy.yml deploy/rolling-relay-recreate.sh /opt/nostream/ +cp -r deploy/haproxy /opt/nostream/ +chmod +x /opt/nostream/rolling-relay-recreate.sh +cd /opt/nostream +docker compose -f docker-compose.haproxy.yml up -d +curl -s http://127.0.0.1:8008/readyz +``` + +HAProxy sets `X-Forwarded-For`. In `.nostr/settings.yaml` (or your settings +overrides), enable forwarded client IPs when using this stack: + +```yaml +network: + remoteIpHeader: x-forwarded-for + trustedProxies: + - "127.0.0.1" + - "::ffff:127.0.0.1" + - "::1" + # HAProxy container on the compose network (get after first up): + # docker inspect -f '{{range .NetworkSettings.Networks}}{{.IPAddress}}{{end}}' nostream-haproxy +``` + +To update, load the new image, run migrations, then replace relays one at a time: + +```bash +cd /opt/nostream +docker compose -f docker-compose.haproxy.yml run --rm nostream-migrate +./rolling-relay-recreate.sh +``` + +The script requires the peer relay to be running and `/readyz` healthy before it +stops either backend. It waits for each replacement to become healthy before +moving to the second relay. HAProxy health-checks `/readyz` every 2s and retries +failed requests on the other backend (`option redispatch`). Set +`STOP_GRACE_PERIOD` (default `45s`) above `WS_DRAIN_TIMEOUT_MS` (default 30s) so +WebSocket drain finishes before Docker sends SIGKILL. diff --git a/deploy/docker-compose.haproxy.yml b/deploy/docker-compose.haproxy.yml new file mode 100644 index 00000000..1d92e1fb --- /dev/null +++ b/deploy/docker-compose.haproxy.yml @@ -0,0 +1,129 @@ +# Blue/green alternative to docker-compose.prod.yml: two relays behind HAProxy. +# Keep the shared services below in sync with docker-compose.prod.yml. +# +# docker compose -f docker-compose.haproxy.yml up -d + +x-nostream-relay: &nostream-relay + image: ghcr.io/cameri/nostream:main + pull_policy: never + env_file: .env + environment: + RELAY_PORT: 8008 + NOSTR_CONFIG_DIR: /home/node/.nostr + DB_HOST: nostream-db + DB_PORT: 5432 + DB_USER: ${DB_USER} + DB_PASSWORD: ${DB_PASSWORD} + DB_NAME: ${DB_NAME} + DB_MIN_POOL_SIZE: ${DB_MIN_POOL_SIZE:-16} + DB_MAX_POOL_SIZE: ${DB_MAX_POOL_SIZE:-64} + DB_ACQUIRE_CONNECTION_TIMEOUT: ${DB_ACQUIRE_CONNECTION_TIMEOUT:-60000} + REDIS_HOST: nostream-cache + REDIS_PORT: 6379 + REDIS_USER: default + REDIS_PASSWORD: ${REDIS_PASSWORD} + READ_REPLICA_ENABLED: 'false' + # Required for blue/green; do not inherit RELAY_BROADCAST_FANOUT=false from bootstrap .env + RELAY_BROADCAST_FANOUT: 'true' + RELAY_BROADCAST_STREAM_KEY: ${RELAY_BROADCAST_STREAM_KEY:-nostream:relay:broadcast} + WORKER_COUNT: ${WORKER_COUNT:-2} + WS_DRAIN_TIMEOUT_MS: ${WS_DRAIN_TIMEOUT_MS:-30000} + user: node:node + volumes: + - ${PWD}/.nostr:/home/node/.nostr + depends_on: + nostream-cache: + condition: service_healthy + nostream-db: + condition: service_healthy + nostream-migrate: + condition: service_completed_successfully + restart: on-failure + stop_grace_period: ${STOP_GRACE_PERIOD:-45s} + # Gates `up --wait` during rolling recreate, so the next relay is only + # replaced once this one serves traffic. + healthcheck: + test: + [ + 'CMD-SHELL', + "node -e \"fetch('http://127.0.0.1:8008/readyz').then(r=>process.exit(r.ok?0:1)).catch(()=>process.exit(1))\"", + ] + interval: 5s + timeout: 5s + retries: 6 + start_period: 60s + +services: + haproxy: + image: haproxy:3.0-alpine + container_name: nostream-haproxy + volumes: + - ./haproxy/haproxy.cfg:/usr/local/etc/haproxy/haproxy.cfg:ro + ports: + - 127.0.0.1:8008:8008 + depends_on: + - nostream-blue + - nostream-green + restart: on-failure + + nostream-blue: + <<: *nostream-relay + container_name: nostream-blue + + nostream-green: + <<: *nostream-relay + container_name: nostream-green + + nostream-db: + image: postgres:15 + container_name: nostream-db + environment: + POSTGRES_DB: ${DB_NAME} + POSTGRES_USER: ${DB_USER} + POSTGRES_PASSWORD: ${DB_PASSWORD} + volumes: + - ${PWD}/.nostr/data:/var/lib/postgresql/data + - ${PWD}/.nostr/db-logs:/var/log/postgresql + - ${PWD}/postgresql.conf:/postgresql.conf + command: postgres -c 'config_file=/postgresql.conf' + restart: always + healthcheck: + test: ['CMD-SHELL', 'pg_isready -U ${DB_USER}'] + interval: 5s + timeout: 5s + retries: 5 + start_period: 360s + + nostream-cache: + image: redis:7.0.5-alpine3.16 + container_name: nostream-cache + environment: + REDIS_PASSWORD: ${REDIS_PASSWORD} + volumes: + - cache:/data + command: sh -c 'redis-server --loglevel warning --requirepass "$$REDIS_PASSWORD"' + restart: always + healthcheck: + test: ['CMD-SHELL', 'redis-cli -a "$$REDIS_PASSWORD" ping | grep PONG'] + interval: 2s + timeout: 5s + retries: 10 + + nostream-migrate: + image: ghcr.io/cameri/nostream:main + pull_policy: never + container_name: nostream-migrate + user: node:node + command: ['node_modules/.bin/knex', 'migrate:latest'] + environment: + DB_HOST: nostream-db + DB_PORT: 5432 + DB_USER: ${DB_USER} + DB_PASSWORD: ${DB_PASSWORD} + DB_NAME: ${DB_NAME} + depends_on: + nostream-db: + condition: service_healthy + +volumes: + cache: diff --git a/deploy/env.example b/deploy/env.example index f0bc3826..1c84bac7 100644 --- a/deploy/env.example +++ b/deploy/env.example @@ -21,3 +21,7 @@ DB_MIN_POOL_SIZE=16 DB_MAX_POOL_SIZE=64 DB_ACQUIRE_CONNECTION_TIMEOUT=60000 WORKER_COUNT=2 + +# Single-relay stack only. HAProxy compose (docker-compose.haproxy.yml) always enables fan-out on relay services. +RELAY_BROADCAST_FANOUT=false +RELAY_BROADCAST_STREAM_KEY=nostream:relay:broadcast diff --git a/deploy/haproxy/haproxy.cfg b/deploy/haproxy/haproxy.cfg new file mode 100644 index 00000000..7a818a60 --- /dev/null +++ b/deploy/haproxy/haproxy.cfg @@ -0,0 +1,41 @@ +global + log stdout format raw local0 info + maxconn 4096 + +defaults + mode http + log global + option httplog + timeout connect 5s + timeout client 30s + timeout server 30s + timeout tunnel 1h + timeout check 5s + retries 3 + option redispatch + retry-on conn-failure empty-response response-timeout 502 503 504 + http-restrict-req-retry-on all + +# Docker's embedded DNS, so recreated relay containers are picked up by IP change. +resolvers docker + nameserver dns 127.0.0.11:53 + resolve_retries 3 + timeout resolve 1s + timeout retry 1s + hold valid 2s + +frontend nostream_in + bind *:8008 + option forwardfor + default_backend nostream_relays + +backend nostream_relays + balance roundrobin + option httpchk + http-check send meth GET uri /readyz + http-check expect status 200 + + default-server check inter 2s fall 2 rise 1 resolvers docker resolve-prefer ipv4 init-addr libc,none + + server blue nostream-blue:8008 + server green nostream-green:8008 diff --git a/deploy/rolling-relay-recreate.sh b/deploy/rolling-relay-recreate.sh new file mode 100755 index 00000000..67b3d471 --- /dev/null +++ b/deploy/rolling-relay-recreate.sh @@ -0,0 +1,74 @@ +#!/usr/bin/env bash +set -euo pipefail + +# Rolling relay recreate for the HAProxy blue/green stack. Replaces one relay +# at a time so the other keeps serving traffic. +# +# Usage: +# ./rolling-relay-recreate.sh [/opt/nostream] +# +# Load the new image and run migrations before calling this. + +TARGET="${1:-/opt/nostream}" +COMPOSE_FILE="${COMPOSE_FILE:-docker-compose.haproxy.yml}" +RELAY_PORT="${RELAY_PORT:-8008}" + +cd "$TARGET" + +if [[ ! -f "$COMPOSE_FILE" ]]; then + echo "error: compose file not found: $TARGET/$COMPOSE_FILE" >&2 + exit 1 +fi + +compose() { + docker compose -f "$COMPOSE_FILE" "$@" +} + +relay_readyz_ok() { + local service=$1 + local cid + cid="$(compose ps -q "$service" 2>/dev/null || true)" + if [[ -z "$cid" ]]; then + return 1 + fi + docker exec "$cid" node -e \ + "fetch('http://127.0.0.1:${RELAY_PORT}/readyz').then(r=>process.exit(r.ok?0:1)).catch(()=>process.exit(1))" +} + +peer_for() { + case "$1" in + nostream-blue) echo nostream-green ;; + nostream-green) echo nostream-blue ;; + *) echo "error: unknown service $1" >&2; exit 1 ;; + esac +} + +require_healthy_peer() { + local peer=$1 + if [[ -z "$(compose ps -q "$peer")" ]]; then + echo "error: peer $peer is not running; start it before replacing the other relay" >&2 + exit 1 + fi + if ! relay_readyz_ok "$peer"; then + echo "error: peer $peer /readyz is not healthy; fix it before continuing" >&2 + exit 1 + fi +} + +for service in nostream-blue nostream-green; do + if [[ -z "$(compose ps -q "$service")" ]]; then + echo "Starting $service (not running)..." + compose rm -f "$service" >/dev/null 2>&1 || true + compose up -d --no-deps --wait "$service" + continue + fi + + require_healthy_peer "$(peer_for "$service")" + + echo "Replacing $service..." + compose stop "$service" + compose rm -f "$service" + compose up -d --no-deps --wait "$service" +done + +echo "Rolling recreate complete" diff --git a/src/app/app.ts b/src/app/app.ts index d1b81ad8..707a5ec4 100644 --- a/src/app/app.ts +++ b/src/app/app.ts @@ -12,7 +12,14 @@ import { Settings } from '../@types/settings' import { SettingsStatic } from '../utils/settings' import { shutdownMetricsTelemetry } from '../telemetry/metrics' import { OperatorNotificationEventType } from '../@types/operator-notifications' +import { RedisRelayBroadcastFanout } from '../relay-broadcast/redis-relay-broadcast-fanout' import { enqueueOperatorNotification } from '../utils/operator-notification-enqueue' +import { RelayBroadcastDeduplicator } from '../utils/relay-broadcast-deduplicator' +import { + isRelayBroadcastFanoutEnabled, + isRelayBroadcastMessage, + RelayBroadcastMessage, +} from '../utils/relay-broadcast-message' import { getPrimaryShutdownDeadlineMs } from '../utils/shutdown-state' const logger = createLogger('app-primary') @@ -20,6 +27,8 @@ const logger = createLogger('app-primary') export class App implements IRunnable { private workers: WeakMap> private watchers: FSWatcher[] | undefined + private relayBroadcastFanout: RedisRelayBroadcastFanout | undefined + private readonly relayBroadcastDeduplicator = new RelayBroadcastDeduplicator() private shuttingDown = false public constructor( @@ -133,6 +142,19 @@ export class App implements IRunnable { logger('settings: %O', settings) + if (isRelayBroadcastFanoutEnabled()) { + this.relayBroadcastFanout = new RedisRelayBroadcastFanout() + void this.relayBroadcastFanout + .start((message) => this.onRelayBroadcastFromPeer(message)) + .then(() => { + logCentered('Relay broadcast fan-out enabled (Redis stream)', width) + }) + .catch((error) => { + logger.error('relay broadcast fan-out failed to start: %o', error) + this.relayBroadcastFanout = undefined + }) + } + // Primary-only: one outbox event per process start (not per client worker). void enqueueOperatorNotification(OperatorNotificationEventType.RELAY_RESTARTED, { version: packageJson.version, @@ -152,8 +174,31 @@ export class App implements IRunnable { private onClusterMessage(source: Worker, message: Serializable) { logger('message received from worker %s: %o', source.process.pid, message) + + if (isRelayBroadcastMessage(message)) { + this.relayBroadcastDeduplicator.mark(message.event.id) + const publishPromise = this.relayBroadcastFanout?.publish(message) + void publishPromise?.catch((error) => { + logger.error('relay broadcast publish failed: %o', error) + }) + } + + this.fanOutClusterMessage(message, source.id) + } + + private onRelayBroadcastFromPeer(message: RelayBroadcastMessage) { + if (this.relayBroadcastDeduplicator.has(message.event.id)) { + logger('skipping duplicate relay broadcast %s', message.event.id) + return + } + + this.relayBroadcastDeduplicator.mark(message.event.id) + this.fanOutClusterMessage(message) + } + + private fanOutClusterMessage(message: Serializable, excludeWorkerId?: number) { for (const worker of Object.values(this.cluster.workers as any) as Worker[]) { - if (source.id === worker.id) { + if (excludeWorkerId !== undefined && worker.id === excludeWorkerId) { continue } @@ -196,12 +241,15 @@ export class App implements IRunnable { let remaining = workers.length let finished = false + let deadline: NodeJS.Timeout | undefined const finishOnce = () => { if (finished) { return } finished = true - clearTimeout(deadline) + if (deadline !== undefined) { + clearTimeout(deadline) + } this.finishExit() } @@ -212,7 +260,7 @@ export class App implements IRunnable { } } - const deadline = setTimeout(() => { + deadline = setTimeout(() => { logger.warn('shutdown deadline exceeded, exiting primary') finishOnce() }, getPrimaryShutdownDeadlineMs()) @@ -238,8 +286,11 @@ export class App implements IRunnable { watcher.close() } } - if (typeof callback === 'function') { - callback() - } + const stopFanout = this.relayBroadcastFanout?.stop() ?? Promise.resolve() + void stopFanout.finally(() => { + if (typeof callback === 'function') { + callback() + } + }) } } diff --git a/src/relay-broadcast/redis-relay-broadcast-fanout.ts b/src/relay-broadcast/redis-relay-broadcast-fanout.ts new file mode 100644 index 00000000..532a0402 --- /dev/null +++ b/src/relay-broadcast/redis-relay-broadcast-fanout.ts @@ -0,0 +1,161 @@ +import { createClient } from 'redis' +import { hostname } from 'os' + +import { CacheClient } from '../@types/cache' +import { createLogger } from '../factories/logger-factory' +import { getCacheConfig } from '../cache/client' +import { + getRelayBroadcastStreamKey, + getRelayBroadcastStreamMaxLen, + isRelayBroadcastMessage, + RelayBroadcastMessage, +} from '../utils/relay-broadcast-message' + +const logger = createLogger('relay-broadcast-fanout') + +const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)) + +type StreamPayload = RelayBroadcastMessage & { originInstanceId: string } + +export class RedisRelayBroadcastFanout { + private publisher: CacheClient | undefined + private subscriber: CacheClient | undefined + private running = false + private readLoopPromise: Promise | undefined + private lastStreamId = '$' + + public constructor( + private readonly streamKey = getRelayBroadcastStreamKey(), + private readonly instanceId = process.env.RELAY_INSTANCE_ID?.trim() || `${hostname()}:${process.pid}`, + ) {} + + public async start(onMessage: (message: RelayBroadcastMessage) => void): Promise { + if (this.running) { + return + } + + this.running = true + const config = getCacheConfig() + + this.publisher = createClient(config) + this.subscriber = createClient(config) + + this.publisher.on('error', (error) => logger.error('publisher error: %o', error)) + this.subscriber.on('error', (error) => logger.error('subscriber error: %o', error)) + + await this.publisher.connect() + await this.subscriber.connect() + + logger('connected (stream=%s, instance=%s)', this.streamKey, this.instanceId) + + this.readLoopPromise = this.readLoop(onMessage) + } + + public async publish(message: RelayBroadcastMessage): Promise { + if (!this.publisher?.isOpen) { + logger.warn('publish skipped: publisher not connected') + return + } + + const payload: StreamPayload = { + ...message, + originInstanceId: this.instanceId, + } + + await this.publisher.xAdd( + this.streamKey, + '*', + { + payload: JSON.stringify(payload), + }, + { + TRIM: { + strategy: 'MAXLEN', + strategyModifier: '~', + threshold: getRelayBroadcastStreamMaxLen(), + }, + }, + ) + } + + public async stop(): Promise { + this.running = false + + if (this.readLoopPromise) { + await this.readLoopPromise.catch(() => undefined) + this.readLoopPromise = undefined + } + + await Promise.all([ + this.subscriber?.isOpen ? this.subscriber.disconnect() : Promise.resolve(), + this.publisher?.isOpen ? this.publisher.disconnect() : Promise.resolve(), + ]) + + this.subscriber = undefined + this.publisher = undefined + } + + private async readLoop(onMessage: (message: RelayBroadcastMessage) => void): Promise { + while (this.running && this.subscriber?.isOpen) { + try { + const response = await this.subscriber.xRead( + { + key: this.streamKey, + id: this.lastStreamId, + }, + { + BLOCK: 5000, + COUNT: 32, + }, + ) + + if (!response) { + continue + } + + for (const stream of response) { + for (const entry of stream.messages) { + this.lastStreamId = entry.id + const rawPayload = entry.message.payload + + if (typeof rawPayload !== 'string') { + continue + } + + let parsed: StreamPayload + try { + parsed = JSON.parse(rawPayload) as StreamPayload + } catch (error) { + logger.warn('invalid stream payload: %o', error) + continue + } + + if (parsed.originInstanceId === this.instanceId) { + continue + } + + const relayMessage: RelayBroadcastMessage = { + eventName: parsed.eventName, + event: parsed.event, + source: parsed.source, + } + + if (!isRelayBroadcastMessage(relayMessage)) { + logger.warn('skipping invalid relay broadcast stream entry') + continue + } + + onMessage(relayMessage) + } + } + } catch (error) { + if (!this.running) { + return + } + + logger.error('stream read failed: %o', error) + await sleep(1000) + } + } + } +} diff --git a/src/utils/relay-broadcast-deduplicator.ts b/src/utils/relay-broadcast-deduplicator.ts new file mode 100644 index 00000000..ec3c7e3f --- /dev/null +++ b/src/utils/relay-broadcast-deduplicator.ts @@ -0,0 +1,25 @@ +export class RelayBroadcastDeduplicator { + private readonly seen = new Map() + + public constructor(private readonly ttlMs = 120_000) {} + + public has(eventId: string): boolean { + this.prune() + const expiresAt = this.seen.get(eventId) + return expiresAt !== undefined && expiresAt > Date.now() + } + + public mark(eventId: string): void { + this.prune() + this.seen.set(eventId, Date.now() + this.ttlMs) + } + + private prune(): void { + const now = Date.now() + for (const [eventId, expiresAt] of this.seen) { + if (expiresAt <= now) { + this.seen.delete(eventId) + } + } + } +} diff --git a/src/utils/relay-broadcast-message.ts b/src/utils/relay-broadcast-message.ts new file mode 100644 index 00000000..7a147043 --- /dev/null +++ b/src/utils/relay-broadcast-message.ts @@ -0,0 +1,51 @@ +import { Serializable } from 'child_process' + +import { Event } from '../@types/event' +import { WebSocketServerAdapterEvent } from '../constants/adapter' + +export type RelayBroadcastMessage = { + eventName: WebSocketServerAdapterEvent.Broadcast + event: Event + source?: string +} + +export const isRelayBroadcastMessage = (message: Serializable): message is RelayBroadcastMessage => { + if (typeof message !== 'object' || message === null) { + return false + } + + const candidate = message as RelayBroadcastMessage + + return ( + candidate.eventName === WebSocketServerAdapterEvent.Broadcast && + typeof candidate.event === 'object' && + candidate.event !== null && + typeof candidate.event.id === 'string' + ) +} + +export const isRelayBroadcastFanoutEnabled = (): boolean => { + const value = process.env.RELAY_BROADCAST_FANOUT?.trim().toLowerCase() + return value === '1' || value === 'true' || value === 'yes' +} + +export const getRelayBroadcastStreamKey = (): string => { + const key = process.env.RELAY_BROADCAST_STREAM_KEY?.trim() + return key && key.length > 0 ? key : 'nostream:relay:broadcast' +} + +const DEFAULT_STREAM_MAXLEN = 50_000 + +export const getRelayBroadcastStreamMaxLen = (): number => { + const raw = process.env.RELAY_BROADCAST_STREAM_MAXLEN?.trim() + if (raw === undefined || raw === '') { + return DEFAULT_STREAM_MAXLEN + } + + const parsed = Number(raw) + if (!Number.isFinite(parsed) || parsed < 1) { + return DEFAULT_STREAM_MAXLEN + } + + return Math.floor(parsed) +} diff --git a/test/unit/app/app.spec.ts b/test/unit/app/app.spec.ts index 9639bd74..5daa30b2 100644 --- a/test/unit/app/app.spec.ts +++ b/test/unit/app/app.spec.ts @@ -80,7 +80,9 @@ describe('App', () => { }) sigtermHandler() - await Promise.resolve() + await new Promise((resolve) => { + setImmediate(resolve) + }) expect(worker.kill).to.have.been.calledOnce expect(fakeProcess.exit).to.have.been.calledOnceWithExactly(0) diff --git a/test/unit/relay-broadcast/redis-relay-broadcast-fanout.spec.ts b/test/unit/relay-broadcast/redis-relay-broadcast-fanout.spec.ts new file mode 100644 index 00000000..a12de108 --- /dev/null +++ b/test/unit/relay-broadcast/redis-relay-broadcast-fanout.spec.ts @@ -0,0 +1,137 @@ +import chai from 'chai' +import Sinon from 'sinon' +import sinonChai from 'sinon-chai' + +import { WebSocketServerAdapterEvent } from '../../../src/constants/adapter' +import * as redis from 'redis' +import { RedisRelayBroadcastFanout } from '../../../src/relay-broadcast/redis-relay-broadcast-fanout' + +chai.use(sinonChai) + +const { expect } = chai + +describe('RedisRelayBroadcastFanout', () => { + let sandbox: Sinon.SinonSandbox + let publisher: any + let subscriber: any + let xReadStub: Sinon.SinonStub + + beforeEach(() => { + sandbox = Sinon.createSandbox() + + publisher = { + connect: sandbox.stub().resolves(), + disconnect: sandbox.stub().resolves(), + on: sandbox.stub(), + isOpen: true, + xAdd: sandbox.stub().resolves('1-0'), + } + + xReadStub = sandbox.stub() + subscriber = { + connect: sandbox.stub().resolves(), + disconnect: sandbox.stub().resolves(), + on: sandbox.stub(), + isOpen: true, + xRead: xReadStub, + } + + sandbox.stub(redis, 'createClient').callsFake(() => { + if (!publisher._used) { + publisher._used = true + return publisher + } + + return subscriber + }) + }) + + afterEach(() => { + sandbox.restore() + }) + + it('publishes with stream trim options', async () => { + const fanout = new RedisRelayBroadcastFanout('test:stream', 'instance-a') + await fanout.start(() => undefined) + + await fanout.publish({ + eventName: WebSocketServerAdapterEvent.Broadcast, + event: { id: 'aa'.repeat(32) } as any, + }) + + expect(publisher.xAdd).to.have.been.calledOnce + const trimOptions = publisher.xAdd.firstCall.args[3] + expect(trimOptions).to.deep.include({ + TRIM: { + strategy: 'MAXLEN', + strategyModifier: '~', + threshold: 50_000, + }, + }) + + await fanout.stop() + }) + + it('skips self-origin and malformed stream entries', async () => { + xReadStub.callsFake(async () => { + if (xReadStub.callCount > 1) { + return null + } + + return [ + { + name: 'test:stream', + messages: [ + { + id: '1-0', + message: { + payload: JSON.stringify({ + originInstanceId: 'instance-a', + eventName: WebSocketServerAdapterEvent.Broadcast, + event: { id: 'bb'.repeat(32) }, + }), + }, + }, + { + id: '2-0', + message: { + payload: JSON.stringify({ + originInstanceId: 'instance-b', + eventName: WebSocketServerAdapterEvent.Broadcast, + event: null, + }), + }, + }, + { + id: '3-0', + message: { + payload: JSON.stringify({ + originInstanceId: 'instance-b', + eventName: WebSocketServerAdapterEvent.Broadcast, + event: { id: 'cc'.repeat(32) }, + }), + }, + }, + ], + }, + ] + }) + + const fanout = new RedisRelayBroadcastFanout('test:stream', 'instance-a') + const received: string[] = [] + + await fanout.start((message) => { + received.push(message.event.id) + }) + + for (let attempt = 0; attempt < 20 && received.length === 0; attempt++) { + await new Promise((resolve) => { + setImmediate(resolve) + }) + } + + expect(received).to.deep.equal(['cc'.repeat(32)]) + + await fanout.stop() + }) +}) diff --git a/test/unit/utils/relay-broadcast-deduplicator.spec.ts b/test/unit/utils/relay-broadcast-deduplicator.spec.ts new file mode 100644 index 00000000..900a176d --- /dev/null +++ b/test/unit/utils/relay-broadcast-deduplicator.spec.ts @@ -0,0 +1,28 @@ +import chai from 'chai' +import sinon from 'sinon' + +import { RelayBroadcastDeduplicator } from '../../../src/utils/relay-broadcast-deduplicator' + +const { expect } = chai + +describe('RelayBroadcastDeduplicator', () => { + let clock: sinon.SinonFakeTimers + + beforeEach(() => { + clock = sinon.useFakeTimers() + }) + + afterEach(() => { + clock.restore() + }) + + it('tracks event ids for the configured ttl', () => { + const deduplicator = new RelayBroadcastDeduplicator(1000) + + deduplicator.mark('abc') + expect(deduplicator.has('abc')).to.equal(true) + + clock.tick(1001) + expect(deduplicator.has('abc')).to.equal(false) + }) +}) diff --git a/test/unit/utils/relay-broadcast-message.spec.ts b/test/unit/utils/relay-broadcast-message.spec.ts new file mode 100644 index 00000000..d6ae09ad --- /dev/null +++ b/test/unit/utils/relay-broadcast-message.spec.ts @@ -0,0 +1,31 @@ +import chai from 'chai' + +import { WebSocketServerAdapterEvent } from '../../../src/constants/adapter' +import { + getRelayBroadcastStreamMaxLen, + isRelayBroadcastMessage, +} from '../../../src/utils/relay-broadcast-message' + +const { expect } = chai + +describe('relay-broadcast-message', () => { + it('detects cluster broadcast messages', () => { + expect( + isRelayBroadcastMessage({ + eventName: WebSocketServerAdapterEvent.Broadcast, + event: { id: '00'.repeat(32) }, + }), + ).to.equal(true) + + expect(isRelayBroadcastMessage({ eventName: 'other' })).to.equal(false) + }) + + it('defaults stream maxlen when env is unset', () => { + const previous = process.env.RELAY_BROADCAST_STREAM_MAXLEN + delete process.env.RELAY_BROADCAST_STREAM_MAXLEN + expect(getRelayBroadcastStreamMaxLen()).to.equal(50_000) + if (previous !== undefined) { + process.env.RELAY_BROADCAST_STREAM_MAXLEN = previous + } + }) +})