Files
web-terminal/relay-run/tests/wiring/revocation-subscriber.test.ts
Yaojia Wang aa1912b962 feat(relay): Phase1 waves A2-E — server entry, shared-store data plane, agent runtime, deploy artifacts
RELAY-PHASE1 Wave A2/B/C/D/E (12-agent workflow, all tsc-clean, 314/314 tests pass):
- A2: control-plane server.ts entry + boot/redis.ts revocation-bus wiring + start script.
- B1: relay-run shared-store EnforceDeps (relay-auth ports over the SAME Postgres+Redis as P3).
- B2: registry-backed MtlsVerifier (verifyAgentCert, fail-closed, INV14).
- B3: store-backed RouteResolver (subdomain->hostId).
- B4: Redis relay:revocations subscriber -> tunnel teardown (INV12).
- B5: main-phase1.ts production entry (public bind, real TLS, async-mTLS prefetch bridge) + staging /auth/mint.
- B6 (PARTIAL): relay-web operator login + browser DPoP; proof offered via term.dpop.<b64u> subprotocol.
- C: agent dist/cli.js build (esbuild) + runTunnel run-loop + CliDeps.
- D1: same-origin static serve of relay-web/public from the browser WSS.
- E: systemd units + gen-ca/gen-capability-key/issue-tls-cert scripts + deploy/RUNBOOK.md.
Adversarial review: all hard invariants PASS. Follow-ups (B7): close DPoP-subprotocol read on
browser-server (blocks browser connect); rate-limit /auth/mint (F1); wire activeSessionCount (F2);
scrub error logs (F5). Excludes unrelated public/style.css (concurrent iOS job).
2026-07-06 16:13:34 +02:00

238 lines
8.5 KiB
TypeScript

