mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
Merge pull request #3469 from vitorpamplona/claude/relay-performance-geode-quartz-v6lbys
feat(quartz): NostrServer.ingest — local write path with per-submission verify skip
This commit is contained in:
@@ -33,12 +33,12 @@ listing (including archived plans), open that folder's `README.md`.
|
||||
| amethyst | 21 | 19 | 1 | 1 | 0 | [amethyst/plans](amethyst/plans/README.md) |
|
||||
| nestsClient | 26 | 23 | 1 | 2 | 0 | [nestsClient/plans](nestsClient/plans/README.md) |
|
||||
| desktopApp | 13 | 10 | 2 | 1 | 0 | [desktopApp/plans](desktopApp/plans/README.md) |
|
||||
| quartz | 9 | 7 | 0 | 2 | 0 | [quartz/plans](quartz/plans/README.md) |
|
||||
| quartz | 10 | 7 | 0 | 3 | 0 | [quartz/plans](quartz/plans/README.md) |
|
||||
| commons | 6 | 2 | 2 | 2 | 0 | [commons/plans](commons/plans/README.md) |
|
||||
| cli | 6 | 5 | 1 | 0 | 0 | [cli/plans](cli/plans/README.md) |
|
||||
| quic | 4 | 3 | 0 | 0 | 1 | [quic/plans](quic/plans/README.md) |
|
||||
| quic/interop | 1 | 1 | 0 | 0 | 0 | [quic/interop/plans](quic/interop/plans/README.md) |
|
||||
| geode | 4 | 4 | 0 | 0 | 0 | [geode/plans](geode/plans/README.md) |
|
||||
| geode | 5 | 4 | 0 | 0 | 0 | [geode/plans](geode/plans/README.md) |
|
||||
| docs (frozen) | 52 | 48 | 2 | 0 | 2 | [docs/plans](docs/plans/README.md) |
|
||||
|
||||
## Live work (not shipped)
|
||||
@@ -67,6 +67,7 @@ them under each folder's `archive/` via the per-module index above.
|
||||
| amethyst | [napplet-inter-applet](amethyst/plans/2026-06-20-napplet-inter-applet.md) | NAP-INC / NAP-INTENT inter-applet messaging; prerequisites (multi-applet hosting, archetype registry, `MESSAGING` capability) not built. |
|
||||
| quartz | [local-headers-explorer](quartz/plans/2026-05-08-local-headers-explorer.md) | Headers-only Bitcoin P2P client to verify NIP-03 OTS attestations without a trusted block explorer. |
|
||||
| quartz | [giftwrap-deletion-requests](quartz/plans/2026-06-12-giftwrap-deletion-requests.md) | Let a recipient-authored kind-5 delete/block a gift wrap (kind 1059) addressed to them. |
|
||||
| quartz | [incremental-negentropy-storage](quartz/plans/2026-07-03-incremental-negentropy-storage.md) | Always-current (created_at, id) index so cold NEG-OPENs stop paying a full scan + seal. |
|
||||
| commons | [event-renderer](commons/plans/2026-04-21-event-renderer.md) | Cross-platform UI-agnostic `RenderedEvent` subsystem shared by Amy, Desktop, Android; not started. |
|
||||
| commons | [amethyst-to-commons-migration](commons/plans/2026-05-30-amethyst-to-commons-migration.md) | Roadmap to move shared `amethyst` Android code into `commons`; keystone `Account`/`LocalCache` extraction not begun. |
|
||||
| desktopApp | [embedded-wallet-phase2-research](desktopApp/plans/2026-05-21-embedded-wallet-phase2-research.md) | Research for an embedded self-custodial Lightning wallet (Breez/ldk-node/lightning-kmp); parked, no code. |
|
||||
|
||||
+7
@@ -20,6 +20,7 @@
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import okhttp3.Dispatcher
|
||||
import okhttp3.OkHttpClient
|
||||
@@ -57,6 +58,12 @@ class OkHttpClientFactoryForRelays(
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
.dispatcher(myDispatcher)
|
||||
// TCP_NODELAY: a CLOSE (never answered by relays) followed by a
|
||||
// REQ — every feed switch — otherwise nagles the REQ behind the
|
||||
// unACKed CLOSE for the peer's delayed-ACK window (~40 ms+).
|
||||
// See quartz TcpNoDelaySocketFactory. Direct connections only;
|
||||
// the Tor SOCKS path is unaffected.
|
||||
.socketFactory(TcpNoDelaySocketFactory)
|
||||
.dns(dns)
|
||||
.eventListenerFactory(DnsInvalidatingEventListener.Factory(dns))
|
||||
.followRedirects(true)
|
||||
|
||||
@@ -50,6 +50,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
@@ -111,7 +112,7 @@ class Context(
|
||||
val identity: Identity,
|
||||
val state: RunState,
|
||||
) : AutoCloseable {
|
||||
private val okhttp = OkHttpClient.Builder().build()
|
||||
private val okhttp = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
|
||||
|
||||
val client: NostrClient =
|
||||
NostrClient(
|
||||
|
||||
@@ -31,6 +31,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
|
||||
import kotlinx.coroutines.selects.select
|
||||
@@ -143,7 +144,7 @@ object NipCommand {
|
||||
timeoutMs: Long,
|
||||
): List<Map<String, Any?>> {
|
||||
if (SEARCH_RELAYS.isEmpty()) return emptyList()
|
||||
val okhttp = OkHttpClient.Builder().build()
|
||||
val okhttp = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
|
||||
val client = NostrClient(websocketBuilder = BasicOkHttpWebSocket.Builder { okhttp })
|
||||
val filter = Filter(kinds = listOf(NIPTEXT_KIND, WIKI_KIND, LONGFORM_KIND), search = "NIP-$slug", limit = 10)
|
||||
val subId = newSubId()
|
||||
|
||||
@@ -35,6 +35,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
|
||||
import com.vitorpamplona.quartz.nip46RemoteSigner.BunkerResponse
|
||||
import com.vitorpamplona.quartz.nip46RemoteSigner.NostrConnectEvent
|
||||
@@ -128,7 +129,7 @@ object NostrConnect {
|
||||
System.err.println("[nostrconnect] paste this into your signer within ${timeoutMs / 1000}s:")
|
||||
System.err.println(offer)
|
||||
|
||||
val okhttp = OkHttpClient.Builder().build()
|
||||
val okhttp = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
|
||||
val client = NostrClient(websocketBuilder = BasicOkHttpWebSocket.Builder { okhttp })
|
||||
val incoming = Channel<NostrConnectEvent>(UNLIMITED)
|
||||
val subId = newSubId()
|
||||
|
||||
+10
@@ -25,6 +25,7 @@ import com.vitorpamplona.amethyst.commons.tor.TorType
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isLocalHost
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isOnion
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
@@ -66,6 +67,9 @@ class DesktopHttpClient(
|
||||
private val directClient: OkHttpClient by lazy {
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
// TCP_NODELAY: keeps CLOSE-then-REQ bursts (feed switches) from
|
||||
// nagling behind the unACKed CLOSE — see TcpNoDelaySocketFactory.
|
||||
.socketFactory(TcpNoDelaySocketFactory)
|
||||
.connectionPool(sharedConnectionPool)
|
||||
.connectTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
|
||||
.readTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
|
||||
@@ -150,6 +154,12 @@ class DesktopHttpClient(
|
||||
private val simpleClient: OkHttpClient by lazy {
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
// Direct relay sockets opened before setInstance() should
|
||||
// get the same TCP_NODELAY as directClient — see
|
||||
// TcpNoDelaySocketFactory. (failClosedClient is SOCKS, and
|
||||
// OkHttp bypasses the socket factory for SOCKS proxies, so
|
||||
// it doesn't need this.)
|
||||
.socketFactory(TcpNoDelaySocketFactory)
|
||||
.connectTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
|
||||
.readTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
|
||||
.writeTimeout(BASE_TIMEOUT_SECONDS, TimeUnit.SECONDS)
|
||||
|
||||
@@ -85,6 +85,11 @@ dependencies {
|
||||
implementation(libs.jackson.module.kotlin)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
|
||||
// Outbound WebSockets for the [[mirror]] upstream streams (quartz's
|
||||
// BasicOkHttpWebSocket transport). Same OkHttp the rest of the repo
|
||||
// already ships (Apache-2.0).
|
||||
implementation(libs.okhttp)
|
||||
|
||||
// Bundled SQLite driver — Relay's default in-memory EventStore creates
|
||||
// an in-memory DB at runtime.
|
||||
implementation(libs.androidx.sqlite.bundled.jvm)
|
||||
|
||||
@@ -46,6 +46,27 @@ path = "/"
|
||||
in_memory = false
|
||||
file = "/var/lib/geode/events.db"
|
||||
|
||||
# NOTE: the four knobs below measured as pure noise in relayBench A/B
|
||||
# runs on 4-core container hardware at 50k events (see
|
||||
# geode/plans/2026-07-04-sqlite-knobs-ab.md) — their value is
|
||||
# hardware-dependent, so benchmark on YOUR box before enabling.
|
||||
#
|
||||
# Reader-connection pool size (file-backed stores only). Default 4.
|
||||
# readers = 4
|
||||
|
||||
# PRAGMA mmap_size in bytes — maps the db file into memory so reads
|
||||
# skip the pread syscall. Off by default (SQLite default).
|
||||
# mmap_size = 268435456
|
||||
|
||||
# PRAGMA temp_store = MEMORY — RAM instead of temp files for large
|
||||
# sorts. Off by default.
|
||||
# temp_store_memory = true
|
||||
|
||||
# Refresh query-planner statistics (PRAGMA optimize) every N seconds.
|
||||
# Incremental and usually a no-op; keeps the planner from drifting
|
||||
# onto the wrong index as the corpus grows. Off by default.
|
||||
# optimize_interval_seconds = 3600
|
||||
|
||||
[options]
|
||||
# Drop events whose Schnorr signature does not verify. Strongly
|
||||
# recommended for any relay accepting traffic from real clients.
|
||||
@@ -86,6 +107,42 @@ require_auth = false
|
||||
# kind_whitelist = [0, 1, 3, 7, 1059, 30023]
|
||||
# kind_blacklist = [4]
|
||||
|
||||
# Mirror upstream relays (strfry-router style, "down" direction): the
|
||||
# relay dials each [[mirror]] url, subscribes to everything newer than
|
||||
# now - backfill_seconds, and ingests the stream alongside client
|
||||
# publishes. Reconnects and re-subscribes automatically.
|
||||
#
|
||||
# `trusted = true` is the relay-to-relay trust switch: events from that
|
||||
# upstream skip Schnorr signature verification (the upstream already
|
||||
# verified its own ingest; re-verifying burns ~8% of ingest CPU). The
|
||||
# trusted identity is the URL *this* relay dialed — TLS-authenticated
|
||||
# for wss:// — so an inbound client can never claim it. Default false:
|
||||
# mirror-but-verify. Only trust relays you operate or whose ingest
|
||||
# discipline you'd stake your own db on.
|
||||
#
|
||||
# `filter` (optional) scopes an upstream, as a NIP-01 filter JSON
|
||||
# object — same idea as strfry-router's per-stream filter. It shapes
|
||||
# the REQ sent upstream AND every delivered event is re-checked
|
||||
# against it before ingest, so even a trusted upstream can only
|
||||
# inject events inside the declared scope. `since`/`limit` inside it
|
||||
# are ignored (backfill_seconds owns the time window). Omit to mirror
|
||||
# everything; for several disjoint scopes, repeat [[mirror]] with the
|
||||
# same url. The keys are validated at boot (a typo like `kindss` or a
|
||||
# scalar where an array belongs fails startup) precisely because this
|
||||
# filter is the trust boundary for `trusted = true`.
|
||||
#
|
||||
# `dir` (strfry-router parity) sets the flow direction: "down" pulls
|
||||
# from the upstream (default), "up" pushes this relay's matching
|
||||
# events to it, "both" does both with echo suppression so the two
|
||||
# directions don't ping-pong the same event.
|
||||
#
|
||||
# [[mirror]]
|
||||
# url = "wss://upstream.example.com/"
|
||||
# dir = "both"
|
||||
# trusted = true
|
||||
# backfill_seconds = 3600
|
||||
# filter = '{"kinds":[0,1,3,7],"#t":["nostr"]}'
|
||||
|
||||
[admin]
|
||||
# NIP-86 relay management API. When `pubkeys` is non-empty, the relay
|
||||
# accepts HTTP POST application/nostr+json+rpc on the same URL,
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
# Server-side SQLite knobs — A/B verdict: no winner on container hardware
|
||||
|
||||
**Status: closed (negative result recorded).** Backlog items 4–5 of the
|
||||
relay performance campaign.
|
||||
|
||||
## What was tested
|
||||
|
||||
The knobs the campaign had previously reverted as noise-inconclusive,
|
||||
re-tested under the A/B protocol — both variants in ONE relayBench run
|
||||
(identical container conditions), then the run repeated with the relay
|
||||
order reversed to expose order bias. 50k synthetic corpus, all four
|
||||
knobs on the candidate at once:
|
||||
|
||||
```toml
|
||||
[database]
|
||||
readers = 8 # reader pool (default 4)
|
||||
mmap_size = 268435456 # 256 MiB
|
||||
temp_store_memory = true
|
||||
optimize_interval_seconds = 15 # PRAGMA optimize between ingest and queries
|
||||
```
|
||||
|
||||
## Result
|
||||
|
||||
Every observed delta flipped with run order. In BOTH runs the
|
||||
second-running relay won most query scenarios and ingest latency —
|
||||
regardless of which variant it was:
|
||||
|
||||
- run 1 (plain → knobs): knobs won 7/10 query scenarios, better ingest
|
||||
p99 (8.2 vs 16.0 ms).
|
||||
- run 2 (knobs → plain): plain won 6/10 query scenarios, better ingest
|
||||
p50/p99/throughput (8,552 vs 8,015 ev/s).
|
||||
|
||||
Storage size identical (±1 MiB). The periodic `PRAGMA optimize` also
|
||||
did not move the author-archive/planner-sensitive scenarios.
|
||||
|
||||
## Verdict
|
||||
|
||||
**All knobs stay off by default.** The `[database]` config plumbing
|
||||
ships anyway (readers / mmap_size / temp_store_memory /
|
||||
optimize_interval_seconds) because their value is hardware-dependent —
|
||||
an operator on NVMe with a large page cache or a memory-constrained VPS
|
||||
should measure on their own box — and `optimize_interval_seconds`
|
||||
remains sensible operationally for long-running relays even without a
|
||||
measurable win at 50k-fresh-database scale.
|
||||
|
||||
**Do not re-benchmark these on container-class 4-core hardware** without
|
||||
first fixing the run-order bias: the protocol note in
|
||||
`quartz/plans/2026-07-03-incremental-negentropy-storage.md` applies —
|
||||
same-run A/B plus order reversal is the minimum, and anything that
|
||||
doesn't survive the order flip is noise.
|
||||
@@ -1,6 +1,6 @@
|
||||
# geode plans
|
||||
|
||||
_Audited 2026-06-30. 4 plans: 4 shipped (archived), 0 in-progress, 0 queued, 0 abandoned._
|
||||
_Audited 2026-06-30. 5 plans: 4 shipped (archived), 0 in-progress, 0 queued, 1 closed (negative result)._
|
||||
|
||||
Performance-focused design docs for future work. Each file is a
|
||||
self-contained sketch — problem statement, observed numbers, proposed
|
||||
@@ -26,3 +26,9 @@ regressions show up in the regular CI matrix once they're enabled.
|
||||
- [archive/2026-05-07-live-broadcast-fanout-index.md](archive/2026-05-07-live-broadcast-fanout-index.md)
|
||||
- [archive/2026-05-07-connection-scaling.md](archive/2026-05-07-connection-scaling.md)
|
||||
- [archive/2026-05-07-negentropy-large-corpus.md](archive/2026-05-07-negentropy-large-corpus.md)
|
||||
|
||||
## Closed (negative result)
|
||||
|
||||
| Plan | Verdict |
|
||||
| ---- | ------- |
|
||||
| [2026-07-04-sqlite-knobs-ab.md](2026-07-04-sqlite-knobs-ab.md) | readers/mmap/temp_store/PRAGMA-optimize knobs: no winner on container hardware; every delta flipped with run order. Knobs ship config-gated, off by default. |
|
||||
|
||||
@@ -21,10 +21,17 @@
|
||||
package com.vitorpamplona.geode
|
||||
|
||||
import com.vitorpamplona.geode.config.BannedEntry
|
||||
import com.vitorpamplona.geode.config.MirrorFilterValidator
|
||||
import com.vitorpamplona.geode.config.RuntimeConfig
|
||||
import com.vitorpamplona.geode.config.RuntimeConfigData
|
||||
import com.vitorpamplona.geode.config.StaticConfig
|
||||
import com.vitorpamplona.geode.mirror.MirrorDirection
|
||||
import com.vitorpamplona.geode.mirror.MirrorUpstream
|
||||
import com.vitorpamplona.geode.mirror.MirrorWorker
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.displayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.FullAuthPolicy
|
||||
@@ -33,9 +40,15 @@ import com.vitorpamplona.quartz.nip01Core.relay.server.policies.OptionalAuthPoli
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.RejectFutureEventsPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyAuthOnlyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyPolicy
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import java.io.File
|
||||
|
||||
/**
|
||||
@@ -85,6 +98,7 @@ fun main(args: Array<String>) {
|
||||
.opt("--config")
|
||||
?.let { StaticConfig.fromFile(File(it)) }
|
||||
?: StaticConfig()
|
||||
config.validate()
|
||||
|
||||
val host = a.opt("--host") ?: config.network.host
|
||||
val port = a.opt("--port")?.toInt() ?: config.network.port
|
||||
@@ -122,7 +136,21 @@ fun main(args: Array<String>) {
|
||||
cliInfoFile?.let { RelayInfo.fromFile(it) }
|
||||
?: config.resolveInfo(fullTextSearch)
|
||||
|
||||
val store: IEventStore = EventStore(dbName = dbFile, relay = advertisedUrl, indexStrategy = relayIndexingStrategy(fullTextSearch))
|
||||
// Deployment tuning from `[database]` — off by default; the quartz
|
||||
// library defaults stay tuned for the app-side stores.
|
||||
val extraPragmas =
|
||||
buildList {
|
||||
config.database.mmap_size?.let { add("PRAGMA mmap_size = $it;") }
|
||||
if (config.database.temp_store_memory) add("PRAGMA temp_store = MEMORY;")
|
||||
}
|
||||
val store =
|
||||
EventStore(
|
||||
dbName = dbFile,
|
||||
relay = advertisedUrl,
|
||||
indexStrategy = relayIndexingStrategy(fullTextSearch, config.negentropy.live_index),
|
||||
numReaders = config.database.readers ?: 4,
|
||||
extraPragmas = extraPragmas,
|
||||
)
|
||||
|
||||
val policyBuilder: () -> IRelayPolicy = {
|
||||
composePolicy(config, advertisedUrl, requireAuth, optionalAuth, verifySigs, parallelVerify)
|
||||
@@ -172,10 +200,99 @@ fun main(args: Array<String>) {
|
||||
callGroupSize = config.network.call_group_size,
|
||||
).start()
|
||||
|
||||
// `[[mirror]]` upstreams: dial each configured relay and stream its
|
||||
// events into the local store. `trusted = true` entries skip Schnorr
|
||||
// verification for that connection (relay-to-relay trust) — only
|
||||
// meaningful while verify is on; with --no-verify nothing verifies
|
||||
// anyway. Never mirror ourselves: a self-URL would echo every local
|
||||
// publish back forever.
|
||||
val upstreams =
|
||||
config.mirror.map { m ->
|
||||
// The optional scope filter is parsed eagerly so a malformed
|
||||
// JSON object fails the boot, not the first delivery. Its
|
||||
// since/limit are stripped: the mirror owns the time window
|
||||
// (backfill_seconds) and never bounds the subscription.
|
||||
val scope =
|
||||
m.filter?.let { json ->
|
||||
// Strict-validate FIRST: the deserializer is tolerant
|
||||
// (unknown keys skipped, wrong-typed entries dropped),
|
||||
// so a typo like `{"kindss":[4]}` would silently parse
|
||||
// to an empty filter — widening a trusted upstream's
|
||||
// scope to the whole firehose. This filter is the
|
||||
// containment boundary for `trusted = true`, so a typo
|
||||
// must fail the boot, not the boundary.
|
||||
MirrorFilterValidator.validate(m.url, json)
|
||||
val parsed =
|
||||
try {
|
||||
OptimizedJsonMapper.fromJsonTo<Filter>(json)
|
||||
} catch (e: Exception) {
|
||||
throw IllegalArgumentException(
|
||||
"[[mirror]] filter for ${m.url} is not a valid NIP-01 filter object: $json",
|
||||
e,
|
||||
)
|
||||
}
|
||||
parsed.copy(since = null, limit = null)
|
||||
}
|
||||
val direction =
|
||||
MirrorDirection.parse(m.dir)
|
||||
?: throw IllegalArgumentException(
|
||||
"[[mirror]] dir for ${m.url} must be \"down\", \"up\" or \"both\" (got \"${m.dir}\")",
|
||||
)
|
||||
MirrorUpstream(
|
||||
url = m.url.normalizeRelayUrl(),
|
||||
trusted = m.trusted,
|
||||
backfillSeconds = m.backfill_seconds,
|
||||
filter = scope,
|
||||
direction = direction,
|
||||
)
|
||||
}
|
||||
// Never mirror ourselves — a self-URL echoes every local publish
|
||||
// back forever. Compare scheme-insensitively (ws:// vs wss:// for the
|
||||
// same host is still us) and ignoring the trailing slash. This can't
|
||||
// catch a public URL that resolves to this bind behind a proxy, nor a
|
||||
// `--port 0` autobind, so it's a guardrail against the obvious typo,
|
||||
// not a proof of non-self-reference.
|
||||
val advertisedIdentity = advertisedUrl.displayUrl()
|
||||
require(upstreams.none { it.url.displayUrl() == advertisedIdentity }) {
|
||||
"[[mirror]] must not list this relay's own URL ($advertisedUrl)"
|
||||
}
|
||||
val mirror =
|
||||
if (upstreams.isEmpty()) {
|
||||
null
|
||||
} else {
|
||||
MirrorWorker(upstreams, relay.server).also { it.start() }
|
||||
}
|
||||
|
||||
// Periodic query-planner statistics refresh (`PRAGMA optimize`).
|
||||
// Incremental and usually a no-op; failures are swallowed — a missed
|
||||
// refresh only means slightly staler planner stats until the next tick.
|
||||
val maintenanceScope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||
config.database.optimize_interval_seconds?.let { secs ->
|
||||
maintenanceScope.launch {
|
||||
while (true) {
|
||||
delay(secs * 1000)
|
||||
try {
|
||||
store.optimize()
|
||||
} catch (e: CancellationException) {
|
||||
throw e // shutdown cancelled us; don't swallow it
|
||||
} catch (e: Exception) {
|
||||
// A missed refresh only means slightly staler planner
|
||||
// stats until the next tick — log and keep the loop.
|
||||
// Errors (OOM, etc.) are NOT swallowed.
|
||||
println("PRAGMA optimize failed: ${e.message}")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Runtime.getRuntime().addShutdownHook(
|
||||
Thread {
|
||||
// Each step wrapped so a throw in `server.stop()` doesn't
|
||||
// skip `relay.close()` (which closes the SQLite store).
|
||||
// Each step wrapped so a throw in any stage doesn't skip
|
||||
// `relay.close()` (which closes the SQLite store). The mirror
|
||||
// and maintenance go first: stop touching the store before the
|
||||
// queue and store beneath them shut down.
|
||||
runCatching { maintenanceScope.cancel() }
|
||||
runCatching { mirror?.close() }
|
||||
runCatching { server.stop() }
|
||||
runCatching { relay.close() }
|
||||
},
|
||||
@@ -183,6 +300,10 @@ fun main(args: Array<String>) {
|
||||
|
||||
println("geode listening on ${server.url}")
|
||||
println("NIP-11 info doc: curl -H 'Accept: application/nostr+json' http://$advertisedHost:$port$path")
|
||||
if (upstreams.isNotEmpty()) {
|
||||
val trusted = upstreams.count { it.trusted }
|
||||
println("mirroring ${upstreams.size} upstream relay(s), $trusted trusted (signature verification skipped)")
|
||||
}
|
||||
|
||||
// Park the main thread; shutdown hook handles teardown.
|
||||
Thread.currentThread().join()
|
||||
|
||||
@@ -40,21 +40,30 @@ import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
|
||||
* per-event tokenization on ingest (relayBench measured it at roughly a
|
||||
* quarter of write cost) at the price of `search` filters matching
|
||||
* nothing — pair it with a NIP-11 doc that doesn't advertise 50.
|
||||
* @param liveNegentropyIndex keep the always-current `(created_at, id)`
|
||||
* set that serves full-corpus NIP-77 NEG-OPENs without a scan + seal
|
||||
* (strfry answers those off its live tree; the scan+seal path measured
|
||||
* ~340 ms per cold open at 50k events). Costs ~140 B/event of JVM heap
|
||||
* (hex-string ids) and one indexed pre-SELECT per replaceable insert.
|
||||
* `[negentropy].live_index = false` turns it off.
|
||||
*/
|
||||
fun relayIndexingStrategy(fullTextSearch: Boolean = true) =
|
||||
DefaultIndexingStrategy(
|
||||
indexEventsByCreatedAtAlone = true,
|
||||
// Authors-only filters (no kinds) are relay-common — archives,
|
||||
// migration tools. strfry maintains the same (pubkey, created_at)
|
||||
// index unconditionally; without it the filter walks the whole
|
||||
// time index.
|
||||
indexEventsByPubkeyAlone = true,
|
||||
indexFullTextSearch = fullTextSearch,
|
||||
// Tokenize off the commit path; NostrServer drives the catch-up
|
||||
// worker and search queries drain it first, so NIP-50 stays
|
||||
// exactly as fresh while publishes stop paying for it.
|
||||
deferFullTextSearchIndexing = fullTextSearch,
|
||||
)
|
||||
fun relayIndexingStrategy(
|
||||
fullTextSearch: Boolean = true,
|
||||
liveNegentropyIndex: Boolean = true,
|
||||
) = DefaultIndexingStrategy(
|
||||
indexEventsByCreatedAtAlone = true,
|
||||
// Authors-only filters (no kinds) are relay-common — archives,
|
||||
// migration tools. strfry maintains the same (pubkey, created_at)
|
||||
// index unconditionally; without it the filter walks the whole
|
||||
// time index.
|
||||
indexEventsByPubkeyAlone = true,
|
||||
indexFullTextSearch = fullTextSearch,
|
||||
// Tokenize off the commit path; NostrServer drives the catch-up
|
||||
// worker and search queries drain it first, so NIP-50 stays
|
||||
// exactly as fresh while publishes stop paying for it.
|
||||
deferFullTextSearchIndexing = fullTextSearch,
|
||||
maintainLiveNegentropyIndex = liveNegentropyIndex,
|
||||
)
|
||||
|
||||
/** Stock relay strategy — everything on, matching geode's defaults. */
|
||||
val RelayIndexingStrategy = relayIndexingStrategy()
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.geode.config
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper
|
||||
import com.fasterxml.jackson.databind.node.ObjectNode
|
||||
|
||||
/**
|
||||
* Strict boot-time check for a `[[mirror]].filter` JSON string. The
|
||||
* NIP-01 [com.vitorpamplona.quartz.nip01Core.relay.filters.Filter]
|
||||
* deserializer is tolerant by design — it skips unknown keys and drops
|
||||
* wrong-typed entries — which is right for untrusted wire input but
|
||||
* wrong for operator config: a typo like `{"kindss":[4]}` would silently
|
||||
* parse to an empty (match-everything) filter. Because this filter is
|
||||
* the containment boundary for a `trusted = true` upstream (it caps what
|
||||
* a skip-verify upstream may inject), a silent mis-parse would widen
|
||||
* that trust to the whole firehose. So we reject the obvious mistakes at
|
||||
* boot instead of degrading the boundary.
|
||||
*/
|
||||
object MirrorFilterValidator {
|
||||
/** NIP-01 filter keys the deserializer recognizes; anything else is a typo. */
|
||||
private val KNOWN_KEYS = setOf("ids", "authors", "kinds", "since", "until", "limit", "search")
|
||||
|
||||
/** Keys whose value must be a JSON array. */
|
||||
private val ARRAY_KEYS = setOf("ids", "authors", "kinds")
|
||||
|
||||
private val mapper = ObjectMapper()
|
||||
|
||||
/**
|
||||
* Throws [IllegalArgumentException] when [json] is not a JSON object,
|
||||
* carries an unrecognized top-level key, or gives a list-typed field
|
||||
* (`ids`/`authors`/`kinds`, or a `#tag`/`&tag`) a non-array value.
|
||||
*/
|
||||
fun validate(
|
||||
url: String,
|
||||
json: String,
|
||||
) {
|
||||
val node =
|
||||
try {
|
||||
mapper.readTree(json)
|
||||
} catch (e: Exception) {
|
||||
throw IllegalArgumentException("[[mirror]] filter for $url is not valid JSON: $json", e)
|
||||
}
|
||||
require(node is ObjectNode) {
|
||||
"[[mirror]] filter for $url must be a JSON object (a NIP-01 filter), got: $json"
|
||||
}
|
||||
node.fieldNames().forEach { field ->
|
||||
val isTagKey = field.length > 1 && (field[0] == '#' || field[0] == '&')
|
||||
require(field in KNOWN_KEYS || isTagKey) {
|
||||
"[[mirror]] filter for $url has unknown key \"$field\" — a NIP-01 filter uses " +
|
||||
"ids/authors/kinds/since/until/limit/search or #tag/&tag"
|
||||
}
|
||||
if (field in ARRAY_KEYS || isTagKey) {
|
||||
require(node.get(field).isArray) {
|
||||
"[[mirror]] filter for $url: \"$field\" must be a JSON array"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -45,6 +45,8 @@ data class StaticConfig(
|
||||
val authorization: AuthorizationSection = AuthorizationSection(),
|
||||
val admin: AdminSection = AdminSection(),
|
||||
val negentropy: NegentropySection = NegentropySection(),
|
||||
/** `[[mirror]]` entries — upstream relays this relay streams from. */
|
||||
val mirror: List<MirrorSection> = emptyList(),
|
||||
) {
|
||||
fun resolveInfo(fullTextSearch: Boolean = true): RelayInfo =
|
||||
RelayInfo(
|
||||
@@ -104,6 +106,30 @@ data class StaticConfig(
|
||||
/** True keeps an in-memory SQLite db (default — events vanish on restart). */
|
||||
val in_memory: Boolean = true,
|
||||
val file: String? = null,
|
||||
/**
|
||||
* Reader-connection pool size. `null` keeps quartz's default (4).
|
||||
* Only meaningful for file-backed stores; in-memory databases
|
||||
* share the single writer connection regardless.
|
||||
*/
|
||||
val readers: Int? = null,
|
||||
/**
|
||||
* `PRAGMA mmap_size` in bytes, e.g. `268435456` for 256 MiB.
|
||||
* Maps the database file into memory so reads skip the pread
|
||||
* syscall + page-cache copy. `null` keeps SQLite's default (off).
|
||||
*/
|
||||
val mmap_size: Long? = null,
|
||||
/**
|
||||
* `PRAGMA temp_store = MEMORY` — keeps sort/temp b-trees for
|
||||
* large queries in RAM instead of temp files.
|
||||
*/
|
||||
val temp_store_memory: Boolean = false,
|
||||
/**
|
||||
* Refresh query-planner statistics (`PRAGMA analysis_limit;
|
||||
* PRAGMA optimize`) every this many seconds. Incremental and
|
||||
* usually a no-op, but keeps the planner from drifting onto the
|
||||
* wrong index as the corpus grows/changes shape. `null` = never.
|
||||
*/
|
||||
val optimize_interval_seconds: Long? = null,
|
||||
)
|
||||
|
||||
data class OptionsSection(
|
||||
@@ -151,6 +177,15 @@ data class StaticConfig(
|
||||
val frame_size_limit: Long = 500_000L,
|
||||
val max_sync_events: Int = 1_000_000,
|
||||
val max_sessions_per_connection: Int = 200,
|
||||
/**
|
||||
* Keep an always-current in-memory `(created_at, id)` set so
|
||||
* full-corpus NEG-OPENs skip the table scan + seal (strfry
|
||||
* parity). ~140 B per stored event of heap; on by default. Only
|
||||
* built once the first full-corpus NEG-OPEN arrives, and only
|
||||
* when the corpus fits `max_sync_events` (an over-cap corpus
|
||||
* answers NEG-ERR from a capped scan instead).
|
||||
*/
|
||||
val live_index: Boolean = true,
|
||||
)
|
||||
|
||||
data class AuthorizationSection(
|
||||
@@ -160,6 +195,46 @@ data class StaticConfig(
|
||||
val kind_blacklist: List<Int> = emptyList(),
|
||||
)
|
||||
|
||||
/**
|
||||
* One upstream relay to mirror, declared as a `[[mirror]]` TOML array
|
||||
* entry. The relay dials [url] itself, subscribes to everything newer
|
||||
* than `now - backfill_seconds`, and ingests the stream through the
|
||||
* same group-commit writer as client publishes (reconnects and
|
||||
* re-subscribes automatically).
|
||||
*
|
||||
* [trusted] is the relay-to-relay trust switch (strfry's model):
|
||||
* `true` skips Schnorr signature verification for events arriving on
|
||||
* this connection — sound only when the upstream verifies its own
|
||||
* ingest, which is why it defaults to `false` (mirror-but-verify).
|
||||
* The identity being trusted is the URL this relay dialed (TLS-
|
||||
* authenticated for `wss://`), never anything a peer claims.
|
||||
*/
|
||||
data class MirrorSection(
|
||||
val url: String,
|
||||
val trusted: Boolean = false,
|
||||
/** How far back the initial subscription reaches. 0 = live-only. */
|
||||
val backfill_seconds: Long = 0L,
|
||||
/**
|
||||
* Optional NIP-01 filter as a JSON object string (strfry-router
|
||||
* parity), e.g. `'{"kinds":[0,1,3],"#t":["nostr"]}'`. Scopes what
|
||||
* this upstream is asked for AND what it is allowed to deliver —
|
||||
* every received event is re-checked against it before ingest, so
|
||||
* even a trusted upstream can't push events outside the declared
|
||||
* scope. `since` is managed by the mirror (see [backfill_seconds])
|
||||
* and `limit` is transport-level, so both are ignored if present.
|
||||
* Omitted = mirror everything. For several disjoint filters, add
|
||||
* several `[[mirror]]` entries with the same url.
|
||||
*/
|
||||
val filter: String? = null,
|
||||
/**
|
||||
* Flow direction, strfry-router's `dir`: `"down"` (pull from the
|
||||
* upstream — the default), `"up"` (push this relay's matching
|
||||
* events to it), or `"both"`. Both-way mirrors suppress echoes
|
||||
* (an event pulled down is not pushed straight back).
|
||||
*/
|
||||
val dir: String = "down",
|
||||
)
|
||||
|
||||
/**
|
||||
* NIP-86 admin. [pubkeys] non-empty opens the POST endpoint at the
|
||||
* relay path; only NIP-98 tokens signed by these pubkeys dispatch.
|
||||
@@ -178,6 +253,25 @@ data class StaticConfig(
|
||||
val state_file: String? = null,
|
||||
)
|
||||
|
||||
/**
|
||||
* Boot-time sanity check for values the TOML types can't constrain.
|
||||
* Throws [IllegalArgumentException] (fail-loud at startup) rather
|
||||
* than letting a nonsensical knob degrade the running relay — a zero
|
||||
* reader pool hangs every query, a non-positive optimize interval
|
||||
* busy-loops the writer. Call once after parsing, before building
|
||||
* the store.
|
||||
*/
|
||||
fun validate() {
|
||||
database.readers?.let {
|
||||
require(it >= 1) { "[database].readers must be >= 1 (got $it); a 0/negative pool can never answer a query" }
|
||||
}
|
||||
database.optimize_interval_seconds?.let {
|
||||
require(it > 0) {
|
||||
"[database].optimize_interval_seconds must be > 0 (got $it); a non-positive interval busy-loops PRAGMA optimize under the writer mutex"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
private val mapper = tomlMapper { }
|
||||
|
||||
|
||||
@@ -0,0 +1,496 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.geode.mirror
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.NostrServer
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import okhttp3.OkHttpClient
|
||||
import java.time.Duration
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
|
||||
/**
|
||||
* Which way events flow between this relay and one upstream —
|
||||
* strfry-router's per-stream `dir`.
|
||||
*/
|
||||
enum class MirrorDirection {
|
||||
/** Pull: subscribe to the upstream and ingest what it sends. */
|
||||
DOWN,
|
||||
|
||||
/** Push: publish this relay's matching events to the upstream. */
|
||||
UP,
|
||||
|
||||
/** Pull and push. Echo suppression keeps the two from ping-ponging. */
|
||||
BOTH,
|
||||
;
|
||||
|
||||
companion object {
|
||||
/** Parses strfry's `"down"` / `"up"` / `"both"`; null if unknown. */
|
||||
fun parse(value: String): MirrorDirection? =
|
||||
when (value.lowercase()) {
|
||||
"down" -> DOWN
|
||||
"up" -> UP
|
||||
"both" -> BOTH
|
||||
else -> null
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* One upstream relay this relay mirrors, from the `[[mirror]]` config.
|
||||
*
|
||||
* [trusted] is the relay-to-relay trust switch: events streamed from this
|
||||
* upstream skip Schnorr signature verification on ingest. The trusted
|
||||
* identity is [url] — the address *this* relay dialed (TLS-authenticated
|
||||
* for `wss://`), never anything the peer claims — so the skip can't be
|
||||
* hijacked by an inbound client. Only meaningful for [MirrorDirection.DOWN]
|
||||
* / [MirrorDirection.BOTH]; the up direction never verifies (the upstream
|
||||
* does its own gatekeeping).
|
||||
*/
|
||||
class MirrorUpstream(
|
||||
val url: NormalizedRelayUrl,
|
||||
val trusted: Boolean,
|
||||
/** How far back the initial replay reaches, in BOTH directions. 0 = live-only from connect. */
|
||||
val backfillSeconds: Long = 0L,
|
||||
/**
|
||||
* Optional scope for this upstream (strfry-router's per-stream
|
||||
* `filter`). Applied symmetrically: down, it shapes the REQ sent
|
||||
* upstream AND every delivered event is re-checked before ingest —
|
||||
* so even a [trusted] upstream can only inject events inside the
|
||||
* declared scope; up, it selects which local events are pushed. Its
|
||||
* `since`/`limit` are ignored ([backfillSeconds] owns the time
|
||||
* window; the subscriptions are unbounded). `null` mirrors
|
||||
* everything.
|
||||
*/
|
||||
val filter: Filter? = null,
|
||||
/** Flow direction — strfry-router's `dir`. Defaults to pull-only. */
|
||||
val direction: MirrorDirection = MirrorDirection.DOWN,
|
||||
)
|
||||
|
||||
/**
|
||||
* Streams events from configured upstream relays into the local relay —
|
||||
* geode's equivalent of `strfry router` in the "down" direction.
|
||||
*
|
||||
* One [NostrClient] holds every upstream connection; the client owns
|
||||
* reconnects, exponential backoff, and re-sending the REQ after a drop.
|
||||
* Each upstream gets its own subscription over an open filter
|
||||
* (`since = now - backfill`), and every EVENT that arrives is handed to
|
||||
* [NostrServer.ingest] — the same group-commit writer and live fanout a
|
||||
* client publish takes — with `skipVerify` set for [MirrorUpstream.trusted]
|
||||
* upstreams (the upstream already verified its ingest; re-verifying here
|
||||
* only burns CPU — Schnorr verify profiles at ~8% of busy ingest CPU).
|
||||
*
|
||||
* After a reconnect the upstream replays everything since the boot-time
|
||||
* `since`; replayed duplicates are rejected by the store's unique id
|
||||
* constraint and only show up in [rejected].
|
||||
*
|
||||
* Listener callbacks can't suspend, so events funnel through an unbounded
|
||||
* [inbound] channel into one consumer coroutine whose [NostrServer.ingest]
|
||||
* call suspends on the ingest queue's backpressure. The buffer is unbounded
|
||||
* for the same reason the client's receive channels are (see
|
||||
* `BasicOkHttpWebSocket`): blocking the socket reader parks the backlog on
|
||||
* infrastructure that isn't ours.
|
||||
*/
|
||||
class MirrorWorker(
|
||||
private val upstreams: List<MirrorUpstream>,
|
||||
private val server: NostrServer,
|
||||
/**
|
||||
* Transport override for tests (e.g. `InProcessRelays`). Defaults to
|
||||
* a real OkHttp WebSocket per upstream.
|
||||
*/
|
||||
websocketBuilder: WebsocketBuilder? = null,
|
||||
) : AutoCloseable {
|
||||
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||
|
||||
/**
|
||||
* The ping interval matters on a long-running daemon: a half-open
|
||||
* upstream connection (network drop with no FIN) never fires
|
||||
* onDisconnected on its own, so without pings the client would
|
||||
* believe it is connected forever and stop mirroring silently. A
|
||||
* missed pong fails the socket, which routes into the normal
|
||||
* disconnect → backoff → re-dial path. Same value the Android app
|
||||
* uses for its relay pool.
|
||||
*/
|
||||
private val okhttp: OkHttpClient? =
|
||||
if (websocketBuilder == null) {
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
.socketFactory(TcpNoDelaySocketFactory)
|
||||
.pingInterval(Duration.ofSeconds(PING_INTERVAL_SECS))
|
||||
.build()
|
||||
} else {
|
||||
null
|
||||
}
|
||||
|
||||
private val client =
|
||||
NostrClient(
|
||||
websocketBuilder = websocketBuilder ?: BasicOkHttpWebSocket.Builder { okhttp!! },
|
||||
parentScope = scope,
|
||||
)
|
||||
|
||||
private class Inbound(
|
||||
val event: Event,
|
||||
val skipVerify: Boolean,
|
||||
)
|
||||
|
||||
private val inbound = Channel<Inbound>(Channel.UNLIMITED)
|
||||
|
||||
/** Events accepted into the local store (excludes duplicates). */
|
||||
val accepted = AtomicLong(0)
|
||||
|
||||
/** Events the store rejected — mostly duplicate replays after a reconnect. */
|
||||
val rejected = AtomicLong(0)
|
||||
|
||||
/**
|
||||
* Deliveries dropped before ever reaching the store: events outside
|
||||
* the [MirrorUpstream.filter] scope (an upstream answering outside
|
||||
* the REQ it was given) and events delivered by a different relay
|
||||
* than the one the subscription dialed (a peer answering with a
|
||||
* subscription id that isn't its own).
|
||||
*/
|
||||
val filtered = AtomicLong(0)
|
||||
|
||||
/** Local events handed to the client's outbox for an up-direction upstream. */
|
||||
val sentUp = AtomicLong(0)
|
||||
|
||||
/**
|
||||
* How many times a down subscription advanced its `since` watermark
|
||||
* on a reconnect (test/observability hook). Each advance is one
|
||||
* avoided full-window replay.
|
||||
*/
|
||||
val sinceAdvances = AtomicLong(0)
|
||||
|
||||
/**
|
||||
* One down subscription's re-subscribe state. The upstream replays
|
||||
* everything at or after the REQ's `since` on every (re)connect; left
|
||||
* at the boot-time value, a month-old daemon re-streams a month of
|
||||
* events on every flap. So we track the newest `created_at` ingested
|
||||
* from this upstream and, on a disconnect, advance the REQ's `since`
|
||||
* to `watermark - overlap` before the reconnect re-sends it. The
|
||||
* overlap re-requests a small tail (dup-safe against the store's
|
||||
* unique-id constraint) to cover out-of-order streaming and clock
|
||||
* skew. Advancing only on disconnect keeps a healthy connection from
|
||||
* ever re-querying — no steady-state cost.
|
||||
*/
|
||||
private inner class DownSub(
|
||||
val subId: String,
|
||||
val up: MirrorUpstream,
|
||||
val scopedBase: Filter,
|
||||
val listener: SubscriptionListener,
|
||||
initialSince: Long,
|
||||
val watermark: AtomicLong,
|
||||
) {
|
||||
@Volatile
|
||||
var issuedSince: Long = initialSince
|
||||
|
||||
fun advanceSinceOnReconnect() {
|
||||
val candidate = watermark.get() - WATERMARK_OVERLAP_SECS
|
||||
if (candidate > issuedSince) {
|
||||
issuedSince = candidate
|
||||
client.subscribe(
|
||||
subId = subId,
|
||||
filters = mapOf(up.url to listOf(scopedBase.copy(since = candidate))),
|
||||
listener = listener,
|
||||
)
|
||||
sinceAdvances.incrementAndGet()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private val downSubs = mutableListOf<DownSub>()
|
||||
|
||||
/**
|
||||
* Recently exchanged event ids, one set per up-capable upstream —
|
||||
* the echo suppressor for [MirrorDirection.BOTH]. An event pulled
|
||||
* DOWN from an upstream must not be pushed straight back UP to it
|
||||
* (and one we pushed up must not be re-ingested when the upstream
|
||||
* fans it back on our down subscription). Bounded LRU: eviction only
|
||||
* costs a wasted round trip that the stores' unique-id constraints
|
||||
* absorb, so correctness never depends on it.
|
||||
*/
|
||||
private class RecentIds(
|
||||
private val capacity: Int,
|
||||
) {
|
||||
private val map =
|
||||
object : LinkedHashMap<String, Boolean>(capacity, 0.75f, true) {
|
||||
override fun removeEldestEntry(eldest: Map.Entry<String, Boolean>) = size > capacity
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun add(id: String) {
|
||||
map[id] = true
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun contains(id: String): Boolean = map.containsKey(id)
|
||||
}
|
||||
|
||||
/** Open in-process sessions feeding the up direction; closed with the worker. */
|
||||
private val upSessions = mutableListOf<AutoCloseable>()
|
||||
|
||||
/** Dials every upstream and starts streaming. Call once. */
|
||||
fun start() {
|
||||
scope.launch {
|
||||
for (msg in inbound) {
|
||||
try {
|
||||
server.ingest(msg.event, msg.skipVerify) { outcome ->
|
||||
when (outcome) {
|
||||
IEventStore.InsertOutcome.Accepted -> accepted.incrementAndGet()
|
||||
is IEventStore.InsertOutcome.Rejected -> {
|
||||
rejected.incrementAndGet()
|
||||
Log.d("MirrorWorker") { "rejected ${msg.event.id}: ${outcome.reason}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Throwable) {
|
||||
// A server shutting down closes its ingest queue while
|
||||
// events may still be buffered here; an uncaught throw
|
||||
// would leak into the scope's default handler (and
|
||||
// poison unrelated runTest tests on CI). Stop pulling —
|
||||
// the relay beneath us is going away.
|
||||
Log.w("MirrorWorker") { "ingest failed, stopping mirror consumer: ${e.message}" }
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
val since = TimeUtils.now()
|
||||
upstreams.forEachIndexed { i, up ->
|
||||
// Echo suppression only matters when events can flow both
|
||||
// ways on the same upstream.
|
||||
val exchanged = if (up.direction == MirrorDirection.BOTH) RecentIds(EXCHANGED_IDS_CAPACITY) else null
|
||||
|
||||
// The operator's filter scopes the subscriptions in both
|
||||
// directions; the mirror owns the time window (since) and
|
||||
// never bounds the result (limit). scopedBase carries no
|
||||
// since — the down path applies (and later advances) it via
|
||||
// the watermark; the up path uses the fixed initial value.
|
||||
val scopedBase = (up.filter ?: Filter()).copy(since = null, limit = null)
|
||||
val initialSince = since - up.backfillSeconds
|
||||
|
||||
if (up.direction != MirrorDirection.UP) {
|
||||
downSubs += startDown(i, up, scopedBase, initialSince, exchanged)
|
||||
}
|
||||
if (up.direction != MirrorDirection.DOWN) {
|
||||
startUp(up, scopedBase.copy(since = initialSince), exchanged)
|
||||
}
|
||||
}
|
||||
client.connect()
|
||||
|
||||
// Watermark advance: when an upstream drops, bump its REQ's since
|
||||
// to the newest event we ingested from it (minus an overlap) so
|
||||
// the reconnect doesn't replay the whole window since boot. Only
|
||||
// fires on the connected→disconnected edge, so a stable link
|
||||
// never re-queries.
|
||||
scope.launch {
|
||||
var prev = emptySet<NormalizedRelayUrl>()
|
||||
client.connectedRelaysFlow().collect { current ->
|
||||
val dropped = prev - current
|
||||
if (dropped.isNotEmpty()) {
|
||||
downSubs.forEach { if (it.up.url in dropped) it.advanceSinceOnReconnect() }
|
||||
}
|
||||
prev = current
|
||||
}
|
||||
}
|
||||
|
||||
// Retry pump. NostrClient re-dials once on disconnect and then
|
||||
// relies on its 60s keep-alive — measured as a 61s mirror blackout
|
||||
// when an upstream restarts and the immediate re-dial races the
|
||||
// port rebind. Poking more often costs nothing: reconnectIfNeedsTo
|
||||
// skips connected relays, and each relay's exponential backoff
|
||||
// (1s doubling, 5min cap) still gates actual dial attempts, so a
|
||||
// long-dead upstream is not hammered — only a briefly-restarting
|
||||
// one is picked back up in seconds instead of a minute.
|
||||
scope.launch {
|
||||
while (true) {
|
||||
delay(RECONNECT_POKE_MS)
|
||||
client.reconnect(onlyIfChanged = true, ignoreRetryDelays = false)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Down direction: subscribe to the upstream, ingest what it sends. */
|
||||
private fun startDown(
|
||||
index: Int,
|
||||
up: MirrorUpstream,
|
||||
scopedBase: Filter,
|
||||
initialSince: Long,
|
||||
exchanged: RecentIds?,
|
||||
): DownSub {
|
||||
// watermark tracks the newest created_at ingested from this
|
||||
// upstream; seeded at initialSince so a still-catching-up
|
||||
// backfill can never advance the since backwards.
|
||||
val watermark = AtomicLong(initialSince)
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
// The trusted identity is the relay THIS subscription
|
||||
// dialed. The pool dispatches EVENTs by subscription
|
||||
// id alone and every upstream shares this one client,
|
||||
// so a hostile co-configured upstream could answer
|
||||
// with another subscription's id and ride its trust —
|
||||
// bind the decision to the delivering relay, not the
|
||||
// sub id.
|
||||
if (relay != up.url) {
|
||||
filtered.incrementAndGet()
|
||||
Log.w("MirrorWorker") { "dropped event delivered by ${relay.url} on ${up.url.url}'s subscription: ${event.id}" }
|
||||
return
|
||||
}
|
||||
// strfry-router parity: never take the upstream's
|
||||
// word for what matched. Re-checking the configured
|
||||
// scope here means even a trusted (skip-verify)
|
||||
// upstream can only inject events the operator
|
||||
// declared — the REQ shapes what we ask for, this
|
||||
// shapes what we accept.
|
||||
if (up.filter != null && !up.filter.match(event)) {
|
||||
filtered.incrementAndGet()
|
||||
Log.d("MirrorWorker") { "out-of-scope from ${relay.url}: ${event.id}" }
|
||||
return
|
||||
}
|
||||
// BOTH: an event we just pushed up is fanned back on
|
||||
// this subscription — it already exists locally.
|
||||
if (exchanged?.contains(event.id) == true) return
|
||||
exchanged?.add(event.id)
|
||||
// Advance the reconnect watermark past this event.
|
||||
watermark.updateAndGet { if (event.createdAt > it) event.createdAt else it }
|
||||
inbound.trySend(Inbound(event, up.trusted))
|
||||
}
|
||||
|
||||
override fun onCannotConnect(
|
||||
relay: NormalizedRelayUrl,
|
||||
message: String,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
Log.w("MirrorWorker") { "cannot reach upstream ${relay.url}: $message" }
|
||||
}
|
||||
}
|
||||
val subId = "geode-mirror-$index"
|
||||
client.subscribe(
|
||||
subId = subId,
|
||||
filters = mapOf(up.url to listOf(scopedBase.copy(since = initialSince))),
|
||||
listener = listener,
|
||||
)
|
||||
return DownSub(subId, up, scopedBase, listener, initialSince, watermark)
|
||||
}
|
||||
|
||||
/**
|
||||
* Up direction: an in-process session on the LOCAL relay subscribes
|
||||
* with the same scoped filter — stored replay covers the backfill
|
||||
* window, the live tail covers everything after — and each matching
|
||||
* event is handed to the client's outbox for [MirrorUpstream.url].
|
||||
* The outbox owns delivery: it re-sends on reconnect until the
|
||||
* upstream OKs, and the upstream's own duplicate handling absorbs
|
||||
* replays. Going through a real session (not the store) means the
|
||||
* relay's policy chain gates what leaves, same as any client.
|
||||
*/
|
||||
private fun startUp(
|
||||
up: MirrorUpstream,
|
||||
scopedFilter: Filter,
|
||||
exchanged: RecentIds?,
|
||||
) {
|
||||
val session =
|
||||
server.connect { json ->
|
||||
if (!json.startsWith("[\"EVENT\"")) return@connect
|
||||
val event =
|
||||
runCatching { (OptimizedJsonMapper.fromJsonToMessage(json) as? EventMessage)?.event }
|
||||
.getOrNull() ?: return@connect
|
||||
// BOTH: don't push back what we just pulled down.
|
||||
if (exchanged?.contains(event.id) == true) return@connect
|
||||
exchanged?.add(event.id)
|
||||
client.publish(event, setOf(up.url))
|
||||
sentUp.incrementAndGet()
|
||||
}
|
||||
upSessions += AutoCloseable { session.close() }
|
||||
scope.launch {
|
||||
session.receive(OptimizedJsonMapper.toJson(ReqCmd("geode-mirror-up", listOf(scopedFilter))))
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Stops pulling from every upstream. In-flight ingest submissions
|
||||
* drain through the server's queue; events still buffered in
|
||||
* [inbound] are dropped — the next boot's `since` overlaps only if
|
||||
* the operator configured a backfill window, which is the documented
|
||||
* trade-off of a live mirror.
|
||||
*/
|
||||
override fun close() {
|
||||
// Close the client first so no listener callback races the
|
||||
// channel close below.
|
||||
upSessions.forEach { runCatching { it.close() } }
|
||||
runCatching { client.close() }
|
||||
inbound.close()
|
||||
scope.cancel()
|
||||
okhttp?.dispatcher?.executorService?.shutdown()
|
||||
okhttp?.connectionPool?.evictAll()
|
||||
}
|
||||
|
||||
private companion object {
|
||||
/** Matches the Android app's relay-pool WebSocket ping interval. */
|
||||
const val PING_INTERVAL_SECS = 120L
|
||||
|
||||
/** How often the retry pump nudges disconnected upstreams. */
|
||||
const val RECONNECT_POKE_MS = 5_000L
|
||||
|
||||
/**
|
||||
* Overlap (seconds) subtracted from the reconnect `since`
|
||||
* watermark. Re-requests a small tail on reconnect to cover
|
||||
* out-of-order streaming and clock skew; the store's unique-id
|
||||
* constraint drops the duplicates. Bounds worst-case replay after
|
||||
* a flap to this window instead of everything since boot.
|
||||
*/
|
||||
const val WATERMARK_OVERLAP_SECS = 300L
|
||||
|
||||
/**
|
||||
* Per-upstream echo-suppression LRU size (BOTH direction only).
|
||||
* Covers the burst window between pulling an event down and the
|
||||
* up-session seeing its local fanout; eviction only costs a
|
||||
* duplicate round trip.
|
||||
*/
|
||||
const val EXCHANGED_IDS_CAPACITY = 8_192
|
||||
}
|
||||
}
|
||||
@@ -20,9 +20,12 @@
|
||||
*/
|
||||
package com.vitorpamplona.geode.config
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import java.io.File
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
@@ -46,6 +49,130 @@ class StaticConfigTest {
|
||||
assertEquals(false, c.options.verify_signatures)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun parsesDatabaseTuningKnobs() {
|
||||
val toml =
|
||||
"""
|
||||
[database]
|
||||
in_memory = false
|
||||
file = "/tmp/x.db"
|
||||
readers = 8
|
||||
mmap_size = 268435456
|
||||
temp_store_memory = true
|
||||
optimize_interval_seconds = 3600
|
||||
""".trimIndent()
|
||||
|
||||
val c = StaticConfig.fromToml(toml)
|
||||
|
||||
assertEquals(8, c.database.readers)
|
||||
assertEquals(268435456L, c.database.mmap_size)
|
||||
assertEquals(true, c.database.temp_store_memory)
|
||||
assertEquals(3600L, c.database.optimize_interval_seconds)
|
||||
|
||||
// And all knobs default to off/null so plain configs are untouched.
|
||||
val d = StaticConfig.fromToml("")
|
||||
assertEquals(null, d.database.readers)
|
||||
assertEquals(null, d.database.mmap_size)
|
||||
assertEquals(false, d.database.temp_store_memory)
|
||||
assertEquals(null, d.database.optimize_interval_seconds)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mirrorSectionDefaultsToEmpty() {
|
||||
assertTrue(StaticConfig.fromToml("").mirror.isEmpty())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun validateRejectsNonPositiveReaders() {
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
StaticConfig.fromToml("[database]\nreaders = 0").validate()
|
||||
}
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
StaticConfig.fromToml("[database]\nreaders = -1").validate()
|
||||
}
|
||||
// A sane pool passes.
|
||||
StaticConfig.fromToml("[database]\nreaders = 1").validate()
|
||||
// Unset passes (quartz default applies).
|
||||
StaticConfig.fromToml("").validate()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun validateRejectsNonPositiveOptimizeInterval() {
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
StaticConfig.fromToml("[database]\noptimize_interval_seconds = 0").validate()
|
||||
}
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
StaticConfig.fromToml("[database]\noptimize_interval_seconds = -5").validate()
|
||||
}
|
||||
StaticConfig.fromToml("[database]\noptimize_interval_seconds = 3600").validate()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mirrorFilterValidatorRejectsTyposAndScalars() {
|
||||
val url = "wss://up.example/"
|
||||
|
||||
// Unknown key (a typo) — must fail, not silently widen scope.
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
MirrorFilterValidator.validate(url, """{"kindss":[4]}""")
|
||||
}
|
||||
// List field given a scalar.
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
MirrorFilterValidator.validate(url, """{"authors":"abc"}""")
|
||||
}
|
||||
// Not an object.
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
MirrorFilterValidator.validate(url, """["kinds",1]""")
|
||||
}
|
||||
// Malformed JSON.
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
MirrorFilterValidator.validate(url, """{"kinds":[1,}""")
|
||||
}
|
||||
|
||||
// Valid shapes pass: recognized scalar + array + tag keys.
|
||||
MirrorFilterValidator.validate(url, """{"kinds":[0,1,3],"#t":["nostr"],"since":123,"limit":5,"search":"x"}""")
|
||||
MirrorFilterValidator.validate(url, """{"&p":["abc"]}""")
|
||||
MirrorFilterValidator.validate(url, "{}")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mirrorFilterJsonParsesToANip01Filter() {
|
||||
// The exact parse Main.kt runs on [[mirror]].filter at boot.
|
||||
val f = OptimizedJsonMapper.fromJsonTo<Filter>("""{"kinds":[0,1],"#t":["nostr"],"since":123,"limit":5}""")
|
||||
assertEquals(listOf(0, 1), f.kinds)
|
||||
assertEquals(listOf("nostr"), f.tags?.get("t"))
|
||||
assertEquals(123L, f.since)
|
||||
assertEquals(5, f.limit)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun parsesMirrorUpstreams() {
|
||||
val toml =
|
||||
"""
|
||||
[[mirror]]
|
||||
url = "wss://trusted.upstream.example/"
|
||||
trusted = true
|
||||
backfill_seconds = 3600
|
||||
filter = '{"kinds":[0,1,3],"#t":["nostr"]}'
|
||||
|
||||
[[mirror]]
|
||||
url = "wss://public.upstream.example/"
|
||||
""".trimIndent()
|
||||
|
||||
val c = StaticConfig.fromToml(toml)
|
||||
|
||||
assertEquals(2, c.mirror.size)
|
||||
assertEquals("wss://trusted.upstream.example/", c.mirror[0].url)
|
||||
assertEquals(true, c.mirror[0].trusted)
|
||||
assertEquals(3600L, c.mirror[0].backfill_seconds)
|
||||
assertEquals("""{"kinds":[0,1,3],"#t":["nostr"]}""", c.mirror[0].filter)
|
||||
// Trust and scoping are opt-in per upstream: the default is
|
||||
// mirror-everything-but-verify.
|
||||
assertEquals("wss://public.upstream.example/", c.mirror[1].url)
|
||||
assertEquals(false, c.mirror[1].trusted)
|
||||
assertEquals(0L, c.mirror[1].backfill_seconds)
|
||||
assertEquals(null, c.mirror[1].filter)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun parsesAllSectionsTogether() {
|
||||
val toml =
|
||||
|
||||
@@ -0,0 +1,164 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.geode.mirror
|
||||
|
||||
import com.vitorpamplona.geode.KtorRelay
|
||||
import com.vitorpamplona.geode.RelayEngine
|
||||
import com.vitorpamplona.geode.testing.preload
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import org.junit.After
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* The mirror is a long-running daemon, so surviving its upstream is not
|
||||
* optional. This drives the REAL production transport (OkHttp WebSocket
|
||||
* against a real Ktor port — no in-process shortcut): the upstream is
|
||||
* stopped mid-mirror and a new instance is brought up on the same port.
|
||||
* The worker must ride through the disconnect (NostrClient's backoff +
|
||||
* keep-alive own the re-dial), re-send its REQ, drop the replayed
|
||||
* duplicate, and pull the event that only exists on the new instance.
|
||||
*
|
||||
* Timing note: after a stable connection drops, the client's backoff is
|
||||
* reset to 1s and re-dial attempts double from there. The worker's 5s
|
||||
* retry pump nudges the pool well ahead of NostrClient's own 60s
|
||||
* keep-alive (without it, this test measured a 61s blackout when the
|
||||
* immediate re-dial raced the port rebind; with it, ~6s). The generous
|
||||
* timeout only covers a worst-case scheduling stall on a loaded CI
|
||||
* runner.
|
||||
*/
|
||||
class MirrorWorkerReconnectTest {
|
||||
private val downstreamStore = EventStore(null)
|
||||
private val downstream =
|
||||
RelayEngine(
|
||||
url = "ws://127.0.0.1:7797/".normalizeRelayUrl(),
|
||||
store = downstreamStore,
|
||||
parallelVerify = true,
|
||||
)
|
||||
|
||||
private var worker: MirrorWorker? = null
|
||||
private var upstreamServer: KtorRelay? = null
|
||||
private var upstreamEngine: RelayEngine? = null
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
worker?.close()
|
||||
upstreamServer?.stop(gracePeriodMillis = 0, timeoutMillis = 1_000)
|
||||
upstreamEngine?.close()
|
||||
downstream.close()
|
||||
}
|
||||
|
||||
private fun startUpstream(port: Int): KtorRelay {
|
||||
val engine = RelayEngine(url = "ws://127.0.0.1:7796/".normalizeRelayUrl())
|
||||
val server = KtorRelay(engine, host = "127.0.0.1", port = port).start()
|
||||
upstreamEngine = engine
|
||||
upstreamServer = server
|
||||
return server
|
||||
}
|
||||
|
||||
private fun forgedEvent(idSeed: Int): Event =
|
||||
Event(
|
||||
id = idSeed.toString().padStart(64, '0'),
|
||||
pubKey = "1".repeat(64),
|
||||
createdAt = TimeUtils.now() - idSeed,
|
||||
kind = 1,
|
||||
tags = emptyArray(),
|
||||
content = "forged $idSeed",
|
||||
sig = "f".repeat(128),
|
||||
)
|
||||
|
||||
private suspend fun awaitDownstreamCount(expected: Int) =
|
||||
withTimeout(120_000) {
|
||||
while (downstreamStore.count(Filter()) < expected) delay(50)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mirrorSurvivesUpstreamRestart() =
|
||||
runBlocking {
|
||||
val beforeRestart = forgedEvent(1)
|
||||
val afterRestart = forgedEvent(2)
|
||||
|
||||
// First upstream instance on an OS-assigned port.
|
||||
val first = startUpstream(port = 0)
|
||||
val port =
|
||||
first.url
|
||||
.normalizeRelayUrl()
|
||||
.url
|
||||
.substringAfterLast(':')
|
||||
.trimEnd('/')
|
||||
.toInt()
|
||||
upstreamEngine!!.preload(beforeRestart)
|
||||
|
||||
// Real OkHttp transport: websocketBuilder deliberately omitted.
|
||||
val mirror =
|
||||
MirrorWorker(
|
||||
upstreams =
|
||||
listOf(
|
||||
MirrorUpstream(
|
||||
url = first.url.normalizeRelayUrl(),
|
||||
trusted = true,
|
||||
backfillSeconds = 3600,
|
||||
),
|
||||
),
|
||||
server = downstream.server,
|
||||
).also { worker = it }
|
||||
mirror.start()
|
||||
awaitDownstreamCount(1)
|
||||
|
||||
// Kill the upstream: every socket drops, the port closes.
|
||||
first.stop(gracePeriodMillis = 0, timeoutMillis = 1_000)
|
||||
upstreamEngine!!.close()
|
||||
|
||||
// Bring up a NEW instance on the SAME port. Its store has the
|
||||
// old event (so the re-sent REQ replays a duplicate) plus one
|
||||
// that only exists post-restart.
|
||||
val second = startUpstream(port = port)
|
||||
upstreamEngine!!.preload(beforeRestart, afterRestart)
|
||||
|
||||
// The worker must reconnect on its own — no external poke —
|
||||
// re-subscribe, and pull the new event.
|
||||
awaitDownstreamCount(2)
|
||||
|
||||
assertEquals(
|
||||
setOf(beforeRestart.id, afterRestart.id),
|
||||
downstreamStore.query<Event>(Filter()).map { it.id }.toSet(),
|
||||
)
|
||||
// The duplicate replay of the pre-restart event was dropped by
|
||||
// the store, not double-inserted.
|
||||
assertEquals(2, downstreamStore.count(Filter()))
|
||||
|
||||
// The disconnect advanced the REQ's since watermark: the
|
||||
// reconnect subscribed from ~(newest event - overlap), not
|
||||
// from the boot-time `now - backfill`. Without this the
|
||||
// reconnect would replay the whole backfill window every flap.
|
||||
assertTrue(mirror.sinceAdvances.get() >= 1)
|
||||
|
||||
second.stop(gracePeriodMillis = 0, timeoutMillis = 1_000)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,261 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.geode.mirror
|
||||
|
||||
import com.vitorpamplona.geode.InProcessRelays
|
||||
import com.vitorpamplona.geode.RelayEngine
|
||||
import com.vitorpamplona.geode.testing.preload
|
||||
import com.vitorpamplona.geode.testing.publish
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.EventAssembler
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import org.junit.After
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* End-to-end `[[mirror]]` behavior over the in-process transport: a
|
||||
* downstream relay (verification ON via the parallel-verify queue) dials
|
||||
* an upstream and streams its events.
|
||||
*
|
||||
* - `trusted = true` — the relay-to-relay trust switch — must land
|
||||
* events the downstream could never verify itself.
|
||||
* - `trusted = false` must keep verify-everything semantics: forged
|
||||
* events from the upstream are dropped, valid ones land.
|
||||
*/
|
||||
class MirrorWorkerTest {
|
||||
private val upstreamUrl = RelayUrlNormalizer.normalize("ws://upstream.relay/")
|
||||
private val downstreamUrl = RelayUrlNormalizer.normalize("ws://downstream.relay/")
|
||||
|
||||
/** Upstream side: no verification (EmptyPolicy), so forged events store fine. */
|
||||
private val hub = InProcessRelays()
|
||||
|
||||
/** Downstream side: signature verification on, in the IngestQueue. */
|
||||
private val downstreamStore = EventStore(null)
|
||||
private val downstream =
|
||||
RelayEngine(
|
||||
url = downstreamUrl,
|
||||
store = downstreamStore,
|
||||
parallelVerify = true,
|
||||
)
|
||||
|
||||
private var worker: MirrorWorker? = null
|
||||
|
||||
private val signer = KeyPair()
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
worker?.close()
|
||||
downstream.close()
|
||||
hub.close()
|
||||
}
|
||||
|
||||
private fun startMirror(
|
||||
trusted: Boolean,
|
||||
filter: Filter? = null,
|
||||
direction: MirrorDirection = MirrorDirection.DOWN,
|
||||
): MirrorWorker =
|
||||
MirrorWorker(
|
||||
upstreams =
|
||||
listOf(
|
||||
MirrorUpstream(
|
||||
upstreamUrl,
|
||||
trusted = trusted,
|
||||
backfillSeconds = 3600,
|
||||
filter = filter,
|
||||
direction = direction,
|
||||
),
|
||||
),
|
||||
server = downstream.server,
|
||||
websocketBuilder = hub,
|
||||
).also {
|
||||
worker = it
|
||||
it.start()
|
||||
}
|
||||
|
||||
private fun forgedEvent(
|
||||
idSeed: Int,
|
||||
kind: Int = 1,
|
||||
): Event =
|
||||
Event(
|
||||
id = idSeed.toString().padStart(64, '0'),
|
||||
pubKey = "1".repeat(64),
|
||||
createdAt = TimeUtils.now() - idSeed,
|
||||
kind = kind,
|
||||
tags = emptyArray(),
|
||||
content = "forged $idSeed",
|
||||
sig = "f".repeat(128),
|
||||
)
|
||||
|
||||
private fun signedEvent(content: String): Event =
|
||||
EventAssembler.hashAndSign(
|
||||
pubKey = signer.pubKey.toHexKey(),
|
||||
createdAt = TimeUtils.now() - 5,
|
||||
kind = 1,
|
||||
tags = emptyArray(),
|
||||
content = content,
|
||||
privKey = signer.privKey!!,
|
||||
)
|
||||
|
||||
private suspend fun awaitDownstreamCount(expected: Int) =
|
||||
withTimeout(15_000) {
|
||||
while (downstreamStore.count(Filter()) < expected) delay(25)
|
||||
}
|
||||
|
||||
private suspend fun await(condition: suspend () -> Boolean) =
|
||||
withTimeout(15_000) {
|
||||
while (!condition()) delay(25)
|
||||
}
|
||||
|
||||
// NOTE: assertions are on the downstream STORE, not on exact counter
|
||||
// values — the client may legitimately re-send its REQ while settling
|
||||
// (connect + filter sync), so an upstream can replay an event twice and
|
||||
// the duplicate shows up as one extra rejection. That's the documented
|
||||
// mirror behavior, not a failure.
|
||||
|
||||
@Test
|
||||
fun trustedUpstreamLandsEventsTheDownstreamCannotVerify() =
|
||||
runBlocking {
|
||||
val stored = forgedEvent(1)
|
||||
val live = forgedEvent(2)
|
||||
|
||||
// Stored replay: exists on the upstream before the mirror dials.
|
||||
hub.getOrCreate(upstreamUrl).preload(stored)
|
||||
|
||||
startMirror(trusted = true)
|
||||
awaitDownstreamCount(1)
|
||||
|
||||
// Live tail: published upstream after the mirror subscribed.
|
||||
hub.getOrCreate(upstreamUrl).publish(live)
|
||||
awaitDownstreamCount(2)
|
||||
|
||||
val ids = downstreamStore.query<Event>(Filter()).map { it.id }.toSet()
|
||||
assertEquals(setOf(stored.id, live.id), ids)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun filterScopesTheMirrorToDeclaredKinds() =
|
||||
runBlocking {
|
||||
val wantedStored = forgedEvent(1, kind = 1)
|
||||
val unwantedStored = forgedEvent(2, kind = 7)
|
||||
hub.getOrCreate(upstreamUrl).preload(wantedStored, unwantedStored)
|
||||
|
||||
startMirror(trusted = true, filter = Filter(kinds = listOf(1)))
|
||||
awaitDownstreamCount(1)
|
||||
|
||||
// Live tail: the out-of-scope kind is published FIRST, so by the
|
||||
// time the in-scope one lands downstream (same connection, same
|
||||
// ordered pipeline), the kind-7 has already had its chance.
|
||||
hub.getOrCreate(upstreamUrl).publish(forgedEvent(3, kind = 7))
|
||||
val wantedLive = forgedEvent(4, kind = 1)
|
||||
hub.getOrCreate(upstreamUrl).publish(wantedLive)
|
||||
awaitDownstreamCount(2)
|
||||
|
||||
val stored = downstreamStore.query<Event>(Filter())
|
||||
assertEquals(setOf(wantedStored.id, wantedLive.id), stored.map { it.id }.toSet())
|
||||
assertTrue(stored.all { it.kind == 1 })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun upDirectionPushesLocalEventsToTheUpstream() =
|
||||
runBlocking {
|
||||
// Pre-existing local event: the up replay (backfill window)
|
||||
// must carry it; then a live local publish must follow.
|
||||
val preexisting = forgedEvent(1)
|
||||
downstream.preload(preexisting)
|
||||
|
||||
startMirror(trusted = false, direction = MirrorDirection.UP)
|
||||
|
||||
val upstreamStore = hub.getOrCreate(upstreamUrl).store
|
||||
await { upstreamStore.count(Filter()) >= 1 }
|
||||
|
||||
// Live events go through the verifying publish path, so they
|
||||
// must be genuinely signed (preload bypasses verification).
|
||||
val live = signedEvent("up live")
|
||||
downstream.publish(live)
|
||||
await { upstreamStore.count(Filter()) >= 2 }
|
||||
|
||||
val ids = upstreamStore.query<Event>(Filter()).map { it.id }.toSet()
|
||||
assertEquals(setOf(preexisting.id, live.id), ids)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun bothDirectionConvergesWithoutPingPong() =
|
||||
runBlocking {
|
||||
// One event only the upstream has, one only the local relay
|
||||
// has. BOTH must converge the two stores; echo suppression
|
||||
// (plus store dedup as the backstop) must keep the shared
|
||||
// events from bouncing.
|
||||
val upstreamOnly = forgedEvent(1)
|
||||
val localOnly = forgedEvent(2)
|
||||
hub.getOrCreate(upstreamUrl).preload(upstreamOnly)
|
||||
downstream.preload(localOnly)
|
||||
|
||||
val mirror = startMirror(trusted = true, direction = MirrorDirection.BOTH)
|
||||
|
||||
val upstreamStore = hub.getOrCreate(upstreamUrl).store
|
||||
await { downstreamStore.count(Filter()) == 2 && upstreamStore.count(Filter()) == 2 }
|
||||
|
||||
val expected = setOf(upstreamOnly.id, localOnly.id)
|
||||
assertEquals(expected, downstreamStore.query<Event>(Filter()).map { it.id }.toSet())
|
||||
assertEquals(expected, upstreamStore.query<Event>(Filter()).map { it.id }.toSet())
|
||||
|
||||
// Live: an event published locally reaches the upstream AND
|
||||
// its echo back down doesn't disturb either store. Signed,
|
||||
// because the local publish path verifies.
|
||||
val live = signedEvent("both live")
|
||||
downstream.publish(live)
|
||||
await { upstreamStore.count(Filter()) == 3 }
|
||||
// Let any echo settle, then confirm counts are exact.
|
||||
delay(500)
|
||||
assertEquals(3, downstreamStore.count(Filter()))
|
||||
assertEquals(3, upstreamStore.count(Filter()))
|
||||
assertTrue(mirror.sentUp.get() >= 2)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun untrustedUpstreamStillVerifiesEverything() =
|
||||
runBlocking {
|
||||
val forged = forgedEvent(1)
|
||||
val valid = signedEvent("the real one")
|
||||
hub.getOrCreate(upstreamUrl).preload(forged, valid)
|
||||
|
||||
val mirror = startMirror(trusted = false)
|
||||
|
||||
// The valid event lands; the forged one is verified and dropped
|
||||
// (a forged event can never land untrusted, so once both have
|
||||
// been processed the store can only hold the valid one).
|
||||
await { mirror.rejected.get() >= 1 && downstreamStore.count(Filter()) == 1 }
|
||||
|
||||
val stored = downstreamStore.query<Event>(Filter()).single()
|
||||
assertEquals(valid.id, stored.id)
|
||||
assertTrue(downstreamStore.query<Event>(Filter(ids = listOf(forged.id))).isEmpty())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,165 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.geode.mirror
|
||||
|
||||
import com.vitorpamplona.geode.InProcessRelays
|
||||
import com.vitorpamplona.geode.RelayEngine
|
||||
import com.vitorpamplona.geode.testing.publish
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebSocketListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.WebsocketBuilder
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import org.junit.After
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* The `trusted = true` skip-verify decision must be bound to the relay
|
||||
* that DELIVERED the event, not to the subscription id it arrived on.
|
||||
* The client pool dispatches EVENT messages by subscription id alone and
|
||||
* every `[[mirror]]` upstream shares one client — so a hostile untrusted
|
||||
* upstream that answers with the *trusted* upstream's subscription id
|
||||
* would otherwise ride its skip-verify straight into the store.
|
||||
*/
|
||||
class MirrorWorkerTrustOriginTest {
|
||||
private val trustedUrl = RelayUrlNormalizer.normalize("ws://trusted.relay/")
|
||||
private val hostileUrl = RelayUrlNormalizer.normalize("ws://hostile.relay/")
|
||||
private val downstreamUrl = RelayUrlNormalizer.normalize("ws://downstream.relay/")
|
||||
|
||||
private val hub = InProcessRelays()
|
||||
|
||||
private val downstreamStore = EventStore(null)
|
||||
private val downstream =
|
||||
RelayEngine(
|
||||
url = downstreamUrl,
|
||||
store = downstreamStore,
|
||||
parallelVerify = true,
|
||||
)
|
||||
|
||||
private var worker: MirrorWorker? = null
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
worker?.close()
|
||||
downstream.close()
|
||||
hub.close()
|
||||
}
|
||||
|
||||
private fun forgedEvent(idSeed: Int): Event =
|
||||
Event(
|
||||
id = idSeed.toString().padStart(64, '0'),
|
||||
pubKey = "1".repeat(64),
|
||||
createdAt = TimeUtils.now() - idSeed,
|
||||
kind = 1,
|
||||
tags = emptyArray(),
|
||||
content = "forged $idSeed",
|
||||
sig = "f".repeat(128),
|
||||
)
|
||||
|
||||
/**
|
||||
* A relay that never answers the REQ it was given; instead, on
|
||||
* connect it injects one EVENT frame carrying a subscription id that
|
||||
* belongs to a DIFFERENT upstream's subscription.
|
||||
*/
|
||||
private class HijackingWebSocket(
|
||||
private val out: WebSocketListener,
|
||||
private val frame: String,
|
||||
) : WebSocket {
|
||||
private var connected = false
|
||||
|
||||
override fun needsReconnect(): Boolean = !connected
|
||||
|
||||
override fun connect() {
|
||||
connected = true
|
||||
out.onOpen(0, false)
|
||||
out.onMessage(frame)
|
||||
}
|
||||
|
||||
override fun disconnect() {
|
||||
connected = false
|
||||
}
|
||||
|
||||
override fun send(msg: String): Boolean = true
|
||||
}
|
||||
|
||||
@Test
|
||||
fun hostileUpstreamCannotRideAnotherSubscriptionsTrust() =
|
||||
runBlocking {
|
||||
val forged = forgedEvent(1)
|
||||
// The trusted upstream is configured FIRST, so its down
|
||||
// subscription id is "geode-mirror-0" — which the hostile
|
||||
// relay claims in its injected frame.
|
||||
val hijackFrame = """["EVENT","geode-mirror-0",${forged.toJson()}]"""
|
||||
|
||||
val builder =
|
||||
object : WebsocketBuilder {
|
||||
override fun build(
|
||||
url: NormalizedRelayUrl,
|
||||
out: WebSocketListener,
|
||||
): WebSocket =
|
||||
if (url == hostileUrl) {
|
||||
HijackingWebSocket(out, hijackFrame)
|
||||
} else {
|
||||
hub.build(url, out)
|
||||
}
|
||||
}
|
||||
|
||||
val mirror =
|
||||
MirrorWorker(
|
||||
upstreams =
|
||||
listOf(
|
||||
MirrorUpstream(trustedUrl, trusted = true, backfillSeconds = 3600),
|
||||
MirrorUpstream(hostileUrl, trusted = false, backfillSeconds = 3600),
|
||||
),
|
||||
server = downstream.server,
|
||||
websocketBuilder = builder,
|
||||
).also { worker = it }
|
||||
mirror.start()
|
||||
|
||||
// The injected frame must be dropped at the origin check —
|
||||
// observable as a `filtered` tick, never as a stored row.
|
||||
withTimeout(15_000) {
|
||||
while (mirror.filtered.get() < 1) delay(25)
|
||||
}
|
||||
assertTrue(downstreamStore.query<Event>(Filter(ids = listOf(forged.id))).isEmpty())
|
||||
|
||||
// The genuinely trusted upstream still works on the same
|
||||
// subscription id the hostile relay tried to claim.
|
||||
val real = forgedEvent(2)
|
||||
hub.getOrCreate(trustedUrl).publish(real)
|
||||
withTimeout(15_000) {
|
||||
while (downstreamStore.count(Filter()) < 1) delay(25)
|
||||
}
|
||||
assertEquals(
|
||||
listOf(real.id),
|
||||
downstreamStore.query<Event>(Filter()).map { it.id },
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,324 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.geode.perf
|
||||
|
||||
import com.vitorpamplona.geode.KtorRelay
|
||||
import com.vitorpamplona.geode.RelayEngine
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.DefaultIndexingStrategy
|
||||
import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore
|
||||
import com.vitorpamplona.quartz.utils.EventFactory
|
||||
import io.ktor.server.application.install
|
||||
import io.ktor.server.cio.CIO
|
||||
import io.ktor.server.engine.embeddedServer
|
||||
import io.ktor.server.routing.routing
|
||||
import io.ktor.server.websocket.WebSockets
|
||||
import io.ktor.server.websocket.webSocket
|
||||
import io.ktor.websocket.Frame
|
||||
import io.ktor.websocket.readText
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import okhttp3.Response
|
||||
import okhttp3.WebSocket
|
||||
import okhttp3.WebSocketListener
|
||||
import java.util.concurrent.ArrayBlockingQueue
|
||||
import java.util.concurrent.TimeUnit
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Splits the wire-level small-REQ floor into transport vs relay work,
|
||||
* over the production stack (Ktor CIO server, OkHttp client, loopback).
|
||||
*
|
||||
* - echo-1 / echo-22: a bare Ktor CIO websocket replying with 1 / 22
|
||||
* ~250 B frames — the transport stack's round-trip floor; the burst
|
||||
* variants show per-frame cost is negligible.
|
||||
* - echo external / ext+1ms: replies produced from a foreign coroutine,
|
||||
* immediately or after the connection loop parks — proves cross-context
|
||||
* sends into CIO are prompt.
|
||||
* - geode REQ / NOTICE / empty REQ / inproc: the relay's own layers.
|
||||
*
|
||||
* Historical note — this benchmark found the TcpNoDelaySocketFactory
|
||||
* issue: without TCP_NODELAY on the client, a CLOSE (which relays never
|
||||
* answer) followed by a REQ nagles the REQ behind the unACKed CLOSE for
|
||||
* the ~40 ms delayed-ACK window, measured here as a flat 43.7 ms per
|
||||
* round. With the factory (now used by every production client), geode's
|
||||
* ~21-row REQ costs ~1.2 ms on the wire — matching relayBench — of which
|
||||
* ~0.6 ms is the transport floor and ~0.5 ms the per-REQ server work
|
||||
* documented in quartz/plans/2026-07-04-small-req-floor.md.
|
||||
*/
|
||||
class WireReqFloorBenchmark {
|
||||
companion object {
|
||||
const val EVENTS = 50_000
|
||||
const val AUTHORS = 2_500
|
||||
const val ROUNDS = 300
|
||||
const val WARMUP = 50
|
||||
const val BURST = 22
|
||||
}
|
||||
|
||||
private fun hexId(seed: Int): String = seed.toString(16).padStart(64, '0')
|
||||
|
||||
private fun pubkey(seed: Int): String = (seed % AUTHORS).toString(16).padStart(64, 'a')
|
||||
|
||||
private val sig = "0".repeat(128)
|
||||
|
||||
private fun event(seed: Int): Event =
|
||||
EventFactory.create(
|
||||
id = hexId(seed),
|
||||
pubKey = pubkey(seed),
|
||||
createdAt = 1_600_000_000L + (seed * 7919) % 1_000_000,
|
||||
kind = 1,
|
||||
tags = emptyArray(),
|
||||
content = "wire req floor benchmark $seed",
|
||||
sig = sig,
|
||||
)
|
||||
|
||||
private fun median(samples: LongArray): Double {
|
||||
samples.sort()
|
||||
return samples[samples.size / 2] / 1e6
|
||||
}
|
||||
|
||||
/** Blocking one-connection driver: send [request], await [expectedFrames] replies. */
|
||||
private class Driver(
|
||||
client: OkHttpClient,
|
||||
url: String,
|
||||
) {
|
||||
private val frames = ArrayBlockingQueue<String>(4096)
|
||||
val socket: WebSocket
|
||||
|
||||
init {
|
||||
val opened = ArrayBlockingQueue<Boolean>(1)
|
||||
socket =
|
||||
client.newWebSocket(
|
||||
Request.Builder().url(url).build(),
|
||||
object : WebSocketListener() {
|
||||
override fun onOpen(
|
||||
webSocket: WebSocket,
|
||||
response: Response,
|
||||
) {
|
||||
opened.put(true)
|
||||
}
|
||||
|
||||
override fun onMessage(
|
||||
webSocket: WebSocket,
|
||||
text: String,
|
||||
) {
|
||||
frames.put(text)
|
||||
}
|
||||
},
|
||||
)
|
||||
check(opened.poll(10, TimeUnit.SECONDS) == true) { "websocket did not open" }
|
||||
}
|
||||
|
||||
var lastFirstFrameNanos: Long = 0
|
||||
private set
|
||||
|
||||
fun roundTrip(
|
||||
request: String,
|
||||
expectedFrames: Int,
|
||||
): Long {
|
||||
val t0 = System.nanoTime()
|
||||
socket.send(request)
|
||||
repeat(expectedFrames) { i ->
|
||||
checkNotNull(frames.poll(10, TimeUnit.SECONDS)) { "timed out waiting for frame" }
|
||||
if (i == 0) lastFirstFrameNanos = System.nanoTime() - t0
|
||||
}
|
||||
return System.nanoTime() - t0
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun transportVsRelayFloor() =
|
||||
runBlocking {
|
||||
// TCP_NODELAY, via the same factory production clients use:
|
||||
// without it, a CLOSE (never answered) followed by a REQ nagles
|
||||
// the REQ behind the unACKed CLOSE for the ~40 ms delayed-ACK
|
||||
// window — this benchmark measured a flat 43.7 ms per round
|
||||
// before the factory existed, which is how the issue was found.
|
||||
val client = OkHttpClient.Builder().socketFactory(TcpNoDelaySocketFactory).build()
|
||||
|
||||
// --- bare Ktor CIO echo server: N frames of ~250 B per request ---
|
||||
val payload = "x".repeat(250)
|
||||
val echo =
|
||||
embeddedServer(CIO, port = 0) {
|
||||
install(WebSockets)
|
||||
routing {
|
||||
webSocket("/") {
|
||||
for (frame in incoming) {
|
||||
if (frame !is Frame.Text) continue
|
||||
val text = frame.readText()
|
||||
if (text.startsWith("x")) {
|
||||
// Reply from an EXTERNAL coroutine — the
|
||||
// cross-context path geode's launched REQ
|
||||
// handler uses.
|
||||
val n = text.drop(1).toIntOrNull() ?: 1
|
||||
launch(Dispatchers.Default) {
|
||||
repeat(n) { outgoing.send(Frame.Text(payload)) }
|
||||
}
|
||||
} else if (text.startsWith("d")) {
|
||||
// Same, but AFTER the connection loop has
|
||||
// gone idle — does a parked CIO connection
|
||||
// pick up a cross-thread send promptly?
|
||||
val n = text.drop(1).toIntOrNull() ?: 1
|
||||
launch(Dispatchers.Default) {
|
||||
delay(1)
|
||||
repeat(n) { outgoing.send(Frame.Text(payload)) }
|
||||
}
|
||||
} else {
|
||||
val n = text.toIntOrNull() ?: 1
|
||||
repeat(n) { outgoing.send(Frame.Text(payload)) }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
echo.start(wait = false)
|
||||
val echoPort =
|
||||
echo.engine
|
||||
.resolvedConnectors()
|
||||
.first()
|
||||
.port
|
||||
|
||||
val echoDriver = Driver(client, "ws://127.0.0.1:$echoPort/")
|
||||
repeat(WARMUP) { echoDriver.roundTrip("1", 1) }
|
||||
val echo1 = LongArray(ROUNDS) { echoDriver.roundTrip("1", 1) }
|
||||
repeat(WARMUP) { echoDriver.roundTrip("$BURST", BURST) }
|
||||
val echoN = LongArray(ROUNDS) { echoDriver.roundTrip("$BURST", BURST) }
|
||||
repeat(WARMUP) { echoDriver.roundTrip("x1", 1) }
|
||||
val echoExternal = LongArray(ROUNDS) { echoDriver.roundTrip("x1", 1) }
|
||||
repeat(WARMUP) { echoDriver.roundTrip("d1", 1) }
|
||||
val echoDelayed = LongArray(ROUNDS) { echoDriver.roundTrip("d1", 1) }
|
||||
|
||||
// --- geode: real REQ over the same stack ---
|
||||
// Same store setup SmallReqFloorBenchmark proved answers this
|
||||
// filter in ~0.18 ms in-process — this benchmark attributes
|
||||
// the wire-vs-in-process delta, so the storage side must be
|
||||
// the known-fast configuration.
|
||||
val store = EventStore(dbName = null, indexStrategy = DefaultIndexingStrategy(indexEventsByPubkeyAlone = true))
|
||||
(1..EVENTS).chunked(2000).forEach { chunk -> store.batchInsert(chunk.map { event(it) }) }
|
||||
val relay = RelayEngine(url = "ws://127.0.0.1:7795/".normalizeRelayUrl(), store = store, parentContext = Dispatchers.IO + SupervisorJob())
|
||||
val server = KtorRelay(relay, host = "127.0.0.1", port = 0).start()
|
||||
|
||||
val geodeDriver = Driver(client, server.url)
|
||||
|
||||
fun req(round: Int) = """["REQ","w$round",{"authors":["${pubkey(round)}"],"kinds":[1],"limit":50}]"""
|
||||
|
||||
// Each REQ answers with (rows + 1) frames — the EVENTs then
|
||||
// EOSE — and the row count is author-dependent, so resolve
|
||||
// every round's expected count from the store up front. Subs
|
||||
// are CLOSEd after each round (silently, per NIP-01) so the
|
||||
// relay's live-subscription registry doesn't grow with the
|
||||
// round count and skew later samples.
|
||||
suspend fun rowsFor(round: Int): Int = store.query<Event>(Filter(authors = listOf(pubkey(round)), kinds = listOf(1), limit = 50)).size
|
||||
|
||||
repeat(WARMUP) { round ->
|
||||
val rows = rowsFor(round)
|
||||
geodeDriver.roundTrip(req(round), rows + 1)
|
||||
geodeDriver.socket.send("""["CLOSE","w$round"]""")
|
||||
}
|
||||
val geodeSamples = LongArray(ROUNDS)
|
||||
val geodeFirst = LongArray(ROUNDS)
|
||||
var totalRows = 0
|
||||
for (i in 0 until ROUNDS) {
|
||||
val round = WARMUP + i
|
||||
val rows = rowsFor(round)
|
||||
totalRows += rows
|
||||
geodeSamples[i] = geodeDriver.roundTrip(req(round), rows + 1)
|
||||
geodeFirst[i] = geodeDriver.lastFirstFrameNanos
|
||||
geodeDriver.socket.send("""["CLOSE","w$round"]""")
|
||||
}
|
||||
|
||||
// Probe A: malformed frame -> NOTICE, answered inline on the
|
||||
// pump coroutine (no launch, no SQL). Probe B: REQ matching
|
||||
// nothing -> lone EOSE (launch + SQL, no rows).
|
||||
repeat(WARMUP) { geodeDriver.roundTrip("not json", 1) }
|
||||
val notice = LongArray(ROUNDS) { geodeDriver.roundTrip("not json", 1) }
|
||||
repeat(WARMUP) { i ->
|
||||
geodeDriver.roundTrip("""["REQ","n$i",{"ids":["${"f".repeat(64)}"]}]""", 1)
|
||||
geodeDriver.socket.send("""["CLOSE","n$i"]""")
|
||||
}
|
||||
val emptyReq =
|
||||
LongArray(ROUNDS) { i ->
|
||||
val t = geodeDriver.roundTrip("""["REQ","m$i",{"ids":["${"f".repeat(64)}"]}]""", 1)
|
||||
geodeDriver.socket.send("""["CLOSE","m$i"]""")
|
||||
t
|
||||
}
|
||||
// Same probe with NO CLOSE between rounds: if this is fast, the
|
||||
// 44 ms is the client's CLOSE frame sitting unACKed (server
|
||||
// sends nothing for a CLOSE) and Nagle holding the next REQ
|
||||
// behind it until the ~40 ms delayed ACK — a client-side TCP
|
||||
// artifact, not the relay.
|
||||
repeat(WARMUP) { i -> geodeDriver.roundTrip("""["REQ","nc$i",{"ids":["${"a".repeat(64)}"]}]""", 1) }
|
||||
val emptyNoClose =
|
||||
LongArray(ROUNDS) { i ->
|
||||
geodeDriver.roundTrip("""["REQ","ncm$i",{"ids":["${"a".repeat(64)}"]}]""", 1)
|
||||
}
|
||||
|
||||
// Probe C: same RelayEngine, no wire — an in-process session
|
||||
// while Ktor keeps serving. Splits engine-state issues from
|
||||
// transport-adjacent ones.
|
||||
val inproc = LongArray(ROUNDS)
|
||||
run {
|
||||
val q = ArrayBlockingQueue<String>(4096)
|
||||
val session = relay.server.connect { q.put(it) }
|
||||
repeat(WARMUP) { i ->
|
||||
session.receive("""["REQ","p$i",{"ids":["${"e".repeat(64)}"]}]""")
|
||||
checkNotNull(q.poll(10, TimeUnit.SECONDS))
|
||||
session.receive("""["CLOSE","p$i"]""")
|
||||
}
|
||||
for (i in 0 until ROUNDS) {
|
||||
val t0 = System.nanoTime()
|
||||
session.receive("""["REQ","q$i",{"ids":["${"e".repeat(64)}"]}]""")
|
||||
checkNotNull(q.poll(10, TimeUnit.SECONDS))
|
||||
inproc[i] = System.nanoTime() - t0
|
||||
session.receive("""["CLOSE","q$i"]""")
|
||||
}
|
||||
session.close()
|
||||
}
|
||||
|
||||
assertTrue(totalRows > 0)
|
||||
println("WireReqFloorBenchmark (loopback, OkHttp client) @ ${EVENTS / 1000}k events, medians of $ROUNDS")
|
||||
println(" echo 1 frame: ${"%6.3f".format(median(echo1))} ms")
|
||||
println(" echo $BURST frames: ${"%6.3f".format(median(echoN))} ms")
|
||||
println(" echo 1 frame (external): ${"%6.3f".format(median(echoExternal))} ms")
|
||||
println(" echo 1 frame (ext+1ms): ${"%6.3f".format(median(echoDelayed))} ms")
|
||||
println(" geode REQ (~${totalRows / ROUNDS} rows+EOSE): ${"%6.3f".format(median(geodeSamples))} ms (first frame ${"%6.3f".format(median(geodeFirst))} ms)")
|
||||
println(" geode NOTICE (inline): ${"%6.3f".format(median(notice))} ms")
|
||||
println(" geode empty REQ (launch): ${"%6.3f".format(median(emptyReq))} ms")
|
||||
println(" geode empty REQ no CLOSE: ${"%6.3f".format(median(emptyNoClose))} ms")
|
||||
println(" geode empty REQ (inproc): ${"%6.3f".format(median(inproc))} ms")
|
||||
|
||||
geodeDriver.socket.close(1000, null)
|
||||
echoDriver.socket.close(1000, null)
|
||||
server.stop(0, 1_000)
|
||||
relay.close()
|
||||
echo.stop(0, 500)
|
||||
client.dispatcher.executorService.shutdown()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,130 @@
|
||||
# Incremental live negentropy storage
|
||||
|
||||
**Status: shipped** — backlog item 3 of the relay performance campaign
|
||||
(PR #3466 follow-up). All four milestones landed on this branch.
|
||||
|
||||
**Measured results.** Micro (LiveNegentropyBenchmark, 50k events,
|
||||
in-container): scan+seal cold path 80–100 ms per open; index post-write
|
||||
open 9–16 ms (~5–10×); first open ~70 ms (pays the one-time rebuild).
|
||||
relayBench A/B (both variants in ONE run — geode with index vs geode
|
||||
with `[negentropy].live_index = false` — 50k corpus, then repeated with
|
||||
the relay order reversed to control for run-order bias):
|
||||
|
||||
- **identical-set reconcile: 56/57 ms with the index vs 110/130 ms
|
||||
without, in both orders — the real ~2.2× win.** This phase runs
|
||||
after the delta transfer wrote events, so the old single-slot cache
|
||||
always misses — the exact gap the index closes.
|
||||
- initial (first-ever) reconcile and ingest throughput: whichever
|
||||
relay ran first won/lost respectively in BOTH runs — pure run-order
|
||||
effect, i.e. no ingest regression and (as designed) no win on the
|
||||
very first open, which pays the lazy rebuild.
|
||||
|
||||
Measurement note for future A/Bs: putting both variants in one
|
||||
relayBench run controls container noise better than alternating runs,
|
||||
but per-phase run-order bias is real (~100 ms on initial reconcile,
|
||||
~10% on ingest) — always repeat with the order reversed.
|
||||
|
||||
**strfry head-to-head** (same container, 50k corpus, `--geode-no-search`,
|
||||
strfry built from source, single run): identical-set reconcile geode
|
||||
41.4 ms vs strfry 30.1 ms — the always-current-tree gap is down to
|
||||
~1.4× from the campaign-opening "full scan + seal per open". Initial
|
||||
(cold) reconcile 177 vs 112 ms: strfry keeps its tree across the whole
|
||||
run while geode's first open pays the lazy rebuild. Ingest measured at
|
||||
parity in this container (7,972 vs 7,860 ev/s — treat as ±noise, the
|
||||
earlier 1.59× gap was measured on different hardware); storage 72.8 vs
|
||||
106.0 MiB.
|
||||
|
||||
**Possible follow-up** if the residual matters: the per-open cost is now
|
||||
the O(n) copy + re-seal into a `StorageVector` (~16 ms at 50k). A custom
|
||||
immutable `IStorage` view over the live entries (copy-on-write chunks,
|
||||
no re-seal) would make an open O(1), matching strfry's zero-copy read of
|
||||
its tree — that is the "interesting part" the original backlog item
|
||||
anticipated, deferred until a benchmark says the remaining ~11 ms is
|
||||
worth it.
|
||||
|
||||
## Problem
|
||||
|
||||
A cold NEG-OPEN pays a full scan + O(n log n) seal: `snapshotIdsForNegentropy`
|
||||
(SQL scan of every matching row) → `NegentropyServerSession.sealVector`
|
||||
(sort + seal). relayBench measured ~340 ms per reconcile at 50k events before
|
||||
the single-slot snapshot cache; strfry answers ~21 ms off its always-current
|
||||
in-memory view.
|
||||
|
||||
The cache added in #3466 (keyed on filter JSON + write generation + 30 s TTL)
|
||||
only helps the *identical-filter, zero-writes-in-between* repeat. On a relay
|
||||
ingesting continuously the generation moves constantly, so mirror heartbeats
|
||||
are effectively always cold; any new filter is cold by definition.
|
||||
|
||||
## Design
|
||||
|
||||
### 1. `LiveNegentropyIndex` (quartz, server-only, opt-in)
|
||||
|
||||
An always-current sorted set of `(createdAt, id₃₂)` maintained from the
|
||||
store's write path — strfry's `MemoryView` equivalent, ~140 B/entry on
|
||||
the JVM (`IdAndTime` keeps the id as a 64-char hex string; 1M events ≈
|
||||
140 MB; the index is only built when the corpus fits
|
||||
`negentropy.max_sync_events`, so that also caps the heap).
|
||||
|
||||
- **Structure**: single sorted array with binary-search insert. Nostr inserts
|
||||
are near-tail (created_at ≈ now), so the memmove is tiny in the common
|
||||
case; measure the out-of-order (backfill) worst case and only move to a
|
||||
chunked layout if it shows.
|
||||
- **Snapshot**: NEG-OPEN copies the array into a sealed `IStorage` — O(n)
|
||||
arraycopy of already-sorted data (~2 MB at 50k, sub-ms), reusing the
|
||||
existing single-slot cache so back-to-back opens share one snapshot.
|
||||
Reconcile only reads (`size/getItem/iterate/indexAtOrBeforeBound`), so one
|
||||
sealed snapshot serves any number of concurrent sessions.
|
||||
- **Serves index-total filters only**: no `ids/authors/kinds/tags/search`
|
||||
constraints. `since/until` ARE served — the structure is time-sorted, so a
|
||||
time window is an index sub-range. Constrained filters keep the scan+seal
|
||||
path (+ cache). The mirror-heartbeat pattern this optimizes is a broad
|
||||
time-window filter, so this covers the case that matters.
|
||||
|
||||
### 2. Removal correctness (the interesting part)
|
||||
|
||||
The index must never advertise ids the store no longer has, or peers fetch
|
||||
dead ids. Removal paths differ in frequency and get different treatment:
|
||||
|
||||
- **Replaceable/addressable overwrite** (frequent — every kind 0/3/1xxxx/
|
||||
3xxxx update): `ReplaceableModule`/`AddressableModule` run
|
||||
`DELETE FROM event_headers WHERE …` inside the insert transaction. Add
|
||||
`RETURNING created_at, id` (bundled SQLite ≥ 3.35) and report displaced
|
||||
rows to the index alongside the insert.
|
||||
- **Wholesale/rare paths** (kind-5 NIP-09, expiration sweep, right-to-vanish,
|
||||
NIP-86 `delete(filter)`, FTS reindex/clear): invalidate the whole index;
|
||||
the next NEG-OPEN rebuilds it lazily from one scan and incremental
|
||||
maintenance resumes. Deletes are rare enough that occasional rebuilds beat
|
||||
threading deltas through every module.
|
||||
|
||||
### 3. Gating
|
||||
|
||||
`IndexingStrategy.maintainLiveNegentropyIndex`, default **false** — library
|
||||
defaults unchanged for app-side stores (campaign ground rule). geode's
|
||||
`RelayIndexingStrategy` turns it on; config kill-switch under
|
||||
`[negentropy]`.
|
||||
|
||||
### 4. Concurrency
|
||||
|
||||
Mutations happen only on the writer path (single-writer mutex — same
|
||||
discipline as the FTS worker). Snapshots swap in via copy-on-write so a
|
||||
reconcile never observes a mid-insert array. Shutdown: the index is memory
|
||||
only, rebuilt on boot from the first NEG-OPEN's scan; no lifecycle beyond
|
||||
the store's own close (ground rule 4: no worker, no uncaught exceptions).
|
||||
|
||||
## Measurement plan
|
||||
|
||||
1. Micro: quartz jvmTest benchmark, cold NEG-OPEN time at 50k events,
|
||||
before/after (expect ~340 ms → single-digit ms).
|
||||
2. Headline: relayBench pairwise sync, alternating A/B runs on the 50k
|
||||
corpus (container noise ±10–30%; never trust a single run). Keep only
|
||||
if it wins end-to-end; document either way.
|
||||
|
||||
## Milestones
|
||||
|
||||
1. `LiveNegentropyIndex` + unit tests (ordering, tail/backfill inserts,
|
||||
removal, snapshot immutability under concurrent insert, cap behavior).
|
||||
2. Displaced-row `RETURNING` plumbing in Replaceable/Addressable modules +
|
||||
wholesale-invalidation hooks on the rare paths.
|
||||
3. `LiveEventStore.sealedNegentropyStorage` wiring: index-total filters
|
||||
from the index; everything else keeps scan+seal+cache.
|
||||
4. Benchmarks (micro then relayBench A/B); revert if not a real win.
|
||||
@@ -0,0 +1,77 @@
|
||||
# Small-REQ dispatch floor — investigated, inline fast path reverted
|
||||
|
||||
**Status: closed (negative result recorded).** Backlog item 2 of the
|
||||
relay performance campaign.
|
||||
|
||||
## The gap
|
||||
|
||||
relayBench at 50k events: geode WINS most 500-event query scenarios
|
||||
(hashtag 5.4 vs 8.4 ms, recent-window 4.1 vs 8.5) but loses ~2.5× on
|
||||
small results — author-archive (19 events) 1.2–1.7 ms vs strfry's
|
||||
~0.5–0.6, thread (23) likewise — and the @8conn throughput inverts
|
||||
(strfry 2–3× geode). With ~20-row responses, throughput ≈ 1/latency:
|
||||
there is a fixed per-REQ floor.
|
||||
|
||||
## Decomposition (SmallReqFloorBenchmark, kept in jvmTest)
|
||||
|
||||
In-process at 50k events, ~21 rows/REQ, medians of 400:
|
||||
|
||||
| stage | ms |
|
||||
|---|---:|
|
||||
| A raw store query (SQL + row decode) | 0.18 |
|
||||
| B + live machinery (FilterIndex reg/unreg, dedupe set) | 0.36 |
|
||||
| C + session dispatch (parse, launch, frames) | 0.60 |
|
||||
|
||||
Note: in-memory DBs have no reader pool (`useReader` falls back to the
|
||||
writer mutex), so absolute numbers are conservative vs the file-DB
|
||||
bench setup.
|
||||
|
||||
## What was tried and why it was reverted
|
||||
|
||||
An inline fast path (`SessionBackend.queryRawInline`): REQs with
|
||||
provably bounded replays (limit or ids-count summing ≤ 512) ran their
|
||||
stored replay on the receive coroutine and kept only a live-tail
|
||||
handle — no per-REQ `launch`, no Job, no dispatcher handoffs. It cut
|
||||
in-process time-to-EOSE ~17% (0.60 → 0.50 ms) with full wire-behavior
|
||||
parity (stored→EOSE order, live tail, CLOSE, same-subId replacement).
|
||||
|
||||
Three relayBench runs (baseline, cap-256 [path not engaged — bench
|
||||
filters carry limit=500 or no limit], cap-512 [engaged for
|
||||
author-archive/by-ids/500-limit feeds]) showed **no movement outside
|
||||
the container drift band** — strfry's own numbers drifted ±30% run to
|
||||
run, and inline-eligible scenarios moved the same as ineligible ones.
|
||||
Reverted per the keep-only-winners rule.
|
||||
|
||||
## Where the floor actually is (corrected after WireReqFloorBenchmark)
|
||||
|
||||
The follow-up wire benchmark (geode's `WireReqFloorBenchmark`, Ktor CIO
|
||||
+ OkHttp on loopback) attributed the full path:
|
||||
|
||||
| leg | ms |
|
||||
|---|---:|
|
||||
| bare Ktor CIO echo round trip (1 or 22 frames — same) | ~0.9–1.0 |
|
||||
| geode NOTICE (inline, full pump + Ktor send) | ~0.5 |
|
||||
| geode empty REQ (launch + SQL, 0 rows) | ~0.6–0.8 |
|
||||
| geode ~21-row REQ, wire | ~1.25 (= relayBench's number) |
|
||||
|
||||
**geode's websocket send path has no latency problem** — per-frame burst
|
||||
cost is negligible (echo-22 ≈ echo-1), the pump adds ~nothing (NOTICE ≈
|
||||
0.5 ms), and the residual vs strfry (~0.5 ms/REQ) is the per-REQ server
|
||||
work already investigated above. Frame batching / permessage-deflate
|
||||
would not move these numbers. Backlog item 6's remaining open angle is
|
||||
the INGEST-side CPU share (13–25% in the JFR profile) — a throughput
|
||||
question, not this latency one.
|
||||
|
||||
**The real find was client-side.** The first wire measurements showed a
|
||||
flat 43.7 ms per REQ — which turned out to be the benchmark's own OkHttp
|
||||
client: OkHttp does not set TCP_NODELAY, and the CLOSE-then-REQ pattern
|
||||
(every feed/filter switch!) nagles the REQ behind the unACKed CLOSE
|
||||
(relays never answer CLOSE) for the ~40 ms delayed-ACK window.
|
||||
relayBench's harness client already carried a no-delay socket factory —
|
||||
which is why bench numbers never showed it — but the production clients
|
||||
(Android relay pool, Desktop, amy, geode's mirror) did not. Fixed by
|
||||
`TcpNoDelaySocketFactory` (quartz jvmAndroid), now used by all of them.
|
||||
|
||||
**Do not retry** relay-side latency work for the small-REQ gap; the
|
||||
addressable remainder is the ~0.4 ms of per-REQ dispatch machinery this
|
||||
doc's revert already covers, and it does not show on the wire.
|
||||
@@ -1,12 +1,14 @@
|
||||
# quartz plans
|
||||
|
||||
_Audited 2026-06-30. 9 plans: 7 shipped (archived), 0 in-progress, 2 queued, 0 abandoned._
|
||||
_Audited 2026-06-30. 11 plans: 7 shipped (archived), 0 in-progress, 3 queued, 1 closed (negative result)._
|
||||
|
||||
## Queued
|
||||
| Plan | Summary |
|
||||
| ---- | ------- |
|
||||
| [2026-05-08-local-headers-explorer.md](2026-05-08-local-headers-explorer.md) | Headers-only Bitcoin P2P client to verify NIP-03 OTS attestations without a trusted block explorer. |
|
||||
| [2026-06-12-giftwrap-deletion-requests.md](2026-06-12-giftwrap-deletion-requests.md) | Let a recipient-authored kind-5 delete/block a gift wrap (kind 1059) addressed to them. |
|
||||
| [2026-07-03-incremental-negentropy-storage.md](2026-07-03-incremental-negentropy-storage.md) | Always-current (created_at, id) index so cold NEG-OPENs stop paying a full scan + seal (~340 ms at 50k vs strfry's ~21 ms). |
|
||||
| [2026-07-04-small-req-floor.md](2026-07-04-small-req-floor.md) | Small-REQ dispatch floor: decomposed, inline fast path tried and reverted (no wire-level win); floor is transport-side. |
|
||||
|
||||
## Archived (shipped)
|
||||
| Plan | Summary |
|
||||
|
||||
+37
-2
@@ -20,6 +20,7 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.server
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.verify
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.IngestQueue
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.server.backend.LiveEventStore
|
||||
@@ -64,7 +65,7 @@ class NostrServer(
|
||||
private val store: IEventStore,
|
||||
policyBuilder: () -> IRelayPolicy = { VerifyPolicy },
|
||||
parentContext: CoroutineContext = SupervisorJob(),
|
||||
parallelVerify: Boolean = false,
|
||||
private val parallelVerify: Boolean = false,
|
||||
negentropySettings: NegentropySettings = NegentropySettings.Default,
|
||||
listener: RelayServerListener = RelayServerListener.None,
|
||||
limits: RelayLimits? = null,
|
||||
@@ -95,7 +96,41 @@ class NostrServer(
|
||||
},
|
||||
)
|
||||
|
||||
override val backend: SessionBackend = LiveEventStore(store, ingest)
|
||||
private val liveStore = LiveEventStore(store, ingest)
|
||||
|
||||
override val backend: SessionBackend = liveStore
|
||||
|
||||
/**
|
||||
* Local ingestion path for events that did not arrive over a client
|
||||
* connection — e.g. a mirror worker streaming a trusted upstream
|
||||
* relay, or an import job. Routes through the same group-commit
|
||||
* [IngestQueue] and live fanout as a client EVENT publish, but skips
|
||||
* the **entire** per-connection policy chain (there is no
|
||||
* connection): no [VerifyPolicy], no allow/deny lists, no size
|
||||
* limits. Callers own that screening — scope what may enter this
|
||||
* path (e.g. geode's per-upstream mirror filters) accordingly.
|
||||
*
|
||||
* [skipVerify] exempts the event from signature verification — the
|
||||
* relay-to-relay trust model: pass `true` only for events from an
|
||||
* explicitly configured upstream that already verified them (Schnorr
|
||||
* verify profiles at ~8% of busy ingest CPU). The default `false`
|
||||
* keeps verify-everything semantics regardless of configuration:
|
||||
* when the [IngestQueue] hook is on ([parallelVerify]) it verifies
|
||||
* there, otherwise this method verifies inline — without this, a
|
||||
* server whose verification lives in the (bypassed) policy chain
|
||||
* would silently ingest forgeries.
|
||||
*/
|
||||
suspend fun ingest(
|
||||
event: Event,
|
||||
skipVerify: Boolean = false,
|
||||
onComplete: (IEventStore.InsertOutcome) -> Unit,
|
||||
) {
|
||||
if (!skipVerify && !parallelVerify && !event.verify()) {
|
||||
onComplete(IEventStore.InsertOutcome.Rejected("invalid: bad signature or id"))
|
||||
return
|
||||
}
|
||||
liveStore.submit(event, skipVerify, onComplete)
|
||||
}
|
||||
|
||||
init {
|
||||
// Deferred-FTS catch-up worker: tokenizes in the gaps between
|
||||
|
||||
+23
-6
@@ -27,7 +27,6 @@ import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.awaitAll
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.ClosedReceiveChannelException
|
||||
@@ -115,9 +114,15 @@ class IngestQueue(
|
||||
/**
|
||||
* One outstanding ingest request: the event to insert plus the
|
||||
* callback the writer fires once the row's outcome is known.
|
||||
* [skipVerify] exempts this row from the [verify] hook — the
|
||||
* relay-to-relay trust model: set by local ingestion paths for
|
||||
* events streamed from an explicitly configured upstream relay
|
||||
* that already verified them (see
|
||||
* [com.vitorpamplona.quartz.nip01Core.relay.server.NostrServer.ingest]).
|
||||
*/
|
||||
class Submission(
|
||||
val event: Event,
|
||||
val skipVerify: Boolean,
|
||||
val onComplete: (IEventStore.InsertOutcome) -> Unit,
|
||||
)
|
||||
|
||||
@@ -159,11 +164,12 @@ class IngestQueue(
|
||||
*/
|
||||
suspend fun submit(
|
||||
event: Event,
|
||||
skipVerify: Boolean = false,
|
||||
onComplete: (IEventStore.InsertOutcome) -> Unit,
|
||||
) {
|
||||
ensureWriterStarted()
|
||||
pending.addAndFetch(1)
|
||||
incoming.send(Submission(event, onComplete))
|
||||
incoming.send(Submission(event, skipVerify, onComplete))
|
||||
}
|
||||
|
||||
private fun ensureWriterStarted() {
|
||||
@@ -234,15 +240,26 @@ class IngestQueue(
|
||||
* configured (skip the stage entirely). For multi-event batches
|
||||
* each verify runs as its own `async(Default)` so they spread
|
||||
* across CPU cores; single-event batches short-circuit to a
|
||||
* direct call to avoid coroutine-scope overhead.
|
||||
* direct call to avoid coroutine-scope overhead. Rows flagged
|
||||
* [Submission.skipVerify] (trusted publishers) pass without
|
||||
* invoking the hook.
|
||||
*/
|
||||
private suspend fun verifyBatch(batch: List<Submission>): BooleanArray? {
|
||||
val hook = verify ?: return null
|
||||
if (batch.size == 1) return BooleanArray(1) { hook(batch[0].event) }
|
||||
if (batch.size == 1) {
|
||||
val sub = batch[0]
|
||||
return BooleanArray(1) { sub.skipVerify || hook(sub.event) }
|
||||
}
|
||||
if (batch.all { it.skipVerify }) return BooleanArray(batch.size) { true }
|
||||
return coroutineScope {
|
||||
batch
|
||||
.map { sub -> async(Dispatchers.Default) { hook(sub.event) } }
|
||||
.awaitAll()
|
||||
.map { sub ->
|
||||
if (sub.skipVerify) {
|
||||
null
|
||||
} else {
|
||||
async(Dispatchers.Default) { hook(sub.event) }
|
||||
}
|
||||
}.map { it?.await() ?: true }
|
||||
.toBooleanArray()
|
||||
}
|
||||
}
|
||||
|
||||
+25
-1
@@ -91,8 +91,23 @@ class LiveEventStore(
|
||||
override suspend fun submit(
|
||||
event: Event,
|
||||
onComplete: (IEventStore.InsertOutcome) -> Unit,
|
||||
) = submit(event, skipVerify = false, onComplete = onComplete)
|
||||
|
||||
/**
|
||||
* [submit] variant for locally-originated traffic (mirror/sync
|
||||
* workers rather than client connections). [skipVerify] exempts
|
||||
* this event from the [IngestQueue]'s signature-verification hook —
|
||||
* the relay-to-relay trust model: set it only for events streamed
|
||||
* from a configured upstream relay that already verified them.
|
||||
* Accepted events fan out to live subscribers exactly like a
|
||||
* client publish.
|
||||
*/
|
||||
suspend fun submit(
|
||||
event: Event,
|
||||
skipVerify: Boolean,
|
||||
onComplete: (IEventStore.InsertOutcome) -> Unit,
|
||||
) {
|
||||
ingest.submit(event) { outcome ->
|
||||
ingest.submit(event, skipVerify) { outcome ->
|
||||
if (outcome is IEventStore.InsertOutcome.Accepted) {
|
||||
writeGeneration.addAndFetch(1L)
|
||||
fanout(event)
|
||||
@@ -366,6 +381,15 @@ class LiveEventStore(
|
||||
filters: List<Filter>,
|
||||
maxEntries: Int,
|
||||
): IStorage? {
|
||||
// Full-set NEG-OPENs — a single unconstrained filter, the shape
|
||||
// relay-relay sync sends by default — are served from the store's
|
||||
// always-current index when it maintains one: no scan, no
|
||||
// O(n log n) seal, cold or not. `null` falls through to the scan
|
||||
// path, which also owns the over-cap NEG-ERR detection.
|
||||
if (filters.size == 1 && filters[0].isEmpty()) {
|
||||
store.liveNegentropySnapshot(maxEntries)?.let { return it }
|
||||
}
|
||||
|
||||
val generation = writeGeneration.load()
|
||||
val key = filters.joinToString(" | ||||