Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion HOW_TO_INSTALL.md
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@ docker compose cp cassandra/<file>.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
Expand Down
26 changes: 18 additions & 8 deletions POMBO.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand All @@ -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
Expand Down Expand Up @@ -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 |
Expand Down
3 changes: 3 additions & 0 deletions deploy/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
Expand All @@ -72,3 +74,4 @@ services:

volumes:
cassandra-data:
node-state:
87 changes: 87 additions & 0 deletions packages/node/src/plugins/storage/KnownPublicStreams.ts
Original file line number Diff line number Diff line change
@@ -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<string>
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 })
}
}
}
36 changes: 32 additions & 4 deletions packages/node/src/plugins/storage/PomboGates.ts
Original file line number Diff line number Diff line change
@@ -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)',
Expand Down Expand Up @@ -122,25 +123,42 @@ export class PomboGates {
private readonly accessCache = new MapWithTtl<string, boolean>(ttlOf)
private readonly subscribePublicCache = new MapWithTtl<string, boolean>(ttlOf)
private readonly subscribeUserCache = new MapWithTtl<string, boolean>(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<GateInfo | null> {
const cached = this.gateCache.get(streamId)
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)
if (linked !== undefined && linked !== streamId) {
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
Expand All @@ -154,17 +172,26 @@ 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<boolean> {
const cached = this.subscribePublicCache.get(streamId)
if (cached !== undefined) {
return cached
}
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<boolean> {
const key = `${streamId}_${user}`
Expand Down Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion packages/node/src/plugins/storage/RetentionScheduler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')

Expand Down
9 changes: 7 additions & 2 deletions packages/node/src/plugins/storage/StoragePlugin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -99,7 +100,11 @@ export class StoragePlugin extends Plugin<StoragePluginConfig> {
})
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) => {
Expand Down
2 changes: 1 addition & 1 deletion packages/node/src/plugins/storage/config.schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
15 changes: 13 additions & 2 deletions packages/node/src/plugins/storage/signedReads.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
10 changes: 7 additions & 3 deletions packages/node/src/plugins/storage/storedEndpoint.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading
Loading