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
3 changes: 2 additions & 1 deletion Dockerfile.node
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@ RUN apt-get update && apt-get --assume-yes --no-install-recommends install \
RUN usermod -d /home/streamr -l streamr node && groupmod -n streamr node
# The node writes its GeoIP database under ~/.streamr; when only ~/.streamr/config
# is bind-mounted, Docker would create the parent as root and the write fails.
RUN mkdir -p /home/streamr/.streamr && chown streamr:streamr /home/streamr/.streamr
# The home itself must be writable too: on arm64 a dependency creates ~/.local.
RUN mkdir -p /home/streamr/.streamr && chown streamr:streamr /home/streamr /home/streamr/.streamr
USER streamr
WORKDIR /home/streamr/network
COPY --chown=root:root --from=build /usr/src/network/ .
Expand Down
3 changes: 1 addition & 2 deletions deploy/config/pombo-node.json.example
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,7 @@
"contracts": {
"rpcs": [
{ "url": "https://polygon.drpc.org" },
{ "url": "https://polygon-bor-rpc.publicnode.com" },
{ "url": "https://rpc.ankr.com/polygon" }
{ "url": "https://polygon-bor-rpc.publicnode.com" }
],
"rpcQuorum": 1
},
Expand Down
2 changes: 1 addition & 1 deletion deploy/install.sh
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
set -euo pipefail
cd "$(dirname "$0")"

