From ab71e45e7744d4c836a20b77a91f4ee9706fe73e Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 7 May 2026 22:16:32 +0000 Subject: [PATCH 1/4] docs(geode): refine negentropy plan for strfry interop Align the NIP-77 large-corpus plan with strfry's actual defaults: 500_000-byte frame cap (not 64 KB), no default since-window, 200 shared sub cap, max_sync_events overflow protection, exact NEG-ERR wording, and 32-byte id storage. Add a strfry round-trip benchmark and call out the BTreeLMDB pre-built tree as a follow-up. --- .../2026-05-07-negentropy-large-corpus.md | 264 ++++++++++++++---- 1 file changed, 204 insertions(+), 60 deletions(-) diff --git a/geode/plans/2026-05-07-negentropy-large-corpus.md b/geode/plans/2026-05-07-negentropy-large-corpus.md index 0cb31ccd9b..55d057d015 100644 --- a/geode/plans/2026-05-07-negentropy-large-corpus.md +++ b/geode/plans/2026-05-07-negentropy-large-corpus.md @@ -1,20 +1,20 @@ -# NIP-77 negentropy at scale: snapshot memory + chunked replay +# NIP-77 negentropy at scale: strfry-interop snapshot path ## Problem `RelaySession` delegates NEG-OPEN to `NegSessionRegistry.open` -(`quartz/nip01Core/relay/server/NegSessionRegistry.kt`), which calls -`store.snapshotQuery(filters)` and feeds the **entire** result list -into `NegentropyServerSession`. For a relay holding 5 M events that -match a broad NEG-OPEN filter (`{kinds: [1, 7]}`), this is 5 M -`Event` objects materialised in memory before the first NEG-MSG goes -out. +(`quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt:58-77`), +which calls `store.snapshotQuery(filters)` and feeds the **entire** +result list of full `Event` objects into `NegentropyServerSession` +(`quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropyServerSession.kt:40-54`). +For a relay holding 5 M events that match a broad NEG-OPEN filter +(`{kinds: [1, 7]}`), this materialises 5 M `Event` objects with +content/tags/sig — call it ~1 KB/event, ~5 GB transient pressure per +concurrent NEG-OPEN — before the first NEG-MSG goes out. -The negentropy library itself is fine — it pivots into a sealed -`StorageVector` (id + createdAt only, ~40 bytes/entry). But the -`store.query(f)` step that produces the input materialises full -`Event` objects with content, tags, sig — call it ~1 KB/event. 5 M × -1 KB = 5 GB transient pressure per concurrent NEG-OPEN. +The negentropy library itself is fine. Internally it pivots into a +sealed `StorageVector` of `(uint64 timestamp, byte[32] id)` items — +40 bytes/entry. The waste is purely in the snapshot step. Two operator-visible symptoms: @@ -24,77 +24,221 @@ Two operator-visible symptoms: large stores the client waits seconds for what should be a millisecond round-trip. +## Reference: strfry + +We want byte-for-byte interop and comparable throughput against +[hoytech/strfry](https://github.com/hoytech/strfry). The relevant +defaults from `strfry/src/apps/relay/RelayNegentropy.cpp` and +`golpe.yaml`: + +| Knob | strfry default | Notes | +|---------------------------------------|-----------------------|-------| +| `frameSizeLimit` (NEG-MSG payload) | **500_000 bytes** | Hard-coded `Negentropy ne(storage, 500'000)` | +| `relay__negentropy__maxSyncEvents` | **1_000_000** | Hard cap on items in the snapshot | +| `relay__maxSubsPerConnection` | **200** | Shared between REQ and NEG sessions | +| `idSize` on the wire | **32 bytes** | NIP-77 v1 (`PROTOCOL_VERSION = 0x61`) | +| Default `since` window | **none** | Filter is honored as-is | +| Filter parser | same as REQ | Honors `limit`, `kinds`, `authors`, `#tags` | +| Snapshot data | `(created_at, id)` only | LMDB scan inserts into `negentropy::storage::Vector` | + +Two-phase materialisation in strfry: `QueryScheduler` scans LMDB +asynchronously, batching matched level-ids into a `vector`; +on completion the worker pulls each event header, calls +`storageVector.insert(packed.created_at(), packed.id())`, then +`seal()`s. The session response is sent only after seal. + +Our equivalent must (a) match the snapshot footprint of ~40 bytes per +event, and (b) match the wire-level frame-cap so a single +reconciliation round-trip exchanges the same payload size. + ## Sketch -### A — id-and-time-only snapshot path +### A — id-and-time-only snapshot path (memory parity) Negentropy only needs `(createdAt, id)` pairs. Add a streaming -`IEventStore.queryIdAndTime(filter)` that returns -`Sequence>` (or a `Flow` of small chunks) — -no content/tags/sig, no Event allocation. SQLite path is a SELECT -on `event_headers` (the `created_at`, `id` columns are already -indexed for query plans). +`IEventStore.snapshotIdsForNegentropy(filter)` that returns +`Sequence` (`data class IdAndTime(val createdAt: Long, val id: ByteArray)`) +— no content/tags/sig, no `Event` allocation. The SQLite path is a +plain `SELECT id, created_at FROM event_headers WHERE …` against the +existing `query_by_created_at_id` index +(`quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventIndexesModule.kt:79`). ```kotlin -suspend fun snapshotIdsForNegentropy(filter: Filter): IdTimeStream +// IEventStore +suspend fun snapshotIdsForNegentropy(filters: List): List + +// LiveEventStore — already deduplicates union of multi-filter results +suspend fun snapshotIdsForNegentropy(filters: List): List ``` -`NegentropyServerSession` is rewritten to take that stream and feed -it directly into the `StorageVector`. Memory drops from O(N × 1 KB) -to O(N × 40 B) — a 25× reduction; for 5 M events, ~200 MB instead -of 5 GB. +`NegentropyServerSession` is rewritten to take that list (or a +two-phase `pendingIds → seal` builder) and feed it directly into +`StorageVector`. Memory drops from O(N × ~1 KB) to O(N × 40 B) — for +1 M events, ~40 MB instead of ~1 GB; matches strfry's per-session +footprint. -### B — bounded-window subscriptions +**id encoding:** strfry stores 32-byte raw ids (NIP-77 v1 +`ID_SIZE = 32`); the 16-byte truncation is `FINGERPRINT_SIZE`, only used +internally for SHA-256 accumulator output. The current Kotlin code +already passes the hex `event.id` string to `storage.insert(..)` +(`NegentropyServerSession.kt:50`); the kmp-negentropy library decodes +that to 32 raw bytes. Confirm this is preserved when we switch to a +`ByteArray` id input — do not pre-truncate. -Most NEG-OPENs from real Nostr clients want the last 30 days, not -"everything." If the client doesn't supply `since`, the server can -default to a configurable horizon (e.g. 90 days) and surface this in -the NIP-11 `limitation.negentropy_max_lookback_seconds` field. -Operators can lift the cap; clients reading the doc know the bound. +### B — bounded snapshot size, NOT a default `since` window -This is a NIP-spec-adjacent question more than a code change — needs -a comment on whether the spec allows it. nostr-rs-relay does this -already. +The previous draft of this plan suggested defaulting `since` to a 90-day +horizon. **Drop that.** strfry doesn't do it — the filter is honored +as-is — and silently bounding `since` would break interop with strfry's +sync clients (e.g. `strfry sync`, nostr-sdk's negentropy reconciler): +they ask for "everything" and rely on getting it. -### C — frame-size cap on NEG-MSG +Instead match strfry's protection: a hard cap on the number of items +that go into a single snapshot. + +```toml +[negentropy] +max_sync_events = 1_000_000 # matches strfry's relay__negentropy__maxSyncEvents +``` + +`NegSessionRegistry.open` checks `count >= max_sync_events` after the +SQLite count or during scan, and on overflow sends: + +``` +["NEG-ERR", "", "blocked: too many query results"] +``` + +Exact wording matches strfry so client error-handling that string-matches +behaves identically. (Spec doesn't normatively define error text; this is +de-facto interop.) + +### C — frame-size cap on NEG-MSG: 500_000 bytes (strfry parity) `NegentropyServerSession` is constructed with `frameSizeLimit = 0` -(no limit). At very large reconciliations the message can grow large. -Set a default `frameSizeLimit = 64 * 1024` (matching the typical WS -frame budget) so NEG-MSGs don't blow past `[limits].max_ws_frame_bytes`. +today. **Set `frameSizeLimit = 500_000`** (matching strfry) so a single +NEG-MSG round-trip carries the same payload as strfry's. 64 KB — what +the previous draft suggested — would force 8× more round-trips for +large reconciliations and noticeably slow sync against strfry-style +clients that expect ~1 MB hex-encoded NEG-MSGs. -The library already supports this — pure config change in -`NegSessionRegistry.open`. +The kmp-negentropy library enforces `frameSizeLimit >= 4096` (or `0` +for unlimited). 500_000 is well above the floor. Pure config change in +`NegSessionRegistry.open`; expose via `[negentropy].frame_size_limit` +for operators who tune it down to fit smaller WS frame budgets. -### D — concurrent NEG-OPEN cap +Note: `LimitsSection.max_ws_frame_bytes` +(`geode/src/main/kotlin/com/vitorpamplona/geode/config/RelayConfig.kt:148`) +applies to the WebSocket layer. After hex-encoding a 500_000-byte +negentropy payload doubles to ~1_000_000 bytes on the wire; ensure +`max_ws_frame_bytes` is at least 2 MB (or unlimited) in default +config so the response isn't truncated by the WS layer. + +### D — concurrent NEG-OPEN cap, shared with REQ A NEG-OPEN holds session state until NEG-CLOSE (or connection close). -Today nothing caps the number of concurrent open negentropy sessions -per connection. A misbehaving (or hostile) client could open thousands -and pin RAM. Add `MAX_NEG_SESSIONS_PER_CONNECTION = 16`, send NEG-ERR -on overflow. +Today `NegSessionRegistry.sessions` is unbounded. strfry caps at 200 +sessions per connection **shared with REQ** — both count against +`relay__maxSubsPerConnection`. + +Implement as: +- Reuse the existing per-connection REQ subscription cap (or introduce + one if absent) and let NEG sessions consume the same budget. +- Default cap: **200** to match strfry. Configurable via + `[limits].max_subs_per_connection`. +- On overflow strfry sends a NOTICE, not NEG-ERR: + `["NOTICE", "too many concurrent NEG requests"]`. We should match + this — `NEG-ERR` is reserved for per-session protocol errors in + strfry's model. + +### E — pre-built fingerprint tree (follow-up, not in scope here) + +strfry's real production-scale advantage is the **`negentropy::storage::BTreeLMDB`** backend: a persistent, incrementally-maintained B-tree +of `(timestamp, id)` keys with per-node fingerprint accumulators. +When NEG-OPEN's filter string matches a pre-registered +`NegentropyFilter` tree, strfry skips the materialise-and-seal step +entirely and reconciles directly off the LMDB B-tree in O(log n) +fingerprint computations. See `env.foreach_NegentropyFilter` and +`addStatelessView` in `RelayNegentropy.cpp`. + +Equivalent for geode: a `quartz/.../nip77Negentropy/PrebuiltStorage` +backed by a SQLite-side incremental fingerprint index, registered per +canonical filter. This is a substantial piece of work and **out of +scope for this plan** — call it out for a follow-up +(`geode/plans/2026-05-08-negentropy-prebuilt-tree.md`). The A+B+C+D +combination is enough to match strfry's `MemoryView` path, which is +what 99% of ad-hoc NEG-OPENs hit. + +## Concrete error-string interop table + +Match strfry verbatim — clients in the wild string-match on these: + +| Condition | strfry frame | +|---------------------------------------|-----------------------------------------------------------| +| Snapshot exceeds `max_sync_events` | `["NEG-ERR", "", "blocked: too many query results"]` | +| NEG-MSG for unknown subId | `["NEG-ERR", "", "closed: unknown subscription handle"]` | +| Library `reconcile()` parse failure | `["NEG-ERR", "", "PROTOCOL-ERROR"]` | +| Per-connection sub cap exceeded | `["NOTICE", "too many concurrent NEG requests"]` | +| NEG-MSG before NEG-OPEN seal complete | `["NOTICE", "negentropy error: got NEG-MSG before NEG-OPEN complete"]` | + +The current Kotlin code in `NegSessionRegistry.kt:82` sends +`"error: no negentropy session for "` for the unknown-subId +case. Update to strfry's wording. + +## Wire-level conformance checks + +Before merging, verify against strfry as ground truth: + +1. **Round-trip with `strfry sync`**: stand up a small geode instance, + point `strfry sync ws://geode-host` at it, confirm sync completes + and converges in ≤ comparable round-trips for an N=100k corpus. +2. **idSize**: assert kmp-negentropy round-trips full 32-byte ids; + the 16-byte fingerprint stays internal to the library. +3. **Protocol version byte**: NIP-77 v1 = `0x61`. Confirm + `NegentropyServerSession.processMessage` neither emits nor accepts + a different version byte. Negotiation happens inside the library; + surface the version on `NIP-11.limitation.negentropy = 1` so clients + know the v1 path is supported. ## How to verify -Add to `geode.perf.LoadBenchmark`: +Add to `geode/src/test/kotlin/com/vitorpamplona/geode/perf/LoadBenchmark.kt`: -- `negentropyOpenLatencyLargeCorpus` — preload 1 M events (use - fixtures), measure NEG-OPEN → first NEG-MSG latency. Target <100 ms. +- `negentropyOpenLatencyLargeCorpus` — preload 1 M events, measure + NEG-OPEN → first NEG-MSG latency. Target **<200 ms** (strfry's + C++ `MemoryView` path on equivalent hardware does ~80–150 ms; + we expect a 1.5–2× JVM tax). - `negentropyMemoryPressure` — open 10 concurrent NEG-OPENs on the - same large corpus; measure RSS delta, target <500 MB. + same large corpus; measure RSS delta. Target **<500 MB** + (10 × ~40 MB session footprint + scan overhead). +- `negentropyStrfryInterop` — programmatic round-trip against a + containerised `hoytech/strfry`: same fixture corpus loaded in both, + cross-sync, assert byte-identical id-set convergence in equal + round-trips ±1. + +Add to `quartz/.../nip77Negentropy/`: + +- `NegentropyServerSessionTest.processMessage_atFrameLimit_splitsAcrossRounds` + — large symmetric difference, assert each NEG-MSG payload is + ≤ 500_000 bytes and the session completes in N rounds (not 1). ## Risks -- **`Sequence`/`Flow` over SQLite cursor**: holding a cursor open - across the full sync is fragile if the client stalls. Materialise - to a smaller in-memory list (just (id, createdAt)) once, reuse for - the lifetime of the session. Memory bound is the same. -- **Defaulting `since` is a behaviour change**: existing clients that - expect "everything" silently get a bounded window. Either (a) make - it opt-in via `RelayConfig.NegentropySection.default_lookback_seconds - = null`, (b) advertise the cap in NIP-11 so well-behaved clients - read it. -- **Frame-size cap can break older clients**: the NIP-77 reference - implementation (kmp-negentropy) handles this gracefully — multi-frame - reconciliation is in spec — but field-test against a known-working - client (e.g. nstart, primal-cache) before flipping the default. +- **Cursor lifetime**: holding a SQLite cursor open across the full + sync is fragile if the client stalls. Materialise to a smaller + in-memory `List` (40 bytes/entry) once at NEG-OPEN time, + reuse for the lifetime of the session. Bounded by `max_sync_events`. +- **Frame-size 500_000 vs WS frame budget**: hex-encoded payload is + ~1 MB. If `LimitsSection.max_ws_frame_bytes` is set lower than + ~1.5 MB the response gets truncated/rejected by the WS layer. Lift + the WS default OR cap `frame_size_limit` to + `max_ws_frame_bytes / 2` at startup; fail-fast log a warning if the + operator's config makes negentropy unusable. +- **No default `since` (intentional, but worth flagging)**: a hostile + client doing `NEG-OPEN {kinds:[1]}` against a large corpus will hit + `max_sync_events` and get NEG-ERR. The cap is the protection; do + not also silently bound `since`. +- **kmp-negentropy library version**: pinned to `v1.0.2` + (`gradle/libs.versions.toml:9`). Confirm v1.0.2 enforces + `frameSizeLimit >= 4096`, supports protocol byte `0x61`, and + internal `ID_SIZE = 32`. If any of those don't match strfry, this + plan needs an upstream fix on kmp-negentropy first. From 4b7e0e880ca335b8d6e4682f7380b910669e7216 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 7 May 2026 22:37:59 +0000 Subject: [PATCH 2/4] feat(negentropy): strfry-parity NIP-77 reconciliation path MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements the A+B+C+D plan in geode/plans/2026-05-07-negentropy-large-corpus.md: - A: id-and-time-only snapshot. Adds IEventStore.snapshotIdsForNegentropy returning IdAndTime(createdAt, id) with no Event materialization. SQLite override projects directly off event_headers (~40 B/entry instead of ~1 KB/entry; matches strfry's MemoryView footprint). - B: snapshot-size cap. New [negentropy] config section with max_sync_events=1_000_000 (mirrors strfry's relay__negentropy__maxSyncEvents). Overflow returns NEG-ERR "blocked: too many query results" (strfry-exact wording). No default since-window (strfry honors filters as-is; bounding silently would break interop). - C: 500_000-byte frame cap. NegentropyServerSession.DEFAULT_FRAME_SIZE_LIMIT matches strfry's hard-coded Negentropy ne(storage, 500'000). Configurable via [negentropy].frame_size_limit. - D: per-connection session cap = 200 (matches strfry). Overflow sends NOTICE "too many concurrent NEG requests" (strfry parity). Error-string interop: - NEG-MSG with unknown subId → "closed: unknown subscription handle" - reconcile() parse failure → "PROTOCOL-ERROR" - snapshot overflow → "blocked: too many query results" - per-conn cap → NOTICE "too many concurrent NEG requests" NegentropySettings flows: RelayConfig.NegentropySection → Relay → NostrServer → RelaySession → NegSessionRegistry. RelayHub takes optional settings for tests that need to exercise the caps. Tests: - SnapshotIdsForNegentropyTest covers projection correctness across all indexing strategies + the maxEntries+1 sentinel contract. - Nip77NegentropyTest gains negOpenSnapshotOverflowReturnsStrFryNegErr and negOpenPerConnectionCapEmitsNotice. - Existing negMsgWithoutOpenReturnsNegErr updated to assert strfry wording. --- .../kotlin/com/vitorpamplona/geode/Main.kt | 8 + .../kotlin/com/vitorpamplona/geode/Relay.kt | 9 +- .../com/vitorpamplona/geode/RelayHub.kt | 8 +- .../vitorpamplona/geode/config/RelayConfig.kt | 29 +++ .../geode/Nip77NegentropyTest.kt | 74 +++++- .../cache/interning/InterningEventStore.kt | 6 + .../nip01Core/relay/server/LiveEventStore.kt | 16 ++ .../relay/server/NegSessionRegistry.kt | 62 ++++- .../nip01Core/relay/server/NostrServer.kt | 6 + .../nip01Core/relay/server/RelaySession.kt | 4 +- .../quartz/nip01Core/store/IEventStore.kt | 32 +++ .../quartz/nip01Core/store/IdAndTime.kt | 39 ++++ .../nip01Core/store/ObservableEventStore.kt | 5 + .../nip01Core/store/sqlite/EventStore.kt | 6 + .../nip01Core/store/sqlite/QueryBuilder.kt | 217 ++++++++++++++++++ .../store/sqlite/SQLiteEventStore.kt | 6 + .../NegentropyServerSession.kt | 48 +++- .../nip77Negentropy/NegentropySettings.kt | 50 ++++ .../sqlite/SnapshotIdsForNegentropyTest.kt | 103 +++++++++ .../nip77Negentropy/NegentropySessionTest.kt | 2 +- 20 files changed, 712 insertions(+), 18 deletions(-) create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IdAndTime.kt create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt create mode 100644 quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt index a731890e88..9290fef8e9 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt @@ -32,6 +32,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyAuthOnlyPo 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 java.io.File /** @@ -109,6 +110,12 @@ fun main(args: Array) { } val stateFile = config.admin.state_file?.let { File(it) } + val negentropySettings = + NegentropySettings( + frameSizeLimit = config.negentropy.frame_size_limit, + maxSyncEvents = config.negentropy.max_sync_events, + maxSessionsPerConnection = config.negentropy.max_sessions_per_connection, + ) val relay = Relay( advertisedUrl, @@ -117,6 +124,7 @@ fun main(args: Array) { policyBuilder, stateFile = stateFile, parallelVerify = parallelVerify, + negentropySettings = negentropySettings, ) // Frame cap honors max_ws_frame_bytes when set; max_ws_message_bytes // is treated as the same cap (Ktor's WebSockets plugin only exposes diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/Relay.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/Relay.kt index 6a8436103b..cccdb70909 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/Relay.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/Relay.kt @@ -33,6 +33,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore import com.vitorpamplona.quartz.nip11RelayInfo.Nip11RelayInformation +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import com.vitorpamplona.quartz.nip86RelayManagement.server.BanListPolicy import com.vitorpamplona.quartz.nip86RelayManagement.server.BanStore import kotlinx.coroutines.SupervisorJob @@ -84,6 +85,11 @@ class Relay( * `Main.kt` skips `VerifyPolicy` when this flag is on. */ parallelVerify: Boolean = false, + /** + * NIP-77 server-side tuning (frame cap, snapshot cap, + * per-connection session cap). Defaults to strfry-parity values. + */ + negentropySettings: NegentropySettings = NegentropySettings.Default, ) : AutoCloseable { private val stateStore: RelayStateStore? = stateFile?.let { RelayStateStore(it) } @@ -168,8 +174,9 @@ class Relay( val user = policyBuilder() if (user === EmptyPolicy) BanListPolicy(banStore) else user + BanListPolicy(banStore) }, - parentContext, + parentContext = parentContext, parallelVerify = parallelVerify, + negentropySettings = negentropySettings, ) /** diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/RelayHub.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/RelayHub.kt index 6a6866306e..12d97451a6 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/RelayHub.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/RelayHub.kt @@ -28,6 +28,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.server.policies.EmptyPolicy 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.nip77Negentropy.NegentropySettings import java.util.concurrent.ConcurrentHashMap /** @@ -50,6 +51,7 @@ import java.util.concurrent.ConcurrentHashMap */ class RelayHub( private val defaultPolicy: () -> IRelayPolicy = { EmptyPolicy }, + private val negentropySettings: NegentropySettings = NegentropySettings.Default, ) : WebsocketBuilder, AutoCloseable { private val relays = ConcurrentHashMap() @@ -60,7 +62,11 @@ class RelayHub( fun getOrCreate(url: NormalizedRelayUrl): Relay { check(!closed) { "RelayHub has been closed" } return relays.getOrPut(url) { - Relay(url = url, policyBuilder = defaultPolicy) + Relay( + url = url, + policyBuilder = defaultPolicy, + negentropySettings = negentropySettings, + ) } } diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/config/RelayConfig.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/config/RelayConfig.kt index 1f4ff44415..b8f162ffbe 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/config/RelayConfig.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/config/RelayConfig.kt @@ -43,6 +43,7 @@ data class RelayConfig( val limits: LimitsSection = LimitsSection(), val authorization: AuthorizationSection = AuthorizationSection(), val admin: AdminSection = AdminSection(), + val negentropy: NegentropySection = NegentropySection(), ) { /** * Maps the `[info]` section into a [RelayInfo] used by the NIP-11 @@ -160,6 +161,34 @@ data class RelayConfig( val max_ws_frame_bytes: Int? = null, ) + /** + * NIP-77 negentropy tuning. Defaults track strfry + * (`hoytech/strfry`) so a Geode relay accepts the same workload + * shape and exchanges the same NEG-MSG round-trip size as + * strfry — the de-facto reference implementation. + * + * - [frame_size_limit] mirrors strfry's hard-coded + * `Negentropy ne(storage, 500'000)` in `RelayNegentropy.cpp`. + * Hex-encoded that's ~1 MB on the wire per NEG-MSG; ensure + * `[limits].max_ws_frame_bytes` (when set) is at least double + * this or NEG-MSGs get truncated by the WS layer. + * - [max_sync_events] mirrors strfry's + * `relay__negentropy__maxSyncEvents`. NEG-OPEN whose snapshot + * exceeds this returns + * `["NEG-ERR", "", "blocked: too many query results"]`. + * - [max_sessions_per_connection] caps concurrent NEG-OPEN + * sessions held by a single connection. strfry shares its + * 200-cap with REQ subs via `relay__maxSubsPerConnection`; + * Geode counts NEG independently for now (REQ has no cap yet). + * Overflow returns NOTICE + * `"too many concurrent NEG requests"` (matches strfry). + */ + data class NegentropySection( + val frame_size_limit: Long = 500_000L, + val max_sync_events: Int = 1_000_000, + val max_sessions_per_connection: Int = 200, + ) + data class AuthorizationSection( val pubkey_whitelist: List = emptyList(), val pubkey_blacklist: List = emptyList(), diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/Nip77NegentropyTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/Nip77NegentropyTest.kt index d399457b02..a205d8759d 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/Nip77NegentropyTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/Nip77NegentropyTest.kt @@ -24,6 +24,7 @@ import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer @@ -33,6 +34,7 @@ import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent import com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage import com.vitorpamplona.quartz.nip77Negentropy.NegMsgMessage import com.vitorpamplona.quartz.nip77Negentropy.NegentropySession +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED import kotlinx.coroutines.runBlocking @@ -237,12 +239,82 @@ class Nip77NegentropyTest { val response = client.nextMessage() assertTrue(response is NegErrMessage, "expected NEG-ERR, got ${response::class.simpleName}") assertEquals("ghost-sub", response.subId) - assertTrue(response.reason.contains("no negentropy session")) + // strfry-parity wording — clients in the wild string-match this. + assertEquals("closed: unknown subscription handle", response.reason) } finally { client.close() } } + @Test + fun negOpenSnapshotOverflowReturnsStrFryNegErr() = + runBlocking { + // Tiny cap so the test is fast. Preload more events than the + // cap so NEG-OPEN must reject — strfry's parity behaviour for + // `relay__negentropy__maxSyncEvents`. + val capped = RelayHub(negentropySettings = NegentropySettings(maxSyncEvents = 5)) + try { + val capUrl = RelayUrlNormalizer.normalize("ws://127.0.0.1:7771/") + val events = makeEvents(20) + capped.getOrCreate(capUrl).preload(events) + + val client = WireClient(capped, capUrl) + try { + val session = + NegentropySession( + subId = "neg-overflow", + filter = Filter(kinds = listOf(1)), + localEvents = emptyList(), + ) + client.send(OptimizedJsonMapper.toJson(session.open())) + + val response = client.nextMessage() + assertTrue(response is NegErrMessage, "expected NEG-ERR, got ${response::class.simpleName}") + assertEquals("neg-overflow", response.subId) + // strfry-parity wording. + assertEquals("blocked: too many query results", response.reason) + } finally { + client.close() + } + } finally { + capped.close() + } + } + + @Test + fun negOpenPerConnectionCapEmitsNotice() = + runBlocking { + // Cap = 2, so the third NEG-OPEN on one connection should + // be rejected with a NOTICE (matching strfry's wording). + val capped = RelayHub(negentropySettings = NegentropySettings(maxSessionsPerConnection = 2)) + try { + val capUrl = RelayUrlNormalizer.normalize("ws://127.0.0.1:7772/") + capped.getOrCreate(capUrl).preload(makeEvents(3)) + + val client = WireClient(capped, capUrl) + try { + repeat(2) { i -> + val s = NegentropySession("ok-$i", Filter(kinds = listOf(1)), localEvents = emptyList()) + client.send(OptimizedJsonMapper.toJson(s.open())) + // Drain the NEG-MSG response so the next OPEN goes + // through cleanly. + client.nextMessage() as NegMsgMessage + } + + // Third OPEN — should be rejected with a NOTICE. + val third = NegentropySession("third", Filter(kinds = listOf(1)), localEvents = emptyList()) + client.send(OptimizedJsonMapper.toJson(third.open())) + val response = client.nextMessage() + assertTrue(response is NoticeMessage, "expected NOTICE, got ${response::class.simpleName}") + assertEquals("too many concurrent NEG requests", response.message) + } finally { + client.close() + } + } finally { + capped.close() + } + } + @Test fun negOpenWithSameSubIdReplacesPriorSession() = runBlocking { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt index d224a17d9a..85fbc1c6fa 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/cache/interning/InterningEventStore.kt @@ -24,6 +24,7 @@ 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.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime /** * Decorator that canonicalises every [Event] returned by the inner @@ -111,6 +112,11 @@ class InterningEventStore( override suspend fun count(filters: List): Int = inner.count(filters) + override suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int?, + ): List = inner.snapshotIdsForNegentropy(filters, maxEntries) + override suspend fun delete(filter: Filter) = inner.delete(filter) override suspend fun delete(filters: List) = inner.delete(filters) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt index 912c1e4d04..6264fc7a82 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/LiveEventStore.kt @@ -23,6 +23,7 @@ package com.vitorpamplona.quartz.nip01Core.relay.server import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.channels.BufferOverflow import kotlinx.coroutines.flow.MutableSharedFlow @@ -161,4 +162,19 @@ class LiveEventStore( } return merged } + + /** + * Lightweight snapshot for NIP-77 negentropy. Returns + * `(created_at, id)` pairs only — no Event materialisation — + * matching strfry's `MemoryView` footprint of ~40 B/entry. + * + * If [maxEntries] is non-null, the underlying store may return + * up to `maxEntries + 1` entries; the +1 sentinel lets the + * caller distinguish "exactly at cap" from "exceeds cap" without + * scanning past the cap. + */ + suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int? = null, + ): List = store.snapshotIdsForNegentropy(filters, maxEntries) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt index f85b723d0c..628895fc3d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NegSessionRegistry.kt @@ -21,12 +21,14 @@ package com.vitorpamplona.quartz.nip01Core.relay.server import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.Message +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd import com.vitorpamplona.quartz.nip77Negentropy.NegCloseCmd import com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage import com.vitorpamplona.quartz.nip77Negentropy.NegMsgCmd import com.vitorpamplona.quartz.nip77Negentropy.NegOpenCmd import com.vitorpamplona.quartz.nip77Negentropy.NegentropyServerSession +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings /** * Per-connection NIP-77 negentropy state and dispatch. @@ -39,21 +41,37 @@ import com.vitorpamplona.quartz.nip77Negentropy.NegentropyServerSession * Plain [HashMap] is sufficient because the registry is mutated only * from [RelaySession.receive] — that path is single-threaded per the * WebSocket handler contract. + * + * Defaults match strfry (`hoytech/strfry`) so a Geode relay reconciles + * with the same round-trip shape and the same operator-visible + * protections — see [NegentropySettings]. */ class NegSessionRegistry( private val store: LiveEventStore, private val send: (Message) -> Unit, + private val settings: NegentropySettings = NegentropySettings.Default, ) { private val sessions = HashMap() /** - * Open a reconciliation session. The relay snapshots its matching - * events at this instant — concurrent inserts during the sync are - * not surfaced; clients re-open if they want fresh state. + * Open a reconciliation session. The relay snapshots the matching + * `(created_at, id)` pairs at this instant — concurrent inserts + * during the sync are not surfaced; clients re-open if they want + * fresh state. * * Access control reuses the REQ policy hook: a relay that requires * AUTH or has kind/pubkey allow-deny lists applies the same rules * to NEG-OPEN as it does to subscription REQs. + * + * Two strfry-parity protections fire here: + * - **Per-connection session cap.** If an OPEN would push the + * map past [NegentropySettings.maxSessionsPerConnection], we + * send a NOTICE (matching strfry's + * `"too many concurrent NEG requests"`) and drop the OPEN. + * - **Snapshot size cap.** The store is asked for at most + * `maxSyncEvents + 1` entries; if the +1 sentinel comes back, + * the corpus exceeds the cap and we send NEG-ERR + * `"blocked: too many query results"` (matching strfry). */ suspend fun open( cmd: NegOpenCmd, @@ -66,20 +84,43 @@ class NegSessionRegistry( } val filters = (gate as PolicyResult.Accepted).cmd.filters + // Per-connection cap. Only fires when this is a NEW subId — + // a same-subId re-open replaces the prior session 1-for-1. + val isReopen = sessions.containsKey(cmd.subId) + if (!isReopen && sessions.size >= settings.maxSessionsPerConnection) { + send(NoticeMessage("too many concurrent NEG requests")) + return + } + // NIP-77: same-subId OPEN replaces any prior session. sessions.remove(cmd.subId) - val events = store.snapshotQuery(filters) - val session = NegentropyServerSession(cmd.subId, events) + val cap = settings.maxSyncEvents + val entries = store.snapshotIdsForNegentropy(filters, maxEntries = cap) + if (entries.size > cap) { + send(NegErrMessage(cmd.subId, "blocked: too many query results")) + return + } + + val session = + NegentropyServerSession( + subId = cmd.subId, + localEntries = entries, + frameSizeLimit = settings.frameSizeLimit, + ) sessions[cmd.subId] = session runMessage(cmd.subId, session) { it.processMessage(cmd.initialMessage) } } + /** + * Process a follow-up NEG-MSG. strfry-parity wording for the + * unknown-subId case: `"closed: unknown subscription handle"`. + */ fun msg(cmd: NegMsgCmd) { val session = sessions[cmd.subId] if (session == null) { - send(NegErrMessage(cmd.subId, "error: no negentropy session for ${cmd.subId}")) + send(NegErrMessage(cmd.subId, "closed: unknown subscription handle")) return } runMessage(cmd.subId, session) { it.processMessage(cmd.message) } @@ -99,6 +140,9 @@ class NegSessionRegistry( sessions.clear() } + /** Test/diagnostic accessor. */ + val activeSessionCount: Int get() = sessions.size + private inline fun runMessage( subId: String, session: NegentropyServerSession, @@ -107,9 +151,11 @@ class NegSessionRegistry( try { val response = block(session) if (response != null) send(response) - } catch (e: Exception) { + } catch (_: Exception) { + // strfry sends `PROTOCOL-ERROR` on library reconcile() + // parse failure and tears the session down. sessions.remove(subId) - send(NegErrMessage(subId, "error: ${e.message ?: e::class.simpleName}")) + send(NegErrMessage(subId, "PROTOCOL-ERROR")) } } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt index 0ec8e602be..422d4b1ee2 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/NostrServer.kt @@ -23,6 +23,7 @@ package com.vitorpamplona.quartz.nip01Core.relay.server import com.vitorpamplona.quartz.nip01Core.crypto.verify import com.vitorpamplona.quartz.nip01Core.relay.server.policies.VerifyPolicy import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import com.vitorpamplona.quartz.utils.cache.LargeCache import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.SupervisorJob @@ -43,12 +44,16 @@ import kotlin.coroutines.CoroutineContext * coroutine inside [VerifyPolicy]. Callers that flip this on should * *omit* `VerifyPolicy` from their [policyBuilder] chain to avoid * double-verifying. + * @param negentropySettings NIP-77 server-side tuning (frame cap, + * snapshot cap, per-connection session cap). Defaults to strfry- + * parity values; see [NegentropySettings]. */ class NostrServer( private val store: IEventStore, private val policyBuilder: () -> IRelayPolicy = { VerifyPolicy }, private val parentContext: CoroutineContext = SupervisorJob(), parallelVerify: Boolean = false, + private val negentropySettings: NegentropySettings = NegentropySettings.Default, ) : AutoCloseable { /** Scope for all subscriptions. */ private val scope = CoroutineScope(parentContext + SupervisorJob()) @@ -87,6 +92,7 @@ class NostrServer( onClose = { session -> connections.remove(session.hashCode()) }, + negentropySettings = negentropySettings, ).also { session -> connections.put(session.hashCode(), session) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt index c1698492a3..7a4d2d973d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/RelaySession.kt @@ -38,6 +38,7 @@ import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip77Negentropy.NegCloseCmd import com.vitorpamplona.quartz.nip77Negentropy.NegMsgCmd import com.vitorpamplona.quartz.nip77Negentropy.NegOpenCmd +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings import com.vitorpamplona.quartz.utils.Log import com.vitorpamplona.quartz.utils.cache.LargeCache import kotlinx.coroutines.CancellationException @@ -56,11 +57,12 @@ class RelaySession( private val scope: CoroutineScope, private val onSend: (String) -> Unit, private val onClose: (RelaySession) -> Unit, + negentropySettings: NegentropySettings = NegentropySettings.Default, ) : AutoCloseable { private val subscriptions = LargeCache() /** NIP-77 negentropy state for this connection. */ - private val negentropy = NegSessionRegistry(store, ::send) + private val negentropy = NegSessionRegistry(store, ::send, negentropySettings) private fun addSubscription( subId: String, diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt index 8984e2fc5c..3300185a43 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IEventStore.kt @@ -95,6 +95,38 @@ interface IEventStore : AutoCloseable { suspend fun count(filters: List): Int + /** + * NIP-77 negentropy snapshot. Returns `(created_at, id)` pairs + * for every event matching [filters], with no content/tags/sig + * decode. Used by the server-side reconciliation path to build a + * `StorageVector` without materialising full [Event] objects — + * ~40 B/entry instead of ~1 KB/entry. Order is unspecified; + * negentropy's `seal()` re-sorts. + * + * If [maxEntries] is non-null, the implementation may return up + * to `maxEntries + 1` entries; the caller compares the result + * size to detect overflow (matching strfry's `maxSyncEvents` + * guard). The +1 sentinel lets the caller distinguish "exactly + * capped" from "too many to fit". + * + * Default implementation falls back to the full-decode path so + * non-SQLite stores stay correct; SQLite overrides with a direct + * `SELECT id, created_at` against the `query_by_created_at_id` + * index. Honors the same filter semantics as [query] including + * any `limit`. + */ + suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int? = null, + ): List { + val all = query(filters).map { IdAndTime(it.createdAt, it.id) } + return if (maxEntries != null && all.size > maxEntries + 1) { + all.subList(0, maxEntries + 1) + } else { + all + } + } + suspend fun delete(filter: Filter) suspend fun delete(filters: List) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IdAndTime.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IdAndTime.kt new file mode 100644 index 0000000000..0533dffd99 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/IdAndTime.kt @@ -0,0 +1,39 @@ +/* + * 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.quartz.nip01Core.store + +import com.vitorpamplona.quartz.nip01Core.core.HexKey + +/** + * Lightweight projection of an event used by NIP-77 negentropy: just + * the two fields the reconciliation library indexes — `created_at` + * and the 32-byte event id. + * + * Returned by [IEventStore.snapshotIdsForNegentropy] so the relay can + * build a [com.vitorpamplona.negentropy.storage.StorageVector] without + * materialising full [com.vitorpamplona.quartz.nip01Core.core.Event] + * objects (content, tags, sig). For a 1 M-event snapshot this drops + * peak heap from ~1 GB to ~40 MB — strfry's `MemoryView` parity. + */ +data class IdAndTime( + val createdAt: Long, + val id: HexKey, +) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt index cb18c13a96..b11d1b27e4 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/ObservableEventStore.kt @@ -189,6 +189,11 @@ class ObservableEventStore( override suspend fun count(filters: List): Int = inner.count(filters) + override suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int?, + ): List = inner.snapshotIdsForNegentropy(filters, maxEntries) + override suspend fun delete(filter: Filter) { inner.delete(filter) _changes.emit(StoreChange.DeleteByFilter(listOf(filter))) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt index 097c5672ac..faad3204d3 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/EventStore.kt @@ -26,6 +26,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime /** * SQLite-backed [IEventStore] with default DB-file name and relay @@ -63,6 +64,11 @@ class EventStore( override suspend fun count(filters: List) = store.count(filters) + override suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int?, + ): List = store.snapshotIdsForNegentropy(filters, maxEntries) + override suspend fun delete(filter: Filter) { store.delete(filter) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/QueryBuilder.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/QueryBuilder.kt index 6a3412e522..46fa1579a4 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/QueryBuilder.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/QueryBuilder.kt @@ -28,6 +28,7 @@ import com.vitorpamplona.quartz.nip01Core.core.Kind import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.core.isAddressable import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip01Core.store.sqlite.sql.where import com.vitorpamplona.quartz.utils.EventFactory @@ -172,6 +173,222 @@ class QueryBuilder( } } + // ----------------------------------------------------------------- + // NIP-77 negentropy snapshot path + // + // Projects only (id, created_at) — no content/tags/sig decode — + // so the relay can build a StorageVector without materialising + // full Event objects. ~40 B/entry instead of ~1 KB/entry. + // No ORDER BY: negentropy's seal() re-sorts. No limit injection: + // the per-session cap is enforced upstream as a count check. + // ----------------------------------------------------------------- + fun snapshotIdsForNegentropy( + filters: List, + db: SQLiteConnection, + maxEntries: Int? = null, + ): List { + val inner = + if (filters.size == 1) { + toSnapshotIdsSql(filters.first(), hasher(db)) + } else { + toSnapshotIdsSql(filters, hasher(db)) + } + // Safety cap: wrap with `LIMIT maxEntries + 1` so we can + // detect overflow without scanning beyond the cap. The +1 + // sentinel lets the caller distinguish "exactly capped" from + // "too many to fit". Matches strfry's `maxSyncEvents` guard. + val query = + if (maxEntries != null) { + QuerySpec( + "SELECT id, created_at FROM (${inner.sql}) LIMIT ${maxEntries + 1}", + inner.args, + ) + } else { + inner + } + return db.runIdAndTimeQuery(query) + } + + private fun toSnapshotIdsSql( + filter: Filter, + hasher: TagNameValueHasher, + ): QuerySpec { + val newFilter = filter.toFilterWithDTags() + + // Simple path — no tag joins, no FTS — collapses to a single + // SELECT against event_headers. + if (newFilter.isSimpleQuery()) { + return makeSimpleIdsQuery( + ids = newFilter.ids, + authors = newFilter.authors, + kinds = newFilter.kinds, + dTags = newFilter.dTags, + since = newFilter.since, + until = newFilter.until, + limit = newFilter.limit, + ) + } + + // Search path — FTS join. Project id+created_at off + // event_headers via the FTS row_id linkage. + if (newFilter.isSimpleSearch()) { + return makeSimpleIdsSearch( + search = newFilter.search!!, + ids = newFilter.ids, + authors = newFilter.authors, + kinds = newFilter.kinds, + dTags = newFilter.dTags, + since = newFilter.since, + until = newFilter.until, + limit = newFilter.limit, + ) + } + + // Tag-join path — reuse the existing row_id subquery and + // join back to event_headers for the projection. + val rowIdSubquery = prepareRowIDSubQueries(filter, hasher) + return if (rowIdSubquery == null) { + QuerySpec("SELECT id, created_at FROM event_headers") + } else { + QuerySpec( + """ + SELECT event_headers.id, event_headers.created_at FROM event_headers + INNER JOIN ( + ${rowIdSubquery.sql} + ) AS filtered + ON event_headers.row_id = filtered.row_id + """.trimIndent(), + rowIdSubquery.args, + ) + } + } + + private fun toSnapshotIdsSql( + filters: List, + hasher: TagNameValueHasher, + ): QuerySpec { + val rowIdSubqueries = unionSubqueriesIfNeeded(filters, hasher) + return if (rowIdSubqueries == null) { + QuerySpec("SELECT id, created_at FROM event_headers") + } else { + QuerySpec( + """ + SELECT DISTINCT event_headers.id, event_headers.created_at FROM event_headers + INNER JOIN ( + ${rowIdSubqueries.sql} + ) AS filtered + ON event_headers.row_id = filtered.row_id + """.trimIndent(), + rowIdSubqueries.args, + ) + } + } + + private fun makeSimpleIdsQuery( + ids: List? = null, + authors: List? = null, + kinds: List? = null, + dTags: List? = null, + since: Long? = null, + until: Long? = null, + limit: Int? = null, + ): QuerySpec { + val clause = + where { + ids?.let { equalsOrIn("id", it) } + kinds?.let { equalsOrIn("kind", it) } + authors?.let { equalsOrIn("pubkey", it) } + dTags?.let { equalsOrIn("d_tag", it) } + since?.let { greaterThanOrEquals("created_at", it) } + until?.let { lessThanOrEquals("created_at", it) } + if (dTags != null && kinds != null) { + if (kinds.all { it.isAddressable() }) { + raw("(kind >= 30000 AND kind < 40000)") + } + } + } + + val sql = + buildString { + append("SELECT id, created_at FROM event_headers") + if (clause.conditions.isNotEmpty()) { + append("\nWHERE ") + append(clause.conditions) + } + // Negentropy honors filter `limit` like REQ does + // (matches strfry's NostrFilterGroup behaviour). + // ORDER BY is required for LIMIT to be meaningful. + if (limit != null) { + append("\nORDER BY created_at DESC") + if (indexStrategy.useAndIndexIdOnOrderBy) { + append(", id ASC") + } + append("\nLIMIT ") + append(limit) + } + } + + return QuerySpec(sql, clause.args) + } + + private fun makeSimpleIdsSearch( + search: String, + ids: List? = null, + authors: List? = null, + kinds: List? = null, + dTags: List? = null, + since: Long? = null, + until: Long? = null, + limit: Int? = null, + ): QuerySpec { + val clause = + where { + ids?.let { equalsOrIn("event_headers.id", it) } + match(fts.tableName, search) + kinds?.let { equalsOrIn("event_headers.kind", it) } + authors?.let { equalsOrIn("event_headers.pubkey", it) } + dTags?.let { equalsOrIn("event_headers.d_tag", it) } + since?.let { greaterThanOrEquals("event_headers.created_at", it) } + until?.let { lessThanOrEquals("event_headers.created_at", it) } + if (dTags != null && kinds != null) { + if (kinds.all { it.isAddressable() }) { + raw("(event_headers.kind >= 30000 AND kind < 40000)") + } + } + } + + val sql = + buildString { + append("SELECT event_headers.id, event_headers.created_at FROM event_headers") + append("\nINNER JOIN ${fts.tableName} ON event_headers.row_id = ${fts.tableName}.${fts.eventHeaderRowIdName}") + if (clause.conditions.isNotEmpty()) { + append("\nWHERE ${clause.conditions}") + } + if (limit != null) { + append("\nORDER BY event_headers.created_at DESC") + if (indexStrategy.useAndIndexIdOnOrderBy) { + append(", event_headers.id ASC") + } + append("\nLIMIT ") + append(limit) + } + } + + return QuerySpec(sql, clause.args) + } + + private fun SQLiteConnection.runIdAndTimeQuery(query: QuerySpec): List = + prepare(query.sql).use { stmt -> + query.args.forEachIndexed { index, arg -> + stmt.bindText(index + 1, arg) + } + val results = ArrayList() + while (stmt.step()) { + results.add(IdAndTime(stmt.getLong(1), stmt.getText(0))) + } + results + } + private fun makeEverythingQuery() = "SELECT id, pubkey, created_at, kind, tags, content, sig FROM event_headers ORDER BY created_at DESC${if (indexStrategy.useAndIndexIdOnOrderBy) ", id ASC" else ""}" private fun makeQueryIn(rowIdQuery: String) = diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt index 78b29bffdc..ddd98ef7b1 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SQLiteEventStore.kt @@ -32,6 +32,7 @@ import com.vitorpamplona.quartz.nip01Core.core.isEphemeral import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.store.IEventStore +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.nip40Expiration.isExpired import com.vitorpamplona.quartz.utils.EventFactory @@ -315,6 +316,11 @@ class SQLiteEventStore( suspend fun count(filters: List): Int = pool.useReader { queryBuilder.count(filters, it) } + suspend fun snapshotIdsForNegentropy( + filters: List, + maxEntries: Int? = null, + ): List = pool.useReader { queryBuilder.snapshotIdsForNegentropy(filters, it, maxEntries) } + suspend fun delete(filter: Filter) = pool.useWriter { queryBuilder.delete(filter, it) } suspend fun delete(filters: List) = pool.useWriter { queryBuilder.delete(filters, it) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropyServerSession.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropyServerSession.kt index 4fe221d002..62cbb65634 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropyServerSession.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropyServerSession.kt @@ -23,6 +23,7 @@ package com.vitorpamplona.quartz.nip77Negentropy import com.vitorpamplona.negentropy.Negentropy import com.vitorpamplona.negentropy.storage.StorageVector import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.store.IdAndTime import com.vitorpamplona.quartz.utils.Hex /** @@ -31,28 +32,65 @@ import com.vitorpamplona.quartz.utils.Hex * Used when acting as a relay (or relay-relay sync) to respond to * incoming NEG-OPEN and NEG-MSG from a client. * + * The constructor takes [IdAndTime] entries (just `created_at` and the + * 32-byte event id) to keep the per-session footprint at ~40 B/entry — + * matching strfry's `MemoryView` path. A [List] overload is + * kept for callers (and tests) that already hold full events. + * * Usage: - * 1. On NEG-OPEN: create a [NegentropyServerSession] with the matching local events + * 1. On NEG-OPEN: create a [NegentropyServerSession] with the matching local entries * 2. Call [processMessage] with the initial hex message from NEG-OPEN * 3. Send back the resulting [NegMsgMessage] * 4. On subsequent NEG-MSG: call [processMessage] again and send the response + * + * @param frameSizeLimit max bytes per NEG-MSG response (raw payload, + * before hex). Default `500_000` matches strfry's hard-coded + * `Negentropy ne(storage, 500'000)` so a single round-trip carries + * the same payload as strfry's reconciliation. */ class NegentropyServerSession( val subId: String, - localEvents: List, - frameSizeLimit: Long = 0, + localEntries: List, + frameSizeLimit: Long = DEFAULT_FRAME_SIZE_LIMIT, ) { private val storage = StorageVector() private val negentropy: Negentropy init { - for (event in localEvents) { - storage.insert(event.createdAt, event.id) + for (entry in localEntries) { + storage.insert(entry.createdAt, entry.id) } storage.seal() negentropy = Negentropy(storage, frameSizeLimit) } + companion object { + /** + * strfry parity: `Negentropy ne(storage, 500'000)` in + * `RelayNegentropy.cpp`. Hex-encoded that's ~1 MB on the wire + * per NEG-MSG, the de-facto sync round-trip size. + */ + const val DEFAULT_FRAME_SIZE_LIMIT: Long = 500_000L + + /** + * Convenience for callers that hold full [Event] objects + * (mostly tests + relay-relay sync paths). Production server + * code should call the [IdAndTime] constructor directly via + * `IEventStore.snapshotIdsForNegentropy` to avoid the full + * Event materialisation that this projection collapses. + */ + fun fromEvents( + subId: String, + localEvents: List, + frameSizeLimit: Long = DEFAULT_FRAME_SIZE_LIMIT, + ): NegentropyServerSession = + NegentropyServerSession( + subId = subId, + localEntries = localEvents.map { IdAndTime(it.createdAt, it.id) }, + frameSizeLimit = frameSizeLimit, + ) + } + fun processMessage(hexMessage: String): NegMsgMessage? { val msgBytes = Hex.decode(hexMessage) val result = negentropy.reconcile(msgBytes) diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt new file mode 100644 index 0000000000..896ec7c472 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySettings.kt @@ -0,0 +1,50 @@ +/* + * 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.quartz.nip77Negentropy + +/** + * Server-side NIP-77 tuning. Defaults track strfry + * (`hoytech/strfry`) so a Quartz-based relay accepts the same + * workload shape and exchanges the same NEG-MSG round-trip size. + * + * @param frameSizeLimit Max bytes per NEG-MSG response payload + * (raw, before hex). 500_000 matches strfry's hard-coded + * `Negentropy ne(storage, 500'000)` in `RelayNegentropy.cpp`. + * The `kmp-negentropy` library enforces `>= 4096` (or `0` for + * unlimited). + * @param maxSyncEvents Hard cap on the snapshot size for a single + * NEG-OPEN. Mirrors strfry's `relay__negentropy__maxSyncEvents`. + * Overflow returns NEG-ERR `"blocked: too many query results"`. + * @param maxSessionsPerConnection Cap on concurrent NEG sessions + * held by one connection. strfry shares 200 with REQ subs; we + * count NEG independently. Overflow sends NOTICE + * `"too many concurrent NEG requests"`. + */ +data class NegentropySettings( + val frameSizeLimit: Long = NegentropyServerSession.DEFAULT_FRAME_SIZE_LIMIT, + val maxSyncEvents: Int = 1_000_000, + val maxSessionsPerConnection: Int = 200, +) { + companion object { + /** strfry-equivalent defaults. */ + val Default = NegentropySettings() + } +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt new file mode 100644 index 0000000000..d2cb5e134b --- /dev/null +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt @@ -0,0 +1,103 @@ +/* + * 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.quartz.nip01Core.store.sqlite + +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync +import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * Verifies the NIP-77 negentropy id-and-time projection against the + * full-event query path. Goal: same result set, ~25× lighter + * footprint per row. Run across every indexing strategy via + * [BaseDBTest.forEachDB] so plan changes don't silently break the + * snapshot path. + */ +class SnapshotIdsForNegentropyTest : BaseDBTest() { + private val signer = NostrSignerSync() + + private fun makeEvents(count: Int) = + List(count) { i -> + signer.sign(TextNoteEvent.build("event-$i", createdAt = 1_700_000_000L + i)) + } + + @Test + fun matchesFullQueryForSimpleKindFilter() = + forEachDB { db -> + val events = makeEvents(50) + for (e in events) db.insert(e) + + val filter = Filter(kinds = listOf(1)) + val full = db.query(filter) + val ids = db.snapshotIdsForNegentropy(listOf(filter)) + + assertEquals(full.size, ids.size, "snapshot must cover the same row set") + assertEquals( + full.map { it.id }.toSet(), + ids.map { it.id }.toSet(), + "snapshot ids must match the full-query ids", + ) + // Every (createdAt, id) pair must round-trip. + val byId = full.associate { it.id to it.createdAt } + for (entry in ids) { + assertEquals(byId[entry.id], entry.createdAt, "createdAt mismatch for ${entry.id}") + } + } + + @Test + fun honorsSinceUntilLimit() = + forEachDB { db -> + val events = makeEvents(20) // createdAt 1_700_000_000..1_700_000_019 + for (e in events) db.insert(e) + + // since/until window: [+5, +14] inclusive + val filter = + Filter( + kinds = listOf(1), + since = 1_700_000_005L, + until = 1_700_000_014L, + ) + val ids = db.snapshotIdsForNegentropy(listOf(filter)) + assertEquals(10, ids.size, "since/until window should yield 10 rows") + } + + @Test + fun maxEntriesPlusOneSentinelMarksOverflow() = + forEachDB { db -> + val events = makeEvents(30) + for (e in events) db.insert(e) + + val filter = Filter(kinds = listOf(1)) + // cap = 10; we have 30 rows, so the result must be 11 + // (cap + 1 sentinel) — matches strfry's `maxSyncEvents` + // overflow-detection idiom. + val capped = db.snapshotIdsForNegentropy(listOf(filter), maxEntries = 10) + assertEquals(11, capped.size) + assertTrue(capped.size > 10, "caller relies on size > cap as overflow signal") + + // cap >= total: returns the whole set unchanged. + val whole = db.snapshotIdsForNegentropy(listOf(filter), maxEntries = 100) + assertEquals(30, whole.size) + } +} diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySessionTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySessionTest.kt index 655956239d..3c849e0d27 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySessionTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip77Negentropy/NegentropySessionTest.kt @@ -285,7 +285,7 @@ class NegentropySessionTest { val openCmd = clientSession.open() // Server processes via NegentropyServerSession - val serverSession = NegentropyServerSession("sub1", serverEvents) + val serverSession = NegentropyServerSession.fromEvents("sub1", serverEvents) val response = serverSession.processMessage(openCmd.initialMessage) assertNotNull(response) From 809955c3601ce14d8314d63720abb899807b642e Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 7 May 2026 23:11:07 +0000 Subject: [PATCH 3/4] test(negentropy): strfry-style interop tests for NIP-77 sync MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit New geode/src/test/kotlin/com/vitorpamplona/geode/interop/ folder mirrors strfry's test/syncTest.pl over real WebSocket frames. - InteropSyncDriver: raw-WebSocket helper that drives NEG-OPEN / NEG-MSG round trips against any NIP-77 relay and returns the symmetric difference. No NostrClient indirection — same wire shape strfry sync uses. - GeodeVsGeodeNegentropySyncTest: boots two LocalRelayServer instances, gives each a partially overlapping corpus, drives a pull-sync, closes the loop with REQ + EVENT, asserts both sides converge to the union. Bounded-rounds test on a 200-event corpus guards against pathological framing regressions. - GeodeVsStrfryNegentropySyncTest: opt-in via STRFRY_BIN env var (or -Dstrfry.bin=...). When set, boots strfry as a subprocess on a free loopback port, preloads the same corpus shape via EVENT, and asserts Geode's client-side NegentropySession reconciles against strfry's NEG server with the same haveIds/needIds split as the Geode-vs-Geode case. When unset, prints [skip] and returns — mirrors how LoadBenchmark gates on -DrunLoadBenchmark. Per the agent-derived plan, this is the highest-signal interoperability test we can run without porting upstream's fuzz harness; it covers the e2e NEG-OPEN/NEG-MSG/NEG-CLOSE lifecycle through NegSessionRegistry over real OkHttp frames, which catches issues unit tests can't see (frame fragmentation, WebSocket close semantics, kmp-negentropy ↔ strfry C++ wire shape). --- .../interop/GeodeVsGeodeNegentropySyncTest.kt | 215 ++++++++++++ .../GeodeVsStrfryNegentropySyncTest.kt | 310 ++++++++++++++++++ .../geode/interop/InteropSyncDriver.kt | 195 +++++++++++ 3 files changed, 720 insertions(+) create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt new file mode 100644 index 0000000000..3c7a49a825 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt @@ -0,0 +1,215 @@ +/* + * 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.interop + +import com.vitorpamplona.geode.LocalRelayServer +import com.vitorpamplona.geode.Relay +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndConfirm +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync +import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.runBlocking +import okhttp3.OkHttpClient +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * The Kotlin counterpart to strfry's `test/syncTest.pl`: stand up two + * Geode relays on real WebSocket endpoints, give each a partially + * overlapping corpus, and converge them via NIP-77. + * + * The driver mirrors `strfry sync ws://other --dir both`: + * + * 1. Read the local relay's snapshot for the negotiated filter. + * 2. Open NEG-OPEN against the remote with that snapshot. + * 3. Drive NEG-MSG round trips until the client-side + * [com.vitorpamplona.quartz.nip77Negentropy.NegentropySession] + * reports completion. + * 4. `needIds`: REQ them from the remote, insert into the local relay. + * 5. `haveIds`: fetch from the local relay, publish to the remote. + * + * After the round, both relays must hold the union of the original + * corpora. We assert via REQ on each side; an `id`-filter that returns + * every event we expect, and nothing more, proves convergence + * end-to-end through the NIP-77 server pipeline (`NegSessionRegistry` + * → `NegentropyServerSession` → `IEventStore.snapshotIdsForNegentropy`). + * + * Equivalent to strfry's `runSyncTests.pl` "full DB sync" case at + * small scale. Larger corpora belong in `LoadBenchmark`. + */ +class GeodeVsGeodeNegentropySyncTest { + private lateinit var relayA: Relay + private lateinit var relayB: Relay + private lateinit var serverA: LocalRelayServer + private lateinit var serverB: LocalRelayServer + private lateinit var scope: CoroutineScope + private lateinit var client: NostrClient + private val httpClient = OkHttpClient.Builder().build() + + @BeforeTest + fun setup() { + // The placeholder URLs are normalised so the relay accepts + // them; ports come from the autobind below via [server.url]. + relayA = Relay(url = "ws://127.0.0.1:7771/".normalizeRelayUrl()) + relayB = Relay(url = "ws://127.0.0.1:7772/".normalizeRelayUrl()) + serverA = LocalRelayServer(relayA, host = "127.0.0.1", port = 0).start() + serverB = LocalRelayServer(relayB, host = "127.0.0.1", port = 0).start() + scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + client = NostrClient(BasicOkHttpWebSocket.Builder { _ -> httpClient }, scope) + } + + @AfterTest + fun teardown() { + client.disconnect() + scope.cancel() + serverA.stop() + serverB.stop() + relayA.close() + relayB.close() + } + + /** Generates [count] signed text notes with monotonic createdAt. */ + private fun makeEvents( + count: Int, + seed: Long = 1_700_000_000L, + ): List { + val signer = NostrSignerSync(KeyPair()) + return List(count) { i -> + signer.sign(TextNoteEvent.build("event-$seed-$i", createdAt = seed + i)) + } + } + + /** + * One-way reconciliation from `srcRelay` ⇒ `dstRelay`'s perspective + * (the destination is the side initiating the sync, mirroring + * `strfry sync ` semantics where the calling instance pulls + * from the remote). + * + * Returns the driver result so callers can assert round counts / + * error states. + */ + private fun pullSync( + sourceWsUrl: String, + destLocalEvents: List, + filter: Filter, + ): InteropSyncDriver.Result = + InteropSyncDriver(httpClient).reconcile( + wsUrl = sourceWsUrl, + filter = filter, + localEvents = destLocalEvents, + ) + + @Test + fun bidirectionalSyncConvergesTwoRelays() = + runBlocking { + // Universe of 20 events. Relay A holds [0..14], Relay B + // holds [5..19]. Overlap [5..14], A-only [0..4], B-only + // [15..19]. After bidirectional sync both must hold [0..19]. + val all = makeEvents(20) + val aEvents = all.subList(0, 15) + val bEvents = all.subList(5, 20) + relayA.preload(aEvents) + relayB.preload(bEvents) + + val filter = Filter(kinds = listOf(1)) + + // --- B pulls from A --- + val pullAtoB = pullSync(serverA.url, bEvents, filter) + assertNull(pullAtoB.error, "A→B reconciliation must not error") + // From B's perspective: needs A-only, has B-only. + assertEquals( + aEvents.subList(0, 5).map { it.id }.toSet(), + pullAtoB.needIds, + "B should NEED [0..4] from A", + ) + assertEquals( + bEvents.subList(10, 15).map { it.id }.toSet(), + pullAtoB.haveIds, + "B should announce HAVE for [15..19]", + ) + + // Close the loop: fetch needs from A, push haves to A. + val needFromA = + client.fetchAll( + relay = serverA.url.normalizeRelayUrl(), + filter = Filter(ids = pullAtoB.needIds.toList()), + ) + for (e in needFromA) relayB.preload(e) + + for (id in pullAtoB.haveIds) { + val ev = bEvents.first { it.id == id } + client.publishAndConfirm(ev, setOf(serverA.url.normalizeRelayUrl())) + } + + // --- Verify convergence --- + val expectedAll = all.map { it.id }.toSet() + val onA = relaySnapshotIds(serverA.url, filter) + val onB = relaySnapshotIds(serverB.url, filter) + assertEquals(expectedAll, onA, "Relay A should hold every event after sync") + assertEquals(expectedAll, onB, "Relay B should hold every event after sync") + } + + @Test + fun negentropyConvergesInBoundedRoundsOnSmallCorpus() = + runBlocking { + val all = makeEvents(200) + val aEvents = all.subList(0, 150) + val bEvents = all.subList(50, 200) + relayA.preload(aEvents) + relayB.preload(bEvents) + + val filter = Filter(kinds = listOf(1)) + val res = pullSync(serverA.url, bEvents, filter) + assertNull(res.error) + assertEquals(50, res.needIds.size, "B should need [0..49]") + assertEquals(50, res.haveIds.size, "B should have [150..199]") + // strfry typically converges these in ≤5 rounds; we leave + // generous headroom but guard against catastrophic regression. + assertTrue(res.rounds <= 16, "expected ≤16 NEG-MSG rounds, got ${res.rounds}") + } + + /** Pulls every id matching [filter] from a relay via REQ. */ + private suspend fun relaySnapshotIds( + wsUrl: String, + filter: Filter, + ): Set = + client + .fetchAll( + relay = wsUrl.normalizeRelayUrl(), + filter = filter, + ).map { it.id } + .toSet() +} diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt new file mode 100644 index 0000000000..cb9575e7db --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt @@ -0,0 +1,310 @@ +/* + * 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.interop + +import com.vitorpamplona.geode.LocalRelayServer +import com.vitorpamplona.geode.Relay +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl +import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync +import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.runBlocking +import okhttp3.OkHttpClient +import java.io.File +import java.net.ServerSocket +import kotlin.io.path.createTempDirectory +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * Reciprocal interop test: a Geode client driving NIP-77 against a + * real `strfry` instance. + * + * **Opt-in.** Skipped unless the `STRFRY_BIN` environment variable + * (or `-Dstrfry.bin=…` system property) points at a `strfry` binary. + * When unset the test prints a `[skip]` line and returns. Mirrors the + * way `LoadBenchmark` handles its own opt-in: + * + * STRFRY_BIN=/usr/local/bin/strfry ./gradlew :geode:test \ + * --tests "*GeodeVsStrfryNegentropySyncTest*" + * + * The strfry process is booted on a free loopback port with a fresh + * temp LMDB dir. We feed events into it via the Nostr `EVENT` wire + * (no need for `strfry import`), then run [InteropSyncDriver] against + * it. Strfry's `RelayNegentropy.cpp` answers the same NIP-77 wire we + * test against Geode — if both pass we have byte-shape interop, not + * just "passes our own tests". + */ +class GeodeVsStrfryNegentropySyncTest { + private val strfryBin: String? = + System.getenv("STRFRY_BIN") ?: System.getProperty("strfry.bin") + private val enabled = strfryBin != null + + private lateinit var geodeRelay: Relay + private lateinit var geodeServer: LocalRelayServer + private lateinit var scope: CoroutineScope + private lateinit var client: NostrClient + private lateinit var strfryDir: File + private var strfryProcess: Process? = null + private var strfryUrl: String? = null + private val httpClient = OkHttpClient.Builder().build() + + @BeforeTest + fun setup() { + if (!enabled) return + geodeRelay = Relay(url = "ws://127.0.0.1:7771/".normalizeRelayUrl()) + geodeServer = LocalRelayServer(geodeRelay, host = "127.0.0.1", port = 0).start() + scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + client = NostrClient(BasicOkHttpWebSocket.Builder { _ -> httpClient }, scope) + } + + @AfterTest + fun teardown() { + if (!enabled) return + strfryProcess?.destroy() + strfryProcess?.waitFor() + if (::strfryDir.isInitialized) strfryDir.deleteRecursively() + client.disconnect() + scope.cancel() + geodeServer.stop() + geodeRelay.close() + } + + /** + * Boots a strfry instance in a temp directory on a free port. + * Writes a minimal config, starts the relay subprocess, and spins + * until the WebSocket port is accepting connections. + */ + private fun startStrfry(): String { + val port = ServerSocket(0).use { it.localPort } + strfryDir = createTempDirectory(prefix = "strfry-interop-").toFile() + val configFile = File(strfryDir, "strfry.conf") + // Minimal strfry config — bind, db dir, NIP-77 enabled. Strfry + // uses libconfig's hcl-ish syntax; this snippet is the smallest + // that boots a relay accepting NEG-OPEN/REQ/EVENT. + configFile.writeText( + """ + db = "${strfryDir.absolutePath}/strfry-db" + relay { + bind = "127.0.0.1" + port = $port + nofiles = 1024000 + negentropy { + enabled = true + maxSyncEvents = 1000000 + } + } + """.trimIndent(), + ) + File(strfryDir, "strfry-db").mkdirs() + val pb = + ProcessBuilder(strfryBin, "--config", configFile.absolutePath, "relay") + .redirectErrorStream(true) + .redirectOutput(File(strfryDir, "strfry.log")) + strfryProcess = pb.start() + + // Wait for strfry to start accepting connections — give up + // after a few seconds. Strfry typically binds in <500 ms. + val deadline = System.currentTimeMillis() + 5_000 + while (System.currentTimeMillis() < deadline) { + runCatching { + java.net.Socket("127.0.0.1", port).close() + strfryUrl = "ws://127.0.0.1:$port/" + return strfryUrl!! + } + Thread.sleep(100) + } + throw IllegalStateException( + "strfry did not start within 5s; log: " + + File(strfryDir, "strfry.log").readText(), + ) + } + + private fun makeEvents(count: Int): List { + val signer = NostrSignerSync(KeyPair()) + val now = 1_700_000_000L + return List(count) { i -> + signer.sign(TextNoteEvent.build("strfry-interop-$i", createdAt = now + i)) + } + } + + /** + * Push an event directly into a relay over a one-shot WebSocket. + * Used for both Geode (via `geodeRelay.preload`) and strfry + * (via this method) so the corpus is byte-identical on both sides. + */ + private fun publishToStrfry( + wsUrl: String, + events: List, + ) { + val ok = + java.util.concurrent.atomic + .AtomicInteger() + val target = events.size + val incoming = kotlinx.coroutines.channels.Channel(kotlinx.coroutines.channels.Channel.UNLIMITED) + val ws = + httpClient.newWebSocket( + okhttp3.Request + .Builder() + .url(wsUrl.replace("ws://", "http://")) + .build(), + object : okhttp3.WebSocketListener() { + override fun onMessage( + webSocket: okhttp3.WebSocket, + text: String, + ) { + incoming.trySend(text) + } + + override fun onFailure( + webSocket: okhttp3.WebSocket, + t: Throwable, + response: okhttp3.Response?, + ) { + incoming.close(t) + } + }, + ) + try { + for (e in events) { + val cmd = """["EVENT",${OptimizedJsonMapper.toJson(e)}]""" + check(ws.send(cmd)) { "publish to strfry failed" } + } + // Drain OK responses. + runBlocking { + kotlinx.coroutines.withTimeout(30_000) { + while (ok.get() < target) { + val raw = incoming.receive() + if (raw.startsWith("[\"OK\"")) ok.incrementAndGet() + } + } + } + } finally { + ws.close(1000, "preload-done") + } + } + + @Test + fun geodeReconcilesAgainstStrfryRelay() = + runBlocking { + if (!enabled) { + println("[skip] GeodeVsStrfryNegentropySyncTest — set STRFRY_BIN=/path/to/strfry to enable") + return@runBlocking + } + val strfryWs = startStrfry() + + // Same overlap shape as the Geode-vs-Geode test so the two + // results are directly comparable: A=[0..14], B=[5..19]. + val all = makeEvents(20) + val strfryEvents = all.subList(0, 15) + val geodeEvents = all.subList(5, 20) + + publishToStrfry(strfryWs, strfryEvents) + geodeRelay.preload(geodeEvents) + + val filter = Filter(kinds = listOf(1)) + + // Drive the negentropy reconciliation from Geode's + // perspective against strfry. This is the wire we care + // about: our client-side `NegentropySession` (kmp-negentropy) + // talking to strfry's server-side `Negentropy ne(storage, + // 500'000)`. Symmetric difference must match the Geode-vs- + // Geode case. + val res = InteropSyncDriver(httpClient).reconcile(strfryWs, filter, geodeEvents) + assertNull(res.error, "Geode↔strfry NEG must not error: ${res.error}") + assertEquals( + strfryEvents.subList(0, 5).map { it.id }.toSet(), + res.needIds, + "Geode should NEED [0..4] from strfry", + ) + assertEquals( + geodeEvents.subList(10, 15).map { it.id }.toSet(), + res.haveIds, + "Geode should announce HAVE for [15..19]", + ) + + // Spot-check the wire-level health: round count is bounded. + assertTrue(res.rounds <= 16, "expected ≤16 NEG-MSG rounds, got ${res.rounds}") + } + + @Test + fun strfryDrivesGeodeAsServer() = + runBlocking { + if (!enabled) { + println("[skip] strfryDrivesGeodeAsServer — set STRFRY_BIN=/path/to/strfry to enable") + return@runBlocking + } + // Reverse direction: the server under test is *Geode*. + // We use kmp-negentropy as the client driver — same role + // strfry's own client takes when running `strfry sync + // ws://geode`. We don't actually shell out to `strfry sync` + // here (its CLI doesn't expose the corpus split we want + // to test); the wire-level equivalence is what matters, + // and InteropSyncDriver uses the same NIP-77 protocol that + // strfry's client speaks. + startStrfry() // unused — we just need to confirm the binary boots + val all = makeEvents(20) + val geodeEvents = all.subList(0, 15) + val driverEvents = all.subList(5, 20) + geodeRelay.preload(geodeEvents) + + val res = + InteropSyncDriver(httpClient).reconcile( + wsUrl = geodeServer.url, + filter = Filter(kinds = listOf(1)), + localEvents = driverEvents, + ) + assertNull(res.error) + assertEquals( + geodeEvents.subList(0, 5).map { it.id }.toSet(), + res.needIds, + ) + assertEquals( + driverEvents.subList(10, 15).map { it.id }.toSet(), + res.haveIds, + ) + + // Cross-check with REQ that Geode actually has what we + // think it has. + val onGeode = + client + .fetchAll( + relay = geodeServer.url.normalizeRelayUrl(), + filter = Filter(kinds = listOf(1)), + ).map { it.id } + .toSet() + assertEquals(geodeEvents.map { it.id }.toSet(), onGeode) + } +} diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt new file mode 100644 index 0000000000..1ee39d8b74 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt @@ -0,0 +1,195 @@ +/* + * 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.interop + +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper +import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.NoticeMessage +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip77Negentropy.NegErrMessage +import com.vitorpamplona.quartz.nip77Negentropy.NegMsgMessage +import com.vitorpamplona.quartz.nip77Negentropy.NegentropySession +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.Response +import okhttp3.WebSocket +import okhttp3.WebSocketListener +import kotlin.test.fail + +/** + * Equivalent of strfry's `test/syncTest.pl` driver for our interop + * tests. Drives one round of NIP-77 reconciliation against a real + * `ws://` endpoint (a `LocalRelayServer` running Geode, or any other + * relay that speaks NIP-77 such as `strfry`). + * + * The driver opens a raw WebSocket — no `NostrClient` overhead — + * because the goal is to keep the wire format under our direct + * control, the same way `strfry sync` does. That way we exercise the + * server's NEG-OPEN / NEG-MSG / NEG-CLOSE handling with no client- + * side framing or filter-management indirection. + * + * Given a server endpoint and a `localEvents` snapshot, drives the + * reconciliation until completion (or [maxRounds] is reached), and + * returns the symmetric difference plus stats. + */ +class InteropSyncDriver( + private val httpClient: OkHttpClient = OkHttpClient.Builder().build(), +) { + /** + * Reconciles `localEvents` against the relay at [wsUrl] under + * [filter]. Returns the symmetric-difference id sets so the + * caller can verify or close the loop with REQ/EVENT. + * + * @param wsUrl the source relay's `ws://…` URL. + * @param filter NEG-OPEN filter — usually the broadest filter the + * sync should cover (e.g. `Filter(kinds = listOf(1))`). + * @param localEvents events the caller already has; the relay + * reconciles these against its own snapshot. + * @param frameSizeLimit `0` lets the relay choose. We pass `0` + * here so the relay's configured cap (500_000 by default) is + * what governs framing — same shape as `strfry sync`. + * @param timeoutMs hard timeout on a single NEG-MSG round trip. + * @param maxRounds upper bound on round trips. Strfry's typical + * converge in ≤5 rounds for 100 k corpora; 64 is a generous + * safety net that catches pathological splits without hanging + * tests forever. + */ + fun reconcile( + wsUrl: String, + filter: Filter, + localEvents: List, + subId: String = "interop-sync", + frameSizeLimit: Long = 0, + timeoutMs: Long = 30_000L, + maxRounds: Int = 64, + ): Result { + val incoming = Channel(UNLIMITED) + val ws = + httpClient.newWebSocket( + Request.Builder().url(wsUrl.replace("ws://", "http://")).build(), + object : WebSocketListener() { + override fun onMessage( + webSocket: WebSocket, + text: String, + ) { + incoming.trySend(text) + } + + override fun onClosing( + webSocket: WebSocket, + code: Int, + reason: String, + ) { + incoming.close() + } + + override fun onFailure( + webSocket: WebSocket, + t: Throwable, + response: Response?, + ) { + incoming.close(t) + } + }, + ) + + return try { + val session = NegentropySession(subId, filter, localEvents, frameSizeLimit) + + // Step 1: NEG-OPEN. + check(ws.send(OptimizedJsonMapper.toJson(session.open()))) { "send NEG-OPEN failed" } + + // Step 2: drive NEG-MSG round trips until the client-side + // session reports completion. + val haveIds = mutableSetOf() + val needIds = mutableSetOf() + var rounds = 0 + while (rounds < maxRounds) { + val raw = + runBlocking { + withTimeout(timeoutMs) { incoming.receive() } + } + val msg = OptimizedJsonMapper.fromJsonToMessage(raw) + rounds++ + when (msg) { + is NegErrMessage -> { + return Result( + haveIds = haveIds, + needIds = needIds, + rounds = rounds, + error = "${msg.subId}: ${msg.reason}", + ) + } + + is NoticeMessage -> { + return Result( + haveIds = haveIds, + needIds = needIds, + rounds = rounds, + error = "NOTICE: ${msg.message}", + ) + } + + is NegMsgMessage -> { + val r = session.processMessage(msg.message) + haveIds += r.haveIds + needIds += r.needIds + if (r.isComplete()) { + return Result(haveIds, needIds, rounds, error = null) + } + check(ws.send(OptimizedJsonMapper.toJson(r.nextCmd!!))) { + "send NEG-MSG failed" + } + } + + else -> { + fail("unexpected message during NEG sync: ${msg::class.simpleName}") + } + } + } + Result(haveIds, needIds, rounds, error = "did not converge in $maxRounds rounds") + } finally { + ws.close(1000, "interop-test-done") + } + } + + /** + * Result of a reconciliation. `error` is non-null on + * NEG-ERR/NOTICE/timeout; otherwise the id-set fields are + * authoritative. + * + * @param haveIds events the client (us) had that the relay did not. + * @param needIds events the relay had that the client (us) lacked. + * @param rounds NEG-MSG round trips, including the one carrying + * the terminator. + */ + data class Result( + val haveIds: Set, + val needIds: Set, + val rounds: Int, + val error: String?, + ) +} From 7c4f2b720a28166ae45c55232cea2aef08e705c8 Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 7 May 2026 23:21:17 +0000 Subject: [PATCH 4/4] refactor(negentropy): audit fixes for interop tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Delete `strfryDrivesGeodeAsServer` — boots strfry but uses kmp-negentropy ↔ Geode, no actual strfry-vs-geode interop. Duplicates `Nip77NegentropyTest.negentropyComputesSymmetricDifference`. - Strip speculative `negentropy { enabled = ... }` and `nofiles` blocks from the strfry config; defaults are what we want to test. - `InteropSyncDriver.reconcile` → `negotiate`. The function computes the symmetric difference; it doesn't move events. Convert to `suspend fun` to drop a nested `runBlocking` that could deadlock under dispatcher pressure. - Inline `pullSync` helper in GeodeVsGeodeNegentropySyncTest — it was a one-line wrapper with a misleading name. - Batch `relayB.preload(needFromA)` instead of looping. - Tighten `idsOnRelay(NormalizedRelayUrl)` signature (was re-parsing the string on every call). - Drop redundant `assertTrue(size > cap)` after `assertEquals(11)`. --- .../interop/GeodeVsGeodeNegentropySyncTest.kt | 77 ++----- .../GeodeVsStrfryNegentropySyncTest.kt | 212 +++++------------- .../geode/interop/InteropSyncDriver.kt | 38 ++-- .../sqlite/SnapshotIdsForNegentropyTest.kt | 5 +- 4 files changed, 105 insertions(+), 227 deletions(-) diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt index 3c7a49a825..dd79a9475b 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsGeodeNegentropySyncTest.kt @@ -29,6 +29,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndConfirm import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync @@ -112,25 +113,7 @@ class GeodeVsGeodeNegentropySyncTest { } } - /** - * One-way reconciliation from `srcRelay` ⇒ `dstRelay`'s perspective - * (the destination is the side initiating the sync, mirroring - * `strfry sync ` semantics where the calling instance pulls - * from the remote). - * - * Returns the driver result so callers can assert round counts / - * error states. - */ - private fun pullSync( - sourceWsUrl: String, - destLocalEvents: List, - filter: Filter, - ): InteropSyncDriver.Result = - InteropSyncDriver(httpClient).reconcile( - wsUrl = sourceWsUrl, - filter = filter, - localEvents = destLocalEvents, - ) + private val driver by lazy { InteropSyncDriver(httpClient) } @Test fun bidirectionalSyncConvergesTwoRelays() = @@ -145,71 +128,59 @@ class GeodeVsGeodeNegentropySyncTest { relayB.preload(bEvents) val filter = Filter(kinds = listOf(1)) + val urlA = serverA.url.normalizeRelayUrl() + val urlB = serverB.url.normalizeRelayUrl() - // --- B pulls from A --- - val pullAtoB = pullSync(serverA.url, bEvents, filter) - assertNull(pullAtoB.error, "A→B reconciliation must not error") + // B negotiates the symmetric difference with A. + val diff = driver.negotiate(serverA.url, filter, bEvents) + assertNull(diff.error, "negotiation must not error: ${diff.error}") // From B's perspective: needs A-only, has B-only. assertEquals( aEvents.subList(0, 5).map { it.id }.toSet(), - pullAtoB.needIds, + diff.needIds, "B should NEED [0..4] from A", ) assertEquals( bEvents.subList(10, 15).map { it.id }.toSet(), - pullAtoB.haveIds, + diff.haveIds, "B should announce HAVE for [15..19]", ) // Close the loop: fetch needs from A, push haves to A. val needFromA = - client.fetchAll( - relay = serverA.url.normalizeRelayUrl(), - filter = Filter(ids = pullAtoB.needIds.toList()), - ) - for (e in needFromA) relayB.preload(e) + client.fetchAll(relay = urlA, filter = Filter(ids = diff.needIds.toList())) + relayB.preload(needFromA) - for (id in pullAtoB.haveIds) { - val ev = bEvents.first { it.id == id } - client.publishAndConfirm(ev, setOf(serverA.url.normalizeRelayUrl())) + for (id in diff.haveIds) { + client.publishAndConfirm(bEvents.first { it.id == id }, setOf(urlA)) } // --- Verify convergence --- - val expectedAll = all.map { it.id }.toSet() - val onA = relaySnapshotIds(serverA.url, filter) - val onB = relaySnapshotIds(serverB.url, filter) - assertEquals(expectedAll, onA, "Relay A should hold every event after sync") - assertEquals(expectedAll, onB, "Relay B should hold every event after sync") + val expected = all.map { it.id }.toSet() + assertEquals(expected, idsOnRelay(urlA, filter), "Relay A should hold every event") + assertEquals(expected, idsOnRelay(urlB, filter), "Relay B should hold every event") } @Test fun negentropyConvergesInBoundedRoundsOnSmallCorpus() = runBlocking { val all = makeEvents(200) - val aEvents = all.subList(0, 150) + relayA.preload(all.subList(0, 150)) val bEvents = all.subList(50, 200) - relayA.preload(aEvents) relayB.preload(bEvents) - val filter = Filter(kinds = listOf(1)) - val res = pullSync(serverA.url, bEvents, filter) + val res = driver.negotiate(serverA.url, Filter(kinds = listOf(1)), bEvents) assertNull(res.error) assertEquals(50, res.needIds.size, "B should need [0..49]") assertEquals(50, res.haveIds.size, "B should have [150..199]") - // strfry typically converges these in ≤5 rounds; we leave - // generous headroom but guard against catastrophic regression. + // strfry typically converges these in ≤5 rounds; ≤16 is + // generous headroom that still catches regressions. assertTrue(res.rounds <= 16, "expected ≤16 NEG-MSG rounds, got ${res.rounds}") } - /** Pulls every id matching [filter] from a relay via REQ. */ - private suspend fun relaySnapshotIds( - wsUrl: String, + /** Every event id matching [filter] visible on the relay at [url] via REQ. */ + private suspend fun idsOnRelay( + url: NormalizedRelayUrl, filter: Filter, - ): Set = - client - .fetchAll( - relay = wsUrl.normalizeRelayUrl(), - filter = filter, - ).map { it.id } - .toSet() + ): Set = client.fetchAll(relay = url, filter = filter).map { it.id }.toSet() } diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt index cb9575e7db..ef6271ce6d 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/GeodeVsStrfryNegentropySyncTest.kt @@ -20,129 +20,100 @@ */ package com.vitorpamplona.geode.interop -import com.vitorpamplona.geode.LocalRelayServer -import com.vitorpamplona.geode.Relay import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair -import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient -import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter -import com.vitorpamplona.quartz.nip01Core.relay.normalizer.normalizeRelayUrl -import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent -import kotlinx.coroutines.CoroutineScope -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.SupervisorJob -import kotlinx.coroutines.cancel +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.Response +import okhttp3.WebSocket +import okhttp3.WebSocketListener import java.io.File import java.net.ServerSocket +import java.net.Socket import kotlin.io.path.createTempDirectory import kotlin.test.AfterTest -import kotlin.test.BeforeTest import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertNull import kotlin.test.assertTrue /** - * Reciprocal interop test: a Geode client driving NIP-77 against a - * real `strfry` instance. + * Reciprocal interop test: Geode's NIP-77 client driving a real + * `strfry` instance. * - * **Opt-in.** Skipped unless the `STRFRY_BIN` environment variable - * (or `-Dstrfry.bin=…` system property) points at a `strfry` binary. - * When unset the test prints a `[skip]` line and returns. Mirrors the - * way `LoadBenchmark` handles its own opt-in: + * **Opt-in.** Skipped unless `STRFRY_BIN` env var (or + * `-Dstrfry.bin=...`) points at a `strfry` binary. When unset, the + * test prints a `[skip]` line and returns. Mirrors the gate + * `LoadBenchmark` uses for `-DrunLoadBenchmark`: * * STRFRY_BIN=/usr/local/bin/strfry ./gradlew :geode:test \ * --tests "*GeodeVsStrfryNegentropySyncTest*" * * The strfry process is booted on a free loopback port with a fresh - * temp LMDB dir. We feed events into it via the Nostr `EVENT` wire - * (no need for `strfry import`), then run [InteropSyncDriver] against - * it. Strfry's `RelayNegentropy.cpp` answers the same NIP-77 wire we - * test against Geode — if both pass we have byte-shape interop, not - * just "passes our own tests". + * temp LMDB dir. We feed events into it via the NIP-01 EVENT wire + * (no `strfry import`), then run [InteropSyncDriver] against it. + * Strfry's `RelayNegentropy.cpp` answers the same NIP-77 wire we test + * against Geode — passing both is byte-shape interop, not just + * "passes our own tests". */ class GeodeVsStrfryNegentropySyncTest { private val strfryBin: String? = System.getenv("STRFRY_BIN") ?: System.getProperty("strfry.bin") private val enabled = strfryBin != null - private lateinit var geodeRelay: Relay - private lateinit var geodeServer: LocalRelayServer - private lateinit var scope: CoroutineScope - private lateinit var client: NostrClient private lateinit var strfryDir: File private var strfryProcess: Process? = null - private var strfryUrl: String? = null - private val httpClient = OkHttpClient.Builder().build() - - @BeforeTest - fun setup() { - if (!enabled) return - geodeRelay = Relay(url = "ws://127.0.0.1:7771/".normalizeRelayUrl()) - geodeServer = LocalRelayServer(geodeRelay, host = "127.0.0.1", port = 0).start() - scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) - client = NostrClient(BasicOkHttpWebSocket.Builder { _ -> httpClient }, scope) - } + private val httpClient by lazy { OkHttpClient.Builder().build() } @AfterTest fun teardown() { - if (!enabled) return strfryProcess?.destroy() strfryProcess?.waitFor() if (::strfryDir.isInitialized) strfryDir.deleteRecursively() - client.disconnect() - scope.cancel() - geodeServer.stop() - geodeRelay.close() } /** - * Boots a strfry instance in a temp directory on a free port. - * Writes a minimal config, starts the relay subprocess, and spins - * until the WebSocket port is accepting connections. + * Boots a strfry instance in a temp LMDB dir on a free port. + * Writes the smallest config strfry accepts, starts the + * subprocess, and polls until the WebSocket port is reachable. + * + * The config is intentionally minimal — strfry's defaults + * (negentropy on, sane limits) are what we want to test against. + * Adding speculative config keys risks failing on schema drift. */ private fun startStrfry(): String { val port = ServerSocket(0).use { it.localPort } strfryDir = createTempDirectory(prefix = "strfry-interop-").toFile() val configFile = File(strfryDir, "strfry.conf") - // Minimal strfry config — bind, db dir, NIP-77 enabled. Strfry - // uses libconfig's hcl-ish syntax; this snippet is the smallest - // that boots a relay accepting NEG-OPEN/REQ/EVENT. configFile.writeText( """ db = "${strfryDir.absolutePath}/strfry-db" relay { bind = "127.0.0.1" port = $port - nofiles = 1024000 - negentropy { - enabled = true - maxSyncEvents = 1000000 - } } """.trimIndent(), ) File(strfryDir, "strfry-db").mkdirs() - val pb = + strfryProcess = ProcessBuilder(strfryBin, "--config", configFile.absolutePath, "relay") .redirectErrorStream(true) .redirectOutput(File(strfryDir, "strfry.log")) - strfryProcess = pb.start() + .start() - // Wait for strfry to start accepting connections — give up - // after a few seconds. Strfry typically binds in <500 ms. val deadline = System.currentTimeMillis() + 5_000 while (System.currentTimeMillis() < deadline) { runCatching { - java.net.Socket("127.0.0.1", port).close() - strfryUrl = "ws://127.0.0.1:$port/" - return strfryUrl!! + Socket("127.0.0.1", port).close() + return "ws://127.0.0.1:$port/" } Thread.sleep(100) } @@ -161,37 +132,33 @@ class GeodeVsStrfryNegentropySyncTest { } /** - * Push an event directly into a relay over a one-shot WebSocket. - * Used for both Geode (via `geodeRelay.preload`) and strfry - * (via this method) so the corpus is byte-identical on both sides. + * Push every event in [events] to the relay at [wsUrl] over a + * one-shot WebSocket, waiting for an `OK` response per event. + * + * Used to seed strfry from the same `Event` objects Geode's + * `Relay.preload` accepts — that way both sides start from a + * byte-identical corpus. */ - private fun publishToStrfry( + private suspend fun publishToStrfry( wsUrl: String, events: List, ) { - val ok = - java.util.concurrent.atomic - .AtomicInteger() - val target = events.size - val incoming = kotlinx.coroutines.channels.Channel(kotlinx.coroutines.channels.Channel.UNLIMITED) + val incoming = Channel(UNLIMITED) val ws = httpClient.newWebSocket( - okhttp3.Request - .Builder() - .url(wsUrl.replace("ws://", "http://")) - .build(), - object : okhttp3.WebSocketListener() { + Request.Builder().url(wsUrl.replace("ws://", "http://")).build(), + object : WebSocketListener() { override fun onMessage( - webSocket: okhttp3.WebSocket, + webSocket: WebSocket, text: String, ) { incoming.trySend(text) } override fun onFailure( - webSocket: okhttp3.WebSocket, + webSocket: WebSocket, t: Throwable, - response: okhttp3.Response?, + response: Response?, ) { incoming.close(t) } @@ -199,16 +166,14 @@ class GeodeVsStrfryNegentropySyncTest { ) try { for (e in events) { - val cmd = """["EVENT",${OptimizedJsonMapper.toJson(e)}]""" - check(ws.send(cmd)) { "publish to strfry failed" } + check(ws.send("""["EVENT",${OptimizedJsonMapper.toJson(e)}]""")) { + "publish to strfry failed" + } } - // Drain OK responses. - runBlocking { - kotlinx.coroutines.withTimeout(30_000) { - while (ok.get() < target) { - val raw = incoming.receive() - if (raw.startsWith("[\"OK\"")) ok.incrementAndGet() - } + var oks = 0 + withTimeout(30_000) { + while (oks < events.size) { + if (incoming.receive().startsWith("[\"OK\"")) oks++ } } } finally { @@ -225,86 +190,31 @@ class GeodeVsStrfryNegentropySyncTest { } val strfryWs = startStrfry() - // Same overlap shape as the Geode-vs-Geode test so the two - // results are directly comparable: A=[0..14], B=[5..19]. + // Same overlap shape as GeodeVsGeodeNegentropySyncTest so + // results are directly comparable: A=[0..14], local=[5..19]. val all = makeEvents(20) val strfryEvents = all.subList(0, 15) - val geodeEvents = all.subList(5, 20) + val localEvents = all.subList(5, 20) publishToStrfry(strfryWs, strfryEvents) - geodeRelay.preload(geodeEvents) val filter = Filter(kinds = listOf(1)) - // Drive the negentropy reconciliation from Geode's - // perspective against strfry. This is the wire we care - // about: our client-side `NegentropySession` (kmp-negentropy) - // talking to strfry's server-side `Negentropy ne(storage, - // 500'000)`. Symmetric difference must match the Geode-vs- - // Geode case. - val res = InteropSyncDriver(httpClient).reconcile(strfryWs, filter, geodeEvents) + // The wire we care about: kmp-negentropy (client) talking + // to strfry's `Negentropy ne(storage, 500'000)` (server). + // Symmetric difference must match the Geode-vs-Geode case. + val res = InteropSyncDriver(httpClient).negotiate(strfryWs, filter, localEvents) assertNull(res.error, "Geode↔strfry NEG must not error: ${res.error}") assertEquals( strfryEvents.subList(0, 5).map { it.id }.toSet(), res.needIds, - "Geode should NEED [0..4] from strfry", + "client should NEED [0..4] from strfry", ) assertEquals( - geodeEvents.subList(10, 15).map { it.id }.toSet(), + localEvents.subList(10, 15).map { it.id }.toSet(), res.haveIds, - "Geode should announce HAVE for [15..19]", + "client should announce HAVE for [15..19]", ) - - // Spot-check the wire-level health: round count is bounded. assertTrue(res.rounds <= 16, "expected ≤16 NEG-MSG rounds, got ${res.rounds}") } - - @Test - fun strfryDrivesGeodeAsServer() = - runBlocking { - if (!enabled) { - println("[skip] strfryDrivesGeodeAsServer — set STRFRY_BIN=/path/to/strfry to enable") - return@runBlocking - } - // Reverse direction: the server under test is *Geode*. - // We use kmp-negentropy as the client driver — same role - // strfry's own client takes when running `strfry sync - // ws://geode`. We don't actually shell out to `strfry sync` - // here (its CLI doesn't expose the corpus split we want - // to test); the wire-level equivalence is what matters, - // and InteropSyncDriver uses the same NIP-77 protocol that - // strfry's client speaks. - startStrfry() // unused — we just need to confirm the binary boots - val all = makeEvents(20) - val geodeEvents = all.subList(0, 15) - val driverEvents = all.subList(5, 20) - geodeRelay.preload(geodeEvents) - - val res = - InteropSyncDriver(httpClient).reconcile( - wsUrl = geodeServer.url, - filter = Filter(kinds = listOf(1)), - localEvents = driverEvents, - ) - assertNull(res.error) - assertEquals( - geodeEvents.subList(0, 5).map { it.id }.toSet(), - res.needIds, - ) - assertEquals( - driverEvents.subList(10, 15).map { it.id }.toSet(), - res.haveIds, - ) - - // Cross-check with REQ that Geode actually has what we - // think it has. - val onGeode = - client - .fetchAll( - relay = geodeServer.url.normalizeRelayUrl(), - filter = Filter(kinds = listOf(1)), - ).map { it.id } - .toSet() - assertEquals(geodeEvents.map { it.id }.toSet(), onGeode) - } } diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt index 1ee39d8b74..1768a0f961 100644 --- a/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/interop/InteropSyncDriver.kt @@ -30,7 +30,6 @@ import com.vitorpamplona.quartz.nip77Negentropy.NegMsgMessage import com.vitorpamplona.quartz.nip77Negentropy.NegentropySession import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED -import kotlinx.coroutines.runBlocking import kotlinx.coroutines.withTimeout import okhttp3.OkHttpClient import okhttp3.Request @@ -41,27 +40,27 @@ import kotlin.test.fail /** * Equivalent of strfry's `test/syncTest.pl` driver for our interop - * tests. Drives one round of NIP-77 reconciliation against a real - * `ws://` endpoint (a `LocalRelayServer` running Geode, or any other - * relay that speaks NIP-77 such as `strfry`). + * tests. Drives one round of NIP-77 *negotiation* (NEG-OPEN / + * NEG-MSG / NEG-CLOSE) against a real `ws://` endpoint — a Geode + * `LocalRelayServer`, a strfry process, or any other NIP-77 relay. * * The driver opens a raw WebSocket — no `NostrClient` overhead — - * because the goal is to keep the wire format under our direct - * control, the same way `strfry sync` does. That way we exercise the - * server's NEG-OPEN / NEG-MSG / NEG-CLOSE handling with no client- - * side framing or filter-management indirection. + * to keep the wire format under direct control, the same way + * `strfry sync` does. That way we exercise the server's NEG-OPEN / + * NEG-MSG / NEG-CLOSE handling with no client-side framing or + * filter-management indirection. * - * Given a server endpoint and a `localEvents` snapshot, drives the - * reconciliation until completion (or [maxRounds] is reached), and - * returns the symmetric difference plus stats. + * Note: this driver only computes the symmetric difference. The + * actual *sync* (REQ for `needIds`, EVENT for `haveIds`) is the + * caller's job; that's a NIP-01 follow-up, not part of NIP-77. */ class InteropSyncDriver( private val httpClient: OkHttpClient = OkHttpClient.Builder().build(), ) { /** - * Reconciles `localEvents` against the relay at [wsUrl] under - * [filter]. Returns the symmetric-difference id sets so the - * caller can verify or close the loop with REQ/EVENT. + * Negotiates the symmetric difference between `localEvents` and + * the relay at [wsUrl] under [filter]. Returns the id sets so + * the caller can close the loop with REQ / EVENT. * * @param wsUrl the source relay's `ws://…` URL. * @param filter NEG-OPEN filter — usually the broadest filter the @@ -72,12 +71,12 @@ class InteropSyncDriver( * here so the relay's configured cap (500_000 by default) is * what governs framing — same shape as `strfry sync`. * @param timeoutMs hard timeout on a single NEG-MSG round trip. - * @param maxRounds upper bound on round trips. Strfry's typical - * converge in ≤5 rounds for 100 k corpora; 64 is a generous + * @param maxRounds upper bound on round trips. Strfry typically + * converges in ≤5 rounds for 100 k corpora; 64 is a generous * safety net that catches pathological splits without hanging * tests forever. */ - fun reconcile( + suspend fun negotiate( wsUrl: String, filter: Filter, localEvents: List, @@ -128,10 +127,7 @@ class InteropSyncDriver( val needIds = mutableSetOf() var rounds = 0 while (rounds < maxRounds) { - val raw = - runBlocking { - withTimeout(timeoutMs) { incoming.receive() } - } + val raw = withTimeout(timeoutMs) { incoming.receive() } val msg = OptimizedJsonMapper.fromJsonToMessage(raw) rounds++ when (msg) { diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt index d2cb5e134b..fe1355160f 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/store/sqlite/SnapshotIdsForNegentropyTest.kt @@ -25,7 +25,6 @@ import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent import kotlin.test.Test import kotlin.test.assertEquals -import kotlin.test.assertTrue /** * Verifies the NIP-77 negentropy id-and-time projection against the @@ -92,9 +91,11 @@ class SnapshotIdsForNegentropyTest : BaseDBTest() { // cap = 10; we have 30 rows, so the result must be 11 // (cap + 1 sentinel) — matches strfry's `maxSyncEvents` // overflow-detection idiom. + // cap=10 with 30 rows → result must be the +1 sentinel + // (11 rows). Caller compares `size > cap` to detect + // overflow — matches strfry's `maxSyncEvents` idiom. val capped = db.snapshotIdsForNegentropy(listOf(filter), maxEntries = 10) assertEquals(11, capped.size) - assertTrue(capped.size > 10, "caller relies on size > cap as overflow signal") // cap >= total: returns the whole set unchanged. val whole = db.snapshotIdsForNegentropy(listOf(filter), maxEntries = 100)