diff --git a/HOW_TO_INSTALL.md b/HOW_TO_INSTALL.md index bb22ae21e..428496ee8 100644 --- a/HOW_TO_INSTALL.md +++ b/HOW_TO_INSTALL.md @@ -200,7 +200,7 @@ docker compose cp cassandra/.cql cassandra:/tmp/x.cql docker compose exec cassandra cqlsh -f /tmp/x.cql ``` -**Stop** the node (the database volume is kept): +**Stop** the node (the database and node-state volumes are kept): ```bash docker compose -f docker-compose.yml -f docker-compose.caddy.yml down diff --git a/POMBO.md b/POMBO.md index 7d935f536..155b69c11 100644 --- a/POMBO.md +++ b/POMBO.md @@ -77,9 +77,17 @@ Reading the streams of a gated channel requires a signed request (`plugins.storage.signedReads.enabled`, on by default), and the signer must have access to the channel right now (owner, moderator, or accepted by the gate's `checkAccess`). The admin stream (`-3`) stays open: channel previews and -the entry screen of non-members are served from it. Streams outside -gated channels are unaffected. While the chain cannot be consulted the -node answers 503 rather than serve the data. +the entry screen of non-members are served from it. A stream outside gated +channels is open when its SUBSCRIBE is public, as on a vanilla node; +otherwise (a DM inbox) the signer must hold SUBSCRIBE on it. + +While the chain cannot be consulted the node answers 503 rather than serve +the data, except for a stream whose latest answer from the chain was no gate +and public SUBSCRIBE: that one is served. Those streams are listed in +`~/.streamr/known-public-streams.json`, so the list holds after a restart +(and after a recreate, in the `node-state` volume of the bundled `docker compose`). +Every answer from the chain replaces the entry, and a stream the node never +asked about still gets 503. Headers: `x-pombo-user`, `x-pombo-issued-at`, `x-pombo-nonce`, `x-pombo-signature`. @@ -105,11 +113,13 @@ phase if any stream errors for another reason (an unstable RPC looks like a deletion otherwise) or if a suspiciously large fraction of streams look deleted, and holds a grace period before removing anything. -The first run comes a minute after the node starts, unless a run started less -than `retention.intervalHours` ago: the start of each run is kept in +The first run comes ten minutes after the node starts, leaving the chain RPC +to the permission lookups of the first reads, unless a run started less than +`retention.intervalHours` ago: the start of each run is kept in `~/.streamr/retention-last-run`, so a node that keeps restarting does not -repeat a full run on every start. Inside the container that file survives a -restart but not a recreate. +repeat a full run on every start. In the bundled `docker compose`, +`~/.streamr` is the `node-state` volume, so the file also survives a +recreate. In a cluster the deletes replicate through Cassandra, so retention runs on **one node only**: the installer leaves `retention.enabled` true on the @@ -171,7 +181,7 @@ Storage plugin keys added to the upstream ones: | `bucket.checkFullBucketsTimeout` | 250 | ms between checks for full buckets | | `read.fetchSize` | 32 | rows per page when streaming stored messages (128 upstream); Cassandra and the node hold a whole page in memory, and a page of 128 file chunks is ~30 MB | | `batch.logErrors` | true | log failed batch inserts (upstream retries them silently) | -| `signedReads.enabled` | true | require signed reads on gated channels | +| `signedReads.enabled` | true | require signed reads on gated channels and on streams whose SUBSCRIBE is not public (a DM inbox) | | `retention.enabled` | true | prune stored data past each stream's storageDays (one machine per cluster) | | `retention.intervalHours` | 6 | how often retention runs | | `retention.graceDays` | 7 | hold before deleting an on-chain-deleted stream's data | diff --git a/deploy/docker-compose.yml b/deploy/docker-compose.yml index f60d1ed7d..c37546d8f 100644 --- a/deploy/docker-compose.yml +++ b/deploy/docker-compose.yml @@ -64,6 +64,8 @@ services: LOG_LEVEL: info NODE_OPTIONS: --max-semi-space-size=128 --max-old-space-size=${NODE_OLD_SPACE_MB:-2048} volumes: + # The node's own state (known public streams, last retention run) must outlive a recreate. + - node-state:/home/streamr/.streamr - ./config:/home/streamr/.streamr/config:ro ports: - "${POMBO_WS_PORT:-32200}:${POMBO_WS_PORT:-32200}" @@ -72,3 +74,4 @@ services: volumes: cassandra-data: + node-state: diff --git a/packages/node/src/plugins/storage/KnownPublicStreams.ts b/packages/node/src/plugins/storage/KnownPublicStreams.ts new file mode 100644 index 000000000..e519a2e17 --- /dev/null +++ b/packages/node/src/plugins/storage/KnownPublicStreams.ts @@ -0,0 +1,87 @@ +import { Logger } from '@streamr/utils' +import { readFileSync, renameSync, writeFileSync } from 'fs' +import os from 'os' +import path from 'path' + +const logger = new Logger('KnownPublicStreams') + +export const DEFAULT_KNOWN_PUBLIC_STREAMS_FILE = path.join(os.homedir(), '.streamr', 'known-public-streams.json') +const SAVE_DELAY_MS = 5 * 1000 + +const load = (file: string): string[] => { + let content: string + try { + content = readFileSync(file, 'utf8') + } catch { + return [] + } + try { + const parsed = JSON.parse(content) + if (Array.isArray(parsed) && parsed.every((id) => typeof id === 'string')) { + return parsed + } + } catch { + // fall through + } + logger.warn('Ignoring an unreadable list of public streams', { file }) + return [] +} + +/** + * The streams whose latest answer from the chain was: no Pombo gate and public + * SUBSCRIBE. Kept on disk, when given a file, so the answer outlives a restart. + */ +export class KnownPublicStreams { + + private readonly file?: string + private readonly streams: Set + private saveTimeout?: NodeJS.Timeout + + constructor(file?: string) { + this.file = file + this.streams = new Set((file !== undefined) ? load(file) : []) + } + + has(streamId: string): boolean { + return this.streams.has(streamId) + } + + set(streamId: string, isPublic: boolean): void { + if (isPublic === this.streams.has(streamId)) { + return + } + if (isPublic) { + this.streams.add(streamId) + } else { + this.streams.delete(streamId) + } + this.scheduleSave() + } + + flush(): void { + if (this.saveTimeout !== undefined) { + clearTimeout(this.saveTimeout) + this.saveTimeout = undefined + this.save() + } + } + + private scheduleSave(): void { + if ((this.file !== undefined) && (this.saveTimeout === undefined)) { + this.saveTimeout = setTimeout(() => { + this.saveTimeout = undefined + this.save() + }, SAVE_DELAY_MS) + } + } + + private save(): void { + const tmpFile = `${this.file!}.tmp` + try { + writeFileSync(tmpFile, JSON.stringify([...this.streams])) + renameSync(tmpFile, this.file!) + } catch (err) { + logger.warn('Could not save the list of public streams', { err, file: this.file }) + } + } +} diff --git a/packages/node/src/plugins/storage/PomboGates.ts b/packages/node/src/plugins/storage/PomboGates.ts index a22dae3ed..1077bf1e5 100644 --- a/packages/node/src/plugins/storage/PomboGates.ts +++ b/packages/node/src/plugins/storage/PomboGates.ts @@ -1,6 +1,7 @@ -import { StreamPermission, StreamrClient } from '@streamr/sdk' +import { StreamMetadata, StreamPermission, StreamrClient } from '@streamr/sdk' import { EthereumAddress, MapWithTtl, toEthereumAddress } from '@streamr/utils' import { Contract } from 'ethers' +import { KnownPublicStreams } from './KnownPublicStreams' const POMBO_GATE_ABI = [ 'function owner() view returns (address)', @@ -122,10 +123,16 @@ export class PomboGates { private readonly accessCache = new MapWithTtl(ttlOf) private readonly subscribePublicCache = new MapWithTtl(ttlOf) private readonly subscribeUserCache = new MapWithTtl(ttlOf) + private readonly knownPublic: KnownPublicStreams - constructor(client: StreamrClient, gateReader: GateReader = createEthersGateReader(client)) { + constructor( + client: StreamrClient, + gateReader: GateReader = createEthersGateReader(client), + knownPublic = new KnownPublicStreams() + ) { this.client = client this.gateReader = gateReader + this.knownPublic = knownPublic } async getGate(streamId: string): Promise { @@ -133,7 +140,15 @@ export class PomboGates { if (cached !== undefined) { return cached } - const metadata = await this.client.getStreamMetadata(streamId) + let metadata: StreamMetadata + try { + metadata = await this.client.getStreamMetadata(streamId) + } catch (err) { + if ((err as { code?: string }).code === 'STREAM_NOT_FOUND') { + this.knownPublic.set(streamId, false) + } + throw err + } let gateAddress = parseGateAddress(metadata) if (gateAddress === undefined) { const linked = parseLinkedStream(metadata) @@ -141,6 +156,9 @@ export class PomboGates { gateAddress = parseGateAddress(await this.client.getStreamMetadata(linked)) } } + if (gateAddress !== undefined) { + this.knownPublic.set(streamId, false) + } const info = (gateAddress !== undefined) ? await this.gateReader.getInfo(gateAddress) : null this.gateCache.set(streamId, info) return info @@ -154,7 +172,7 @@ export class PomboGates { return this.cachedLookup(this.accessCache, gateAddress, user, () => this.gateReader.checkAccess(gateAddress, user)) } - /** Whether the stream can be subscribed to by anyone (a public stream). */ + /** Whether the stream can be subscribed to by anyone (a public stream). Asked only of streams with no gate. */ async isPublicSubscribe(streamId: string): Promise { const cached = this.subscribePublicCache.get(streamId) if (cached !== undefined) { @@ -162,9 +180,18 @@ export class PomboGates { } const result = await this.client.hasPermission({ streamId, permission: StreamPermission.SUBSCRIBE, public: true }) this.subscribePublicCache.set(streamId, result) + this.knownPublic.set(streamId, result) return result } + /** + * Whether the chain's latest answer about the stream, however old, was no + * gate and public SUBSCRIBE. Only for when the chain cannot answer now. + */ + wasLastSeenPublic(streamId: string): boolean { + return this.knownPublic.has(streamId) + } + /** Whether `user` holds SUBSCRIBE on the stream (not counting a public grant). */ async hasSubscribe(streamId: string, user: EthereumAddress): Promise { const key = `${streamId}_${user}` @@ -196,6 +223,7 @@ export class PomboGates { this.accessCache.clear() this.subscribePublicCache.clear() this.subscribeUserCache.clear() + this.knownPublic.flush() } // eslint-disable-next-line class-methods-use-this diff --git a/packages/node/src/plugins/storage/RetentionScheduler.ts b/packages/node/src/plugins/storage/RetentionScheduler.ts index 29761415a..8469507f5 100644 --- a/packages/node/src/plugins/storage/RetentionScheduler.ts +++ b/packages/node/src/plugins/storage/RetentionScheduler.ts @@ -17,7 +17,8 @@ const listStreamParts = async (client: Client): Promise<{ streamId: string, part const DAY_MS = 24 * 60 * 60 * 1000 const HOUR_MS = 60 * 60 * 1000 -const STARTUP_DELAY_MS = 60 * 1000 +// Leaves the chain RPC to the read-permission lookups while their caches warm up. +const STARTUP_DELAY_MS = 10 * 60 * 1000 const CLASSIFY_CONCURRENCY = 5 const DEFAULT_STATE_FILE = path.join(os.homedir(), '.streamr', 'retention-last-run') diff --git a/packages/node/src/plugins/storage/StoragePlugin.ts b/packages/node/src/plugins/storage/StoragePlugin.ts index 5e4f464fc..8006d755b 100644 --- a/packages/node/src/plugins/storage/StoragePlugin.ts +++ b/packages/node/src/plugins/storage/StoragePlugin.ts @@ -5,8 +5,9 @@ import { ApiPluginConfig, Plugin } from '../../Plugin' import { Storage, startCassandraStorage } from './Storage' import { CassandraWatchdog } from './CassandraWatchdog' import { IngestValidator } from './IngestValidator' +import { DEFAULT_KNOWN_PUBLIC_STREAMS_FILE, KnownPublicStreams } from './KnownPublicStreams' import { LoadBalancingPolicyFactory, cassandraContactPoints, createLocalHostPolicyFactory } from './localHostPolicy' -import { PomboGates } from './PomboGates' +import { PomboGates, createEthersGateReader } from './PomboGates' import { SignedRequestVerifier } from './SignedRequest' import { RetentionScheduler } from './RetentionScheduler' import { createCapabilitiesEndpoint } from './capabilitiesEndpoint' @@ -99,7 +100,11 @@ export class StoragePlugin extends Plugin { }) this.cassandraWatchdog.start() this.storageConfig = await this.startStorageConfig(clusterId, assignmentStream) - this.gates = new PomboGates(this.streamrClient) + this.gates = new PomboGates( + this.streamrClient, + createEthersGateReader(this.streamrClient), + new KnownPublicStreams(DEFAULT_KNOWN_PUBLIC_STREAMS_FILE) + ) this.ingestValidator = new IngestValidator(this.streamrClient, metricsContext, this.gates) this.signedRequestVerifier = new SignedRequestVerifier() this.messageListener = (msg) => { diff --git a/packages/node/src/plugins/storage/config.schema.json b/packages/node/src/plugins/storage/config.schema.json index 3b82682b1..7b28020ec 100644 --- a/packages/node/src/plugins/storage/config.schema.json +++ b/packages/node/src/plugins/storage/config.schema.json @@ -141,7 +141,7 @@ }, "signedReads": { "type": "object", - "description": "Require a signed request to read the streams of a gated Pombo channel (all but the admin stream).", + "description": "Require a signed request to read the streams of a gated Pombo channel (all but the admin stream) and streams whose SUBSCRIBE is not public, such as a DM inbox.", "additionalProperties": false, "properties": { "enabled": { diff --git a/packages/node/src/plugins/storage/signedReads.ts b/packages/node/src/plugins/storage/signedReads.ts index f7f254caa..fb98d9f07 100644 --- a/packages/node/src/plugins/storage/signedReads.ts +++ b/packages/node/src/plugins/storage/signedReads.ts @@ -46,8 +46,9 @@ const readEnvelope = (req: Request): unknown => { * now: for a gated channel the gate owner, a moderator, or an account the gate * accepts; for a non-gated stream, SUBSCRIBE on it. A stream whose SUBSCRIBE is * public is served without a signature, as a vanilla node serves it. When the - * chain cannot be consulted the read is refused: a read can be retried, a leak - * cannot be undone. + * chain cannot be consulted, a stream it last described as public is served + * and every other read is refused: a read can be retried, a leak cannot be + * undone. */ export const createSignedReadGuard = ( enabled: boolean, @@ -75,6 +76,11 @@ export const createSignedReadGuard = ( res.status(404).json({ error: 'Stream not found' }) return } + if (gates.wasLastSeenPublic(streamId)) { + logger.warn('Could not read gate, serving a stream last seen as public', { streamId, err }) + next() + return + } logger.warn('Could not read gate, refusing read', { streamId, err }) res.status(503).json({ error: 'Cannot verify access right now' }) return @@ -89,6 +95,11 @@ export const createSignedReadGuard = ( try { isPublic = await gates.isPublicSubscribe(streamId) } catch (err) { + if (gates.wasLastSeenPublic(streamId)) { + logger.warn('Could not read permissions, serving a stream last seen as public', { streamId, err }) + next() + return + } logger.warn('Could not read permissions, refusing read', { streamId, err }) res.status(503).json({ error: 'Cannot verify access right now' }) return diff --git a/packages/node/src/plugins/storage/storedEndpoint.ts b/packages/node/src/plugins/storage/storedEndpoint.ts index 5eb1c608c..7d4bc04ad 100644 --- a/packages/node/src/plugins/storage/storedEndpoint.ts +++ b/packages/node/src/plugins/storage/storedEndpoint.ts @@ -65,9 +65,13 @@ const createHandler = (storage: Storage, gates: PomboGates, client: StreamrClien try { canRead = await gates.canRead(streamId, signer) } catch (err) { - logger.warn('Could not verify access, refusing', { streamId, signer, err }) - res.status(503).json({ error: 'Cannot verify access right now' }) - return + if (!gates.wasLastSeenPublic(streamId)) { + logger.warn('Could not verify access, refusing', { streamId, signer, err }) + res.status(503).json({ error: 'Cannot verify access right now' }) + return + } + logger.warn('Could not verify access, answering for a stream last seen as public', { streamId, err }) + canRead = true } const results: (PurgeTarget & { result: StoredResult })[] = [] for (const target of targets) { diff --git a/packages/node/test/unit/plugins/storage/KnownPublicStreams.test.ts b/packages/node/test/unit/plugins/storage/KnownPublicStreams.test.ts new file mode 100644 index 000000000..0e2b407d9 --- /dev/null +++ b/packages/node/test/unit/plugins/storage/KnownPublicStreams.test.ts @@ -0,0 +1,69 @@ +import { existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'fs' +import os from 'os' +import path from 'path' +import { KnownPublicStreams } from '../../../../src/plugins/storage/KnownPublicStreams' + +const STREAM = '0x1234567890123456789012345678901234567890/public-1' +const OTHER = '0x1234567890123456789012345678901234567890/other-1' + +describe('KnownPublicStreams', () => { + + let dir: string + let file: string + + beforeEach(() => { + jest.useFakeTimers() + dir = mkdtempSync(path.join(os.tmpdir(), 'known-public-')) + file = path.join(dir, 'known-public-streams.json') + }) + + afterEach(() => { + jest.useRealTimers() + rmSync(dir, { recursive: true, force: true }) + }) + + it('keeps the latest answer per stream', () => { + const known = new KnownPublicStreams() + known.set(STREAM, true) + expect(known.has(STREAM)).toBe(true) + expect(known.has(OTHER)).toBe(false) + known.set(STREAM, false) + expect(known.has(STREAM)).toBe(false) + }) + + it('saves a few seconds after a change, and a new instance reads it back', () => { + const known = new KnownPublicStreams(file) + known.set(STREAM, true) + known.set(OTHER, true) + expect(existsSync(file)).toBe(false) + jest.advanceTimersByTime(5 * 1000) + expect(JSON.parse(readFileSync(file, 'utf8'))).toEqual([STREAM, OTHER]) + const restarted = new KnownPublicStreams(file) + expect(restarted.has(STREAM)).toBe(true) + expect(restarted.has(OTHER)).toBe(true) + }) + + it('saves a stream that stopped being public', () => { + const known = new KnownPublicStreams(file) + known.set(STREAM, true) + known.flush() + known.set(STREAM, false) + known.flush() + expect(new KnownPublicStreams(file).has(STREAM)).toBe(false) + }) + + it('writes nothing when no answer changed', () => { + const known = new KnownPublicStreams(file) + known.set(STREAM, false) + known.flush() + expect(existsSync(file)).toBe(false) + }) + + it('starts empty when the file is missing or unreadable', () => { + expect(new KnownPublicStreams(file).has(STREAM)).toBe(false) + writeFileSync(file, '{not json') + expect(new KnownPublicStreams(file).has(STREAM)).toBe(false) + writeFileSync(file, JSON.stringify({ [STREAM]: true })) + expect(new KnownPublicStreams(file).has(STREAM)).toBe(false) + }) +}) diff --git a/packages/node/test/unit/plugins/storage/PomboGates.test.ts b/packages/node/test/unit/plugins/storage/PomboGates.test.ts index f0047d378..96bf806a1 100644 --- a/packages/node/test/unit/plugins/storage/PomboGates.test.ts +++ b/packages/node/test/unit/plugins/storage/PomboGates.test.ts @@ -109,6 +109,81 @@ describe('PomboGates', () => { }) }) + describe('the last answer about a public stream', () => { + const PUBLIC_STREAM = '0x1234567890123456789012345678901234567890/public-1' + + const askChain = async (streamId: string): Promise => { + if (await gates.getGate(streamId) === null) { + await gates.isPublicSubscribe(streamId) + } + } + + beforeEach(() => { + jest.useFakeTimers() + client.hasPermission.mockResolvedValue(true) + }) + + afterEach(() => { + jest.useRealTimers() + }) + + it('remembers a stream with no gate and public SUBSCRIBE', async () => { + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(false) + await askChain(PUBLIC_STREAM) + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(true) + }) + + it('never remembers a private or a gated stream', async () => { + client.hasPermission.mockResolvedValue(false) + await askChain(PUBLIC_STREAM) + await askChain(CONVERSATION_STREAM) + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(false) + expect(gates.wasLastSeenPublic(CONVERSATION_STREAM)).toBe(false) + }) + + it('keeps the last answer while the chain fails', async () => { + await askChain(PUBLIC_STREAM) + jest.advanceTimersByTime(11 * 60 * 1000) + client.getStreamMetadata.mockRejectedValue(new Error('RPC unavailable')) + client.hasPermission.mockRejectedValue(new Error('RPC unavailable')) + await expect(askChain(PUBLIC_STREAM)).rejects.toThrow('RPC unavailable') + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(true) + }) + + it('forgets it as soon as the chain answers that SUBSCRIBE is no longer public', async () => { + await askChain(PUBLIC_STREAM) + jest.advanceTimersByTime(11 * 60 * 1000) + client.hasPermission.mockResolvedValue(false) + await askChain(PUBLIC_STREAM) + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(false) + }) + + it('forgets it as soon as the chain answers that the stream has a gate', async () => { + await askChain(PUBLIC_STREAM) + jest.advanceTimersByTime(21 * 1000) + client.getStreamMetadata.mockResolvedValue(gatedMetadata(GATE)) + await askChain(PUBLIC_STREAM) + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(false) + }) + + it('forgets it when the metadata names a gate, even if the gate itself cannot be read', async () => { + await askChain(PUBLIC_STREAM) + jest.advanceTimersByTime(21 * 1000) + client.getStreamMetadata.mockResolvedValue(gatedMetadata(GATE)) + gateReader.getInfo.mockRejectedValue(new Error('RPC unavailable')) + await expect(askChain(PUBLIC_STREAM)).rejects.toThrow('RPC unavailable') + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(false) + }) + + it('forgets it when the stream is deleted on chain', async () => { + await askChain(PUBLIC_STREAM) + jest.advanceTimersByTime(21 * 1000) + client.getStreamMetadata.mockRejectedValue(Object.assign(new Error('Stream not found'), { code: 'STREAM_NOT_FOUND' })) + await expect(askChain(PUBLIC_STREAM)).rejects.toThrow('Stream not found') + expect(gates.wasLastSeenPublic(PUBLIC_STREAM)).toBe(false) + }) + }) + describe('metadata parsing', () => { it('reads the gate from the Pombo description', () => { expect(parseGateAddress(gatedMetadata(GATE))).toBe(GATE) diff --git a/packages/node/test/unit/plugins/storage/RetentionScheduler.test.ts b/packages/node/test/unit/plugins/storage/RetentionScheduler.test.ts index 82bb8155d..69c1fc73b 100644 --- a/packages/node/test/unit/plugins/storage/RetentionScheduler.test.ts +++ b/packages/node/test/unit/plugins/storage/RetentionScheduler.test.ts @@ -201,9 +201,9 @@ describe('RetentionScheduler first run after a start', () => { rmSync(dir, { recursive: true, force: true }) }) - it('runs a minute after starting when it never ran', () => { + it('runs ten minutes after starting when it never ran', () => { const runOnce = startScheduler() - jest.advanceTimersByTime(59 * 1000) + jest.advanceTimersByTime(10 * 60 * 1000 - 1000) expect(runOnce).not.toHaveBeenCalled() jest.advanceTimersByTime(2 * 1000) expect(runOnce).toHaveBeenCalledTimes(1) @@ -219,9 +219,9 @@ describe('RetentionScheduler first run after a start', () => { }) it('records when a run starts', () => { - const runStart = Date.now() + 60 * 1000 + const runStart = Date.now() + 10 * 60 * 1000 startScheduler() - jest.advanceTimersByTime(61 * 1000) + jest.advanceTimersByTime(10 * 60 * 1000 + 1000) expect(Number(readFileSync(stateFile, 'utf8'))).toBe(runStart) }) }) diff --git a/packages/node/test/unit/plugins/storage/signedReads.test.ts b/packages/node/test/unit/plugins/storage/signedReads.test.ts index f0251e266..f832fdcc7 100644 --- a/packages/node/test/unit/plugins/storage/signedReads.test.ts +++ b/packages/node/test/unit/plugins/storage/signedReads.test.ts @@ -5,6 +5,7 @@ import { BaseWallet, Wallet } from 'ethers' import express from 'express' import { mock } from 'jest-mock-extended' import request from 'supertest' +import { KnownPublicStreams } from '../../../../src/plugins/storage/KnownPublicStreams' import { GateInfo, GateReader, PomboGates } from '../../../../src/plugins/storage/PomboGates' import { SignedRequestVerifier, createSignedRequestMessage } from '../../../../src/plugins/storage/SignedRequest' import { SIGNED_READ_HEADERS, canonicalQuery, createSignedReadGuard } from '../../../../src/plugins/storage/signedReads' @@ -177,6 +178,59 @@ describe('signed reads', () => { }) }) + describe('when the chain cannot answer', () => { + const PUBLIC = '0x1234567890123456789012345678901234567890/public-1' + const RPC_DOWN = new Error('RPC unavailable') + let known: KnownPublicStreams + + // A node restarted with the list on disk: its caches are empty, so every lookup goes to the chain. + const restartNode = () => { + gates.destroy() + gates = new PomboGates(client, gateReader, known) + } + + beforeEach(() => { + known = new KnownPublicStreams() + restartNode() + }) + + it('serves a stream the chain last described as public', async () => { + await read(createApp(true), PUBLIC, { count: '5' }).expect(200) + restartNode() + client.getStreamMetadata.mockRejectedValue(RPC_DOWN) + await read(createApp(true), PUBLIC, { count: '5' }).expect(200) + restartNode() + client.getStreamMetadata.mockResolvedValue({ partitions: 1 }) + client.hasPermission.mockRejectedValue(RPC_DOWN) + await read(createApp(true), PUBLIC, { count: '5' }).expect(200) + }) + + it('refuses a public stream it never asked the chain about', async () => { + client.getStreamMetadata.mockRejectedValue(RPC_DOWN) + await read(createApp(true), PUBLIC, { count: '5' }).expect(503) + }) + + it('refuses a stream the chain last described as no longer public', async () => { + await read(createApp(true), PUBLIC, { count: '5' }).expect(200) + restartNode() + client.hasPermission.mockResolvedValue(false) + await read(createApp(true), PUBLIC, { count: '5' }).expect(401) + restartNode() + client.hasPermission.mockRejectedValue(RPC_DOWN) + await read(createApp(true), PUBLIC, { count: '5' }).expect(503) + }) + + it('refuses a stream the chain last described as gated', async () => { + await read(createApp(true), PUBLIC, { count: '5' }).expect(200) + restartNode() + client.getStreamMetadata.mockResolvedValue({ partitions: 1, description: JSON.stringify({ a: 'pombo', t: 'gated', g: GATE }) }) + await read(createApp(true), PUBLIC, { count: '5' }).expect(401) + restartNode() + client.getStreamMetadata.mockRejectedValue(RPC_DOWN) + await read(createApp(true), PUBLIC, { count: '5' }).expect(503) + }) + }) + it('canonicalises the query string by parameter name', () => { expect(canonicalQuery({ toTimestamp: '2', fromTimestamp: '1', format: 'raw' })).toBe('format=raw&fromTimestamp=1&toTimestamp=2') expect(canonicalQuery({})).toBe('') diff --git a/packages/node/test/unit/plugins/storage/storedEndpoint.test.ts b/packages/node/test/unit/plugins/storage/storedEndpoint.test.ts index 830037880..a3b20c4e7 100644 --- a/packages/node/test/unit/plugins/storage/storedEndpoint.test.ts +++ b/packages/node/test/unit/plugins/storage/storedEndpoint.test.ts @@ -13,6 +13,7 @@ import { BaseWallet, Wallet } from 'ethers' import express from 'express' import { mock } from 'jest-mock-extended' import request from 'supertest' +import { KnownPublicStreams } from '../../../../src/plugins/storage/KnownPublicStreams' import { GateInfo, GateReader, PomboGates } from '../../../../src/plugins/storage/PomboGates' import { SignedRequestVerifier, createSignedRequestMessage } from '../../../../src/plugins/storage/SignedRequest' import { Storage } from '../../../../src/plugins/storage/Storage' @@ -70,6 +71,12 @@ describe('storedEndpoint', () => { return { user, issuedAt, nonce, signature: await signingWallet.signMessage(message), targets } } + const mountEndpoint = () => { + app = express() + const endpoint = createStoredEndpoint(storage, gates, client, new SignedRequestVerifier()) + app.route(endpoint.path)[endpoint.method](endpoint.requestHandlers) + } + const stored = (body: any) => { return request(app).post(`/streams/${encodeURIComponent(STREAM_ID)}/data/partitions/${PARTITION}/stored`).send(body) } @@ -89,9 +96,7 @@ describe('storedEndpoint', () => { client.getMessageSigner.mockReturnValue(randomEthereumAddress()) // not the signer by default gateReader.getInfo.mockResolvedValue(gate()) gates = new PomboGates(client, gateReader) - app = express() - const endpoint = createStoredEndpoint(storage, gates, client, new SignedRequestVerifier()) - app.route(endpoint.path)[endpoint.method](endpoint.requestHandlers) + mountEndpoint() }) afterEach(() => { @@ -148,4 +153,15 @@ describe('storedEndpoint', () => { const body = await signedBody([{ timestamp: 1000, sequenceNumber: 0 }]) await stored(body).expect(503) }) + + it('answers for a stream the chain last described as public while it cannot be consulted', async () => { + const known = new KnownPublicStreams() + known.set(STREAM_ID, true) + gates.destroy() + gates = new PomboGates(client, gateReader, known) + mountEndpoint() + client.getStreamMetadata.mockRejectedValue(new Error('RPC unavailable')) + const res = await stored(await signedBody([{ timestamp: 1000, sequenceNumber: 0 }])).expect(200) + expect(res.body.results[0].result).toBe('present') + }) })