RPCS='[ { "url": "https://polygon.drpc.org" }, { "url": "https://polygon-bor-rpc.publicnode.com" }, { "url": "https://rpc.ankr.com/polygon" } ]'
RPCS='[ { "url": "https://polygon.drpc.org" }, { "url": "https://polygon-bor-rpc.publicnode.com" } ]'
RPC0="https://polygon.drpc.org"
MIN_WEI="20000000000000000" # 0.02 POL: enough for the assignment stream + registration on Polygon
CONFIG_IN_CONTAINER="/home/streamr/.streamr/config/pombo-node.json"
Expand Down
5 changes: 4 additions & 1 deletion packages/node/src/httpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,10 @@ export const startServer = async (
const app = express()
app.use(cors({
origin: true, // Access-Control-Allow-Origin: request origin. The default '*' is invalid if credentials included.
credentials: true // Access-Control-Allow-Credentials: true
credentials: true, // Access-Control-Allow-Credentials: true
// Signed reads carry a custom header, so every one of them is preceded
// by a preflight. Without this the browser default is a few seconds.
maxAge: 600
}))
endpoints.forEach((endpoint: Endpoint) => {
const handlers = [createAuthenticatorMiddleware(endpoint.apiAuthentication)].concat(endpoint.requestHandlers)
Expand Down
15 changes: 14 additions & 1 deletion packages/node/src/plugins/storage/IngestValidator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ export class IngestValidator {
await this.client.validateMessage(msg)
} catch (err: any) {
if (DEFINITIVE_REJECTIONS.has(err?.code)) {
return this.reject(msg, err.code)
return this.reject(msg, err.code, this.signerOf(msg))
}
logger.warn('Could not validate message, storing it', {
streamId: msg.getStreamId(),
Expand Down Expand Up @@ -91,6 +91,19 @@ export class IngestValidator {
return this.reject(msg, 'READ_ONLY', signer)
}

/**
* Who signed, for the log line. In a gated channel the publisher is the
* gate contract, the same for every member, so without this a rejection
* does not say whose message was dropped.
*/
private signerOf(msg: StreamMessage): EthereumAddress | undefined {
try {
return this.client.getMessageSigner(msg)
} catch {
return undefined
}
}

private reject(msg: StreamMessage, reason: string, signer?: EthereumAddress): IngestVerdict {
this.metrics.rejectedMessagesPerSecond.record(1)
logger.info('Rejected message at ingest', {
Expand Down
1 change: 1 addition & 0 deletions packages/sdk/src/StreamrClientError.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ export type StreamrClientErrorCode =
'INVALID_MESSAGE_CONTENT' |
'INVALID_STREAM_METADATA' |
'INVALID_SIGNATURE' |
'CHAIN_UNAVAILABLE' |
'INVALID_PARTITION' |
'DECRYPT_ERROR' |
'STORAGE_NODE_ERROR' |
Expand Down
25 changes: 21 additions & 4 deletions packages/sdk/src/contracts/ERC1271ContractFacade.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { BrandedString, EthereumAddress, EcdsaSecp256k1Evm, MapWithTtl, toUserId, UserID } from '@streamr/utils'
import { Lifecycle, scoped } from 'tsyringe'
import { RpcProviderSource } from '../RpcProviderSource'
import { StreamrClientError } from '../StreamrClientError'
import type { IERC1271 as ERC1271Contract } from '../ethereumArtifacts/IERC1271'
import ERC1271ContractArtifact from '../ethereumArtifacts/IERC1271Abi.json'
import { createLazyMap, Mapping } from '../utils/Mapping'
Expand All @@ -10,7 +11,13 @@ export const SUCCESS_MAGIC_VALUE = '0x1626ba7e' // Magic value for success as de

export type CacheKey = BrandedString<string>

const CACHE_TTL = 10 * 60 * 1000 // 10 minutes
const CACHE_TTL = 10 * 60 * 1000
/**
* A refusal is remembered for seconds. An account publishes the moment it pays
* its way past a gate, and a held "no" drops everything it writes until the
* entry expires.
*/
const DENIAL_TTL = 20 * 1000 // 10 minutes

const signingUtil = new EcdsaSecp256k1Evm()

Expand All @@ -22,7 +29,8 @@ function formCacheKey(contractAddress: EthereumAddress, signerUserId: UserID): C
export class ERC1271ContractFacade {

private readonly contractsByAddress: Mapping<EthereumAddress, ERC1271Contract>
private readonly publisherCache = new MapWithTtl<CacheKey, boolean>(() => CACHE_TTL)
private readonly publisherCache = new MapWithTtl<CacheKey, boolean>(
(isValid) => (isValid ? CACHE_TTL : DENIAL_TTL))

constructor(
contractFactory: ContractFactory,
Expand All @@ -47,8 +55,17 @@ export class ERC1271ContractFacade {
if (cachedValue !== undefined) {
return cachedValue
} else {
const contract = await this.contractsByAddress.get(contractAddress)
const result = await contract.isValidSignature(signingUtil.keccakHash(payload), signature)
let result: string
try {
const contract = await this.contractsByAddress.get(contractAddress)
result = await contract.isValidSignature(signingUtil.keccakHash(payload), signature)
} catch (err) {
// Not an answer about the signature: the caller decides what an
// unreachable chain means, and must not read it as a refusal.
const reason = (err instanceof Error) ? err.message : String(err)
throw new StreamrClientError(
`Could not ask ${contractAddress} whether the signature is valid: ${reason}`, 'CHAIN_UNAVAILABLE')
}
const isValid = result === SUCCESS_MAGIC_VALUE
this.publisherCache.set(cacheKey, isValid)
return isValid
Expand Down
6 changes: 6 additions & 0 deletions packages/sdk/src/signature/SignatureValidator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,12 @@ export class SignatureValidator {
try {
success = await this.validate(streamMessage)
} catch (err) {
// A failure to reach the chain is not a verdict on the signature,
// and callers that drop invalid messages must be able to tell the
// two apart.
if (err instanceof StreamrClientError) {
throw err
}
// eslint-disable-next-line @typescript-eslint/restrict-template-expressions
throw new StreamrClientError(`An error occurred during address recovery from signature: ${err}`, 'INVALID_SIGNATURE', streamMessage)
}
Expand Down
34 changes: 34 additions & 0 deletions packages/sdk/test/unit/ERC1271ContractFacade.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,40 @@ describe('ERC1271ContractFacade', () => {
expect(contractOne.isValidSignature).toHaveBeenCalledTimes(1)
})

it('isValidSignature: an unreachable chain is not a verdict, and is not cached', async () => {
contractOne.isValidSignature.mockRejectedValue(new Error('RPC is down'))
const err = await contractFacade.isValidSignature(CONTRACT_ADDRESS_ONE, PAYLOAD, signature).catch((e) => e)
expect(err.code).toEqual('CHAIN_UNAVAILABLE')
await contractFacade.isValidSignature(CONTRACT_ADDRESS_ONE, PAYLOAD, signature).catch(() => {})
expect(contractOne.isValidSignature).toHaveBeenCalledTimes(2)
})

it('isValidSignature: an invalid result is re-checked within seconds', async () => {
jest.useFakeTimers()
try {
contractOne.isValidSignature.mockResolvedValue('0xaaaaaaaa')
await contractFacade.isValidSignature(CONTRACT_ADDRESS_ONE, PAYLOAD, signature)
jest.advanceTimersByTime(30 * 1000)
await contractFacade.isValidSignature(CONTRACT_ADDRESS_ONE, PAYLOAD, signature)
expect(contractOne.isValidSignature).toHaveBeenCalledTimes(2)
} finally {
jest.useRealTimers()
}
})

it('isValidSignature: a valid result outlives that', async () => {
jest.useFakeTimers()
try {
contractOne.isValidSignature.mockResolvedValue(SUCCESS_MAGIC_VALUE)
await contractFacade.isValidSignature(CONTRACT_ADDRESS_ONE, PAYLOAD, signature)
jest.advanceTimersByTime(30 * 1000)
await contractFacade.isValidSignature(CONTRACT_ADDRESS_ONE, PAYLOAD, signature)
expect(contractOne.isValidSignature).toHaveBeenCalledTimes(1)
} finally {
jest.useRealTimers()
}
})

it('differentiates between different contracts based on contract address', async () => {
contractOne.isValidSignature.mockResolvedValue(SUCCESS_MAGIC_VALUE)
contractTwo.isValidSignature.mockResolvedValue('0xaaaaaaaa')
Expand Down
Loading