import { describe, it, expect, vi } from 'vitest'
import { randomUUID } from 'node:crypto'
import { RELAY_REVOCATIONS_CHANNEL, type KillSignal } from 'relay-contracts'
import {
startRevocationSubscriber,
type ActiveTunnelRef,
type RedisSubscriber,
} from '../../src/wiring/revocation-subscriber.js'
/** Fake ioredis subscriber-mode client: records subscribe/unsubscribe and lets a test emit frames. */
class FakeRedisSubscriber implements RedisSubscriber {
readonly subscribed: string[] = []
readonly unsubscribed: string[] = []
private readonly handlers = new Set<(channel: string, message: string) => void>()
async subscribe(channel: string): Promise<number> {
this.subscribed.push(channel)
return this.subscribed.length
}
async unsubscribe(channel: string): Promise<number> {
this.unsubscribed.push(channel)
return this.unsubscribed.length
}
on(_event: 'message', listener: (channel: string, message: string) => void): this {
this.handlers.add(listener)
return this
}
off(_event: 'message', listener: (channel: string, message: string) => void): this {
this.handlers.delete(listener)
return this
}
/** Test driver: deliver a raw pub/sub frame to every registered listener. */
emit(channel: string, message: string): void {
for (const h of [...this.handlers]) h(channel, message)
}
get listenerCount(): number {
return this.handlers.size
}
}
/** Fake relay node: a live-tunnel map + a closeStream that records + removes the host (models the
* real whole-host teardown mutating the tunnel set, so snapshot-safety is exercised). */
function makeNode(initial: readonly ActiveTunnelRef[]) {
const tunnels = new Map(initial.map((t) => [t.hostId, t]))
const closed: string[] = []
return {
closed,
// Live iterator on purpose: the subscriber must snapshot before tearing down.
activeTunnels: () => tunnels.values(),
closeStream: (hostId: string): void => {
closed.push(hostId)
tunnels.delete(hostId)
},
}
}
const AT = 1_700_000_000
function killMessage(signal: KillSignal): string {
return JSON.stringify(signal)
}
describe('startRevocationSubscriber', () => {
it('subscribes to the relay:revocations channel on start', () => {
const redisSubscriber = new FakeRedisSubscriber()
startRevocationSubscriber({ redisSubscriber, node: makeNode([]) })
expect(redisSubscriber.subscribed).toEqual([RELAY_REVOCATIONS_CHANNEL])
expect(redisSubscriber.listenerCount).toBe(1)
})
it('tears down only the host a host-scoped signal names', () => {
const hostA = randomUUID()
const hostB = randomUUID()
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([
{ hostId: hostA, accountId: randomUUID() },
{ hostId: hostB, accountId: randomUUID() },
])
const onApplied = vi.fn()
startRevocationSubscriber({ redisSubscriber, node, onApplied })
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
killMessage({ scope: { kind: 'host', hostId: hostA }, at: AT, reason: 'compromised' }),
)
expect(node.closed).toEqual([hostA])
expect(onApplied).toHaveBeenCalledTimes(1)
expect(onApplied).toHaveBeenCalledWith(expect.objectContaining({ at: AT }), 1)
})
it('tears down every host under an account-scoped signal, leaving other accounts running', () => {
const acct1 = randomUUID()
const acct2 = randomUUID()
const hostA = randomUUID()
const hostB = randomUUID()
const hostC = randomUUID()
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([
{ hostId: hostA, accountId: acct1 },
{ hostId: hostB, accountId: acct1 },
{ hostId: hostC, accountId: acct2 },
])
startRevocationSubscriber({ redisSubscriber, node })
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
killMessage({ scope: { kind: 'account', accountId: acct1 }, at: AT, reason: 'billing' }),
)
expect(node.closed.sort()).toEqual([hostA, hostB].sort())
expect(node.closed).not.toContain(hostC)
})
it('tears down every live host on a global-scoped signal', () => {
const hosts = [
{ hostId: randomUUID(), accountId: randomUUID() },
{ hostId: randomUUID(), accountId: randomUUID() },
{ hostId: randomUUID(), accountId: randomUUID() },
]
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode(hosts)
startRevocationSubscriber({ redisSubscriber, node })
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
killMessage({ scope: { kind: 'global' }, at: AT, reason: 'kill-switch' }),
)
expect(node.closed.sort()).toEqual(hosts.map((h) => h.hostId).sort())
})
it('is a no-op for a host-scoped signal naming a host this node does not serve', () => {
const served = randomUUID()
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([{ hostId: served, accountId: randomUUID() }])
const onApplied = vi.fn()
startRevocationSubscriber({ redisSubscriber, node, onApplied })
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
killMessage({ scope: { kind: 'host', hostId: randomUUID() }, at: AT, reason: 'other-node' }),
)
expect(node.closed).toEqual([])
expect(onApplied).toHaveBeenCalledWith(expect.anything(), 0)
})
it('drops a malformed (non-JSON) message: no teardown, counted as dropped', () => {
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([{ hostId: randomUUID(), accountId: randomUUID() }])
const onDropped = vi.fn()
const onApplied = vi.fn()
startRevocationSubscriber({ redisSubscriber, node, onDropped, onApplied })
redisSubscriber.emit(RELAY_REVOCATIONS_CHANNEL, '{ this is not json')
expect(node.closed).toEqual([])
expect(onDropped).toHaveBeenCalledTimes(1)
expect(onApplied).not.toHaveBeenCalled()
})
it('drops a schema-invalid signal (non-uuid host / missing fields): no teardown', () => {
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([{ hostId: randomUUID(), accountId: randomUUID() }])
const onDropped = vi.fn()
startRevocationSubscriber({ redisSubscriber, node, onDropped })
// hostId is not a UUID → RevocationScopeSchema rejects.
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
JSON.stringify({ scope: { kind: 'host', hostId: 'not-a-uuid' }, at: AT, reason: 'x' }),
)
// Missing `at` → KillSignalSchema rejects.
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
JSON.stringify({ scope: { kind: 'global' }, reason: 'x' }),
)
expect(node.closed).toEqual([])
expect(onDropped).toHaveBeenCalledTimes(2)
})
it('ignores messages published on a different channel', () => {
const hostA = randomUUID()
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([{ hostId: hostA, accountId: randomUUID() }])
const onDropped = vi.fn()
const onApplied = vi.fn()
startRevocationSubscriber({ redisSubscriber, node, onDropped, onApplied })
redisSubscriber.emit(
'some:other:channel',
killMessage({ scope: { kind: 'host', hostId: hostA }, at: AT, reason: 'wrong-channel' }),
)
expect(node.closed).toEqual([])
expect(onDropped).not.toHaveBeenCalled()
expect(onApplied).not.toHaveBeenCalled()
})
it('close() removes the message listener, unsubscribes, and is idempotent', () => {
const hostA = randomUUID()
const redisSubscriber = new FakeRedisSubscriber()
const node = makeNode([{ hostId: hostA, accountId: randomUUID() }])
const sub = startRevocationSubscriber({ redisSubscriber, node })
sub.close()
sub.close() // idempotent — no throw, no double-unsubscribe
expect(redisSubscriber.unsubscribed).toEqual([RELAY_REVOCATIONS_CHANNEL])
expect(redisSubscriber.listenerCount).toBe(0)
// A frame delivered after close must not tear anything down.
redisSubscriber.emit(
RELAY_REVOCATIONS_CHANNEL,
killMessage({ scope: { kind: 'host', hostId: hostA }, at: AT, reason: 'after-close' }),
)
expect(node.closed).toEqual([])
})
it('routes a rejected subscribe() to onError instead of swallowing it', async () => {
const failing: RedisSubscriber = {
subscribe: () => Promise.reject(new Error('redis down')),
unsubscribe: async () => 0,
on: () => failing,
off: () => failing,
}
const onError = vi.fn()
startRevocationSubscriber({ redisSubscriber: failing, node: makeNode([]), onError })
await Promise.resolve() // let the rejected subscribe() microtask settle
expect(onError).toHaveBeenCalledTimes(1)
expect(onError.mock.calls[0][0]).toBeInstanceOf(Error)
})
})