diff --git a/geode/config.example.toml b/geode/config.example.toml index 2ba85e0ad4..4651105623 100644 --- a/geode/config.example.toml +++ b/geode/config.example.toml @@ -112,6 +112,15 @@ require_auth = false # the future. Enforced by RejectFutureEventsPolicy. # reject_future_seconds = 1800 +# Path for the JSON file that remembers what the [[mirror]] catch-up +# has already synced (per upstream, per scope), so a restart resumes +# instead of re-downloading each upstream's whole backfill window. +# Defaults to ".sync-coverage.json" next to the event +# store; only written when the store itself is file-backed (an +# in-memory store keeps no resume state — saved coverage would +# describe events that no longer exist). +# mirror_sync_state_file = "/var/lib/geode/events.db.sync-coverage.json" + [authorization] # Allow / deny lists. Allow is a permissive ceiling; deny still # removes specific entries inside it. Enforced by Pubkey/KindAllowDenyPolicy. diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt index 2026a7f3ba..1ed17bbcdb 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt @@ -45,6 +45,7 @@ import com.vitorpamplona.quartz.nip01Core.store.IEventStore import com.vitorpamplona.quartz.nip01Core.store.NdjsonImportExport import com.vitorpamplona.quartz.nip01Core.store.sqlite.EventStore import com.vitorpamplona.quartz.nip77Negentropy.NegentropySettings +import com.vitorpamplona.quartz.utils.Log import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers @@ -224,7 +225,8 @@ private fun runImport(args: Array) { } System.err.println( "geode import: read=${stats.read} imported=${stats.imported} " + - "rejected=${stats.rejected} invalid-sig=${stats.invalid} malformed=${stats.malformed} " + + "rejected=${stats.rejected} failed=${stats.failed} " + + "invalid-sig=${stats.invalid} malformed=${stats.malformed} " + "→ ${ctx.dbFile ?: "(in-memory — not persisted; pass --db)"}", ) } finally { @@ -413,14 +415,25 @@ private fun serve(args: Array) { "[[mirror]] must not list this relay's own URL ($advertisedUrl)" } // Resume state for the mirror catch-up, following the admin state-file - // convention: next to the event database unless configured. An in-memory - // store keeps none — bands only pay off across restarts. + // convention: next to the event database unless configured. Keyed to the + // store's ACTUAL persistence, not to `database.file` being set: a + // volatile store with a persistent coverage file would claim, on the + // next boot, that an empty database already holds the backfill window — + // and the mirror would never fetch it. + val sqliteBackend = + config.database.backend + .trim() + .lowercase() in StoreFactory.SQLITE_BACKEND_KEYWORDS + val persistentLocation = + a.opt("--db") ?: config.database.file?.takeUnless { sqliteBackend && config.database.in_memory } val syncCoverage = - if (upstreams.isEmpty()) { + if (upstreams.isEmpty() || persistentLocation == null) { + if (upstreams.isNotEmpty() && config.options.mirror_sync_state_file != null) { + Log.w("Main") { "mirror_sync_state_file ignored: the event store is in-memory, so saved coverage would outlive the events it describes" } + } null } else { - (config.options.mirror_sync_state_file ?: config.database.file?.let { "$it.sync-coverage.json" }) - ?.let { SyncCoverageFile(File(it)) } + SyncCoverageFile(File(config.options.mirror_sync_state_file ?: "$persistentLocation.sync-coverage.json")) } val mirror = if (upstreams.isEmpty()) { diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt index 5916a5f794..c98e050e37 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt @@ -23,9 +23,11 @@ package com.vitorpamplona.geode.mirror import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropySyncException import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.SyncCoverage +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPages import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropyReconcile -import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySyncOrFetch +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropySync import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd @@ -41,12 +43,15 @@ import com.vitorpamplona.quartz.utils.TimeUtils import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.trySendBlocking import kotlinx.coroutines.delay import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeoutOrNull import okhttp3.OkHttpClient import java.time.Duration import java.util.concurrent.atomic.AtomicLong @@ -155,9 +160,9 @@ class MirrorWorker( * historical window before the live REQ tail takes over. The geode binary * turns this on (see `Main`); it defaults **off** so the many existing * MirrorWorker tests keep exercising the pure live-REQ path unchanged. - * When on, `negentropySyncOrFetch` automatically falls back to paged REQ - * against an upstream that doesn't speak NIP-77 — so "either mode" is - * transparent and needs no separate toggle. + * When on, the catch-up automatically falls back to paged REQ against an + * upstream that doesn't speak NIP-77 — so "either mode" is transparent + * and needs no separate toggle. */ private val negentropyBackfill: Boolean = false, /** @@ -422,9 +427,9 @@ class MirrorWorker( * **client-paced** so a fast upstream can't overrun the sink: a plain REQ * backfill of a large set dies here (strfry kills a slow REQ client once its * unsent-outbound buffer crosses `maxPendingOutboundBytes`), which is exactly - * why this uses negentropy. [INostrClient.negentropySyncOrFetch] falls back - * to paged REQ automatically when the upstream doesn't speak NIP-77, so the - * mirror is compatible with either kind of upstream with no config. + * why this uses negentropy. When the upstream doesn't speak NIP-77 the + * catch-up falls back to paged REQ, so the mirror is compatible with + * either kind of upstream with no config. * * A failure here is non-fatal: the live subscription keeps the mirror current * and the reconnect watermark narrows any residual gap. @@ -437,13 +442,12 @@ class MirrorWorker( ) { // Resume memory: coverage is keyed on the STABLE scoped filter, never // on the boot window — whose since/until change every start and would - // never match a stored band. The window only slides forward - // (initialSince = boot − backfillSeconds), so a band recorded against - // an earlier boot's lower floor never licenses skipping a range this - // boot can ask about but the last one could not. + // never match a stored band. The window itself rides along as the + // floor argument, so a band recorded against a shallower window + // re-opens the older span when the operator deepens the backfill. val legs = coverage - ?.legs(up.url, scopedBase) + ?.legs(up.url, scopedBase, initialSince) ?.mapNotNull { clampToWindow(it, initialSince, until) } ?: listOf(scopedBase.copy(since = initialSince, until = until)) if (legs.isEmpty()) { @@ -452,9 +456,10 @@ class MirrorWorker( } // Bounded hand-off → one ingest consumer. `onEvent` can't suspend, so it - // blocks here when the sink falls behind; because negentropySyncOrFetch's - // own delivery pipeline is bounded, that backpressure reaches all the way - // to the upstream — no unbounded buffering (unlike the live-tail path). + // blocks here when the sink falls behind; because negentropySync's own + // delivery pipeline is bounded (and a paged REQ is paced by its pages), + // that backpressure reaches all the way to the upstream — no unbounded + // buffering (unlike the live-tail path). val handoff = Channel(capacity = CATCHUP_HANDOFF) val consumer = scope.launch { @@ -493,41 +498,81 @@ class MirrorWorker( val syncStartedAt = TimeUtils.now() var seenMin: Long? = null var seenMax: Long? = null - val result = - client.negentropySyncOrFetch( - relay = up.url, - filter = leg, - localEntries = localEntries, - onEvent = { event -> - // Same containment as the live path: even a trusted - // upstream may only inject events inside the declared scope. - if (up.filter == null || up.filter.match(event)) { - // Only plausible stamps widen a band — one - // misdated event must not discard the rest. - if (SyncCoverage.isPlausible(event.createdAt)) { - seenMin = minOf(seenMin ?: event.createdAt, event.createdAt) - seenMax = maxOf(seenMax ?: event.createdAt, event.createdAt) - } - handoff.trySendBlocking(event) - } else { - filtered.incrementAndGet() - } - }, - ) - downloaded += result.downloaded - paged = paged || result.pagedFallback + + fun observe(event: Event) { + // Same containment as the live path: even a trusted + // upstream may only inject events inside the declared scope. + if (up.filter == null || up.filter.match(event)) { + // Only plausible stamps widen a band — one + // misdated event must not discard the rest. + if (SyncCoverage.isPlausible(event.createdAt)) { + seenMin = minOf(seenMin ?: event.createdAt, event.createdAt) + seenMax = maxOf(seenMax ?: event.createdAt, event.createdAt) + } + handoff.trySendBlocking(event) + } else { + filtered.incrementAndGet() + } + } + + // The two phases run by hand rather than through + // negentropySyncOrFetch, for two reasons. The combinator keeps + // every delivered id for cross-phase dedup — a multi-million- + // event catch-up cannot afford that heap, and the store's + // unique-id constraint dedups anyway. And the fallback must + // reset the observed span: a half-finished reconcile delivers + // events scattered across the whole leg, and a band built from + // that scatter would claim interior ranges nobody walked. + val legPaged = + try { + downloaded += + client + .negentropySync( + relay = up.url, + filter = leg, + localEntries = localEntries, + onEvent = ::observe, + ).downloaded + false + } catch (e: NegentropySyncException) { + seenMin = null + seenMax = null + // The watchdog matches negentropySync's default rather + // than fetchAllPages' shorter one: a paged catch-up + // sits behind the same slow upstreams. + downloaded += client.fetchAllPages(up.url, listOf(leg), idleTimeoutMs = 120_000L) { observe(it) } + true + } + paged = paged || legPaged // Recorded per leg, so a failure between legs keeps the ground - // the first one gained. A clean reconcile is complete through - // the instant its snapshot was read; a paged fallback earns - // only the span it actually saw. - coverage?.record( - up.url, - scopedBase, - seenMin, - seenMax, - paged = result.pagedFallback, - reconciledThrough = if (result.pagedFallback) null else syncStartedAt, - ) + // the first one gained — and no more: a reconcile compared only + // its own leg, so completeness reaches the leg's ceiling, never + // "now" while a later leg is still pending. A paged fallback + // earns only the span it actually saw, capped at the snapshot + // instant so one future-dated event cannot lift the band's + // ceiling past what was asked. + if (legPaged) { + coverage?.record( + up.url, + scopedBase, + seenMin, + seenMax?.coerceAtMost(syncStartedAt), + paged = true, + ) + } else { + val legFloor = leg.since ?: initialSince + coverage?.record( + up.url, + scopedBase, + // The compared range starts at the leg's floor whether + // or not anything was observed there — that is what a + // clean reconcile proves. + observedMin = minOf(seenMin ?: legFloor, legFloor), + observedMax = null, + paged = false, + reconciledThrough = minOf(leg.until ?: syncStartedAt, syncStartedAt), + ) + } } Log.i("MirrorWorker") { val how = if (paged) "paged REQ (upstream has no NIP-77)" else "negentropy" @@ -743,11 +788,21 @@ class MirrorWorker( runCatching { client.close() } inbound.close() scope.cancel() + // Wait (bounded) for the workers to land: a coverage.record racing + // past the state file's final flush would be recorded and lost. + // Bounded because a worker parked in a blocking hand-off does not + // feel the cancel, and shutdown must not hang on it. + runBlocking { + withTimeoutOrNull(CLOSE_JOIN_MS) { scope.coroutineContext[Job]?.join() } + } okhttp?.dispatcher?.executorService?.shutdown() okhttp?.connectionPool?.evictAll() } private companion object { + /** How long [close] waits for the worker coroutines to land. */ + const val CLOSE_JOIN_MS = 5_000L + /** Matches the Android app's relay-pool WebSocket ping interval. */ const val PING_INTERVAL_SECS = 120L @@ -773,7 +828,7 @@ class MirrorWorker( /** * Depth of the catch-up hand-off between the negentropy download and the - * ingest consumer. Small: negentropySyncOrFetch is already internally + * ingest consumer. Small: the negentropy download is already internally * backpressured, so this only smooths the seam — the bounded IngestQueue * behind `server.ingest` is the real limiter. */ diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt index f02d8f56f5..b31c51a60d 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt @@ -31,6 +31,7 @@ import kotlinx.serialization.json.jsonPrimitive import kotlinx.serialization.json.long import kotlinx.serialization.json.put import java.io.File +import java.nio.file.AtomicMoveNotSupportedException import java.nio.file.Files import java.nio.file.StandardCopyOption @@ -58,7 +59,8 @@ class SyncCoverageFile( init { load() - // Loading marks every restored band dirty; the file already has them. + // restore() bypasses onChange, but stay defensive: reopening a file + // must never count as a change, or every boot rewrites it. dirty = false flusher = Thread { @@ -130,7 +132,16 @@ class SyncCoverageFile( file.parentFile?.mkdirs() val tmp = File(file.parentFile ?: File("."), "${file.name}.tmp") tmp.writeText(json.encodeToString(JsonObject.serializer(), doc)) - Files.move(tmp.toPath(), file.toPath(), StandardCopyOption.REPLACE_EXISTING) + // ATOMIC_MOVE requested explicitly: without it the JVM may + // legally fall back to copy+delete, and a reader could see a + // half map. Same-directory rename, so support is the norm; a + // filesystem that truly can't gets the plain move (and the + // corrupt-file recovery absorbs the residual risk). + try { + Files.move(tmp.toPath(), file.toPath(), StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE) + } catch (_: AtomicMoveNotSupportedException) { + Files.move(tmp.toPath(), file.toPath(), StandardCopyOption.REPLACE_EXISTING) + } }.onFailure { Log.w("SyncCoverageFile") { "could not write ${file.path}: ${it.message}" } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt index c6ddff4dce..35ba9b330a 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt @@ -87,8 +87,10 @@ class SyncCoverage( // filter -> its canonical json. Filter.toJson() runs to tens of thousands // of characters for author-scoped filters, and a fan-out keys once per // relay per cycle over the SAME handful of filter instances. Filter - // compares by identity, so this map is an identity cache; an - // equal-but-distinct filter still keys correctly, just without the cache. + // compares by identity, so this map is an identity cache — and an + // identity cache retains every distinct instance it is handed. A caller + // that rebuilds its filter each cycle would grow it forever, so past + // MAX_FINGERPRINTS new instances key correctly but are not cached. private val fingerprints = ConcurrentMap() /** @@ -106,6 +108,7 @@ class SyncCoverage( fun legs( url: NormalizedRelayUrl, filter: Filter, + floor: Long? = null, ): List { val band = bands[key(url, filter)] ?: return listOf(filter) // Time for another full pass: relays gain old events, and without @@ -114,9 +117,20 @@ class SyncCoverage( val legs = mutableListOf() // Older: up to and including the band's floor, but not past the - // filter's. A complete band has no older leg at all — the reconcile - // already compared the whole range. - if (!band.complete && (filter.since == null || band.minCreatedAt >= filter.since)) { + // filter's (or, when the filter has no `since`, the caller's + // [floor] — a sync window the filter itself must not carry, or it + // would change the band's key every run). A complete band compared + // its whole range already, but only down to the floor it ran + // against: a caller now reaching deeper — a raised backfill window + // — re-opens the span below the band. + val since = filter.since ?: floor + val wantsOlder = + if (band.complete) { + since != null && since < band.minCreatedAt + } else { + since == null || band.minCreatedAt >= since + } + if (wantsOlder) { legs.add(filter.copy(until = minOf(band.minCreatedAt, filter.until ?: Long.MAX_VALUE))) } @@ -211,12 +225,17 @@ class SyncCoverage( var since = Long.MAX_VALUE for (url in urls) { val legs = legs(url, filter) + // Nothing outside its band: this relay asks nothing of the + // snapshot at all — the best case must not widen the window. + if (legs.isEmpty()) continue // More than one leg means an older gap this relay still wants, so // the snapshot cannot start above the filter's own floor. val only = legs.singleOrNull() ?: return filter val legSince = only.since ?: return filter since = minOf(since, legSince) } + // Every relay fully covered: any window would do; the unnarrowed + // filter is merely safe, and callers usually skip the sync entirely. return if (since == Long.MAX_VALUE) filter else filter.copy(since = since) } @@ -245,9 +264,23 @@ class SyncCoverage( private fun key( url: NormalizedRelayUrl, filter: Filter, - ): String = "${url.url} ${fingerprints.getOrPut(filter) { filter.toJson() }}" + ): String { + val fingerprint = + fingerprints[filter] + ?: filter.toJson().also { + // Bounded: stable callers hit the cache at any size a real + // config produces; a caller minting fresh instances just + // pays the toJson each time instead of growing the heap. + if (fingerprints.size() < MAX_FINGERPRINTS) fingerprints[filter] = it + } + return "${url.url} $fingerprint" + } companion object { + // More filter instances than any deliberate configuration holds; only + // a caller rebuilding filters per cycle ever reaches it. + private const val MAX_FINGERPRINTS = 1_000 + /** * A week. Long enough that the narrow path is the normal one, short * enough that anything a band is wrong about is wrong for days, not diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt index 0f515d0ca5..348a27544d 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt @@ -62,16 +62,25 @@ class PagingWindowProgress( private val windows = ConcurrentMap() - /** Begin a pagination over `[bottom, top]` seconds. An inverted window is not one. */ + /** + * Begin a pagination over `[bottom, top]` seconds. An inverted window is + * not one; a single-second window (`top == bottom`) is — coverage legs + * that re-read a band's edge second are exactly that shape. + */ fun begin( key: String, top: Long, bottom: Long, ) { - if (top > bottom) windows[key] = Window(top, bottom, nowMillis(), top) + if (top >= bottom) windows[key] = Window(top, bottom, nowMillis(), top) } - /** The pagination reached [until]; monotonic, so a page that jumps back cannot un-advance it. */ + /** + * The pagination reached [until]; monotonic, so a page that jumps back + * cannot un-advance it. The check-then-set is unsynchronized on purpose: + * one pagination is one coroutine, and a display racing a mark can only + * ever read a value one page stale. + */ fun mark( key: String, until: Long, diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/IngestQueue.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/IngestQueue.kt index d955ed927c..18c0dfa34e 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/IngestQueue.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/server/backend/IngestQueue.kt @@ -36,6 +36,7 @@ import kotlin.concurrent.atomics.AtomicBoolean import kotlin.concurrent.atomics.AtomicInt import kotlin.concurrent.atomics.ExperimentalAtomicApi import kotlin.coroutines.CoroutineContext +import kotlin.coroutines.cancellation.CancellationException /** * Group-commit writer for incoming EVENT publishes. @@ -290,6 +291,11 @@ class IngestQueue( } else { try { store.batchInsert(toInsert) + } catch (e: CancellationException) { + // Shutdown, not a store failure: rethrow so the loop + // stops instead of stamping the batch Failed and + // carrying on while cancelled. + throw e } catch (e: Throwable) { Log.w("IngestQueue") { "batchInsert failed for ${toInsert.size} events: ${e.message}" } val reason = e.message ?: e::class.simpleName ?: "insert failed" 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 16b82132e9..ae49e25749 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 @@ -27,6 +27,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent +import kotlin.coroutines.cancellation.CancellationException /** * Storage contract for Nostr events: insert, filter-query, count, delete, @@ -136,16 +137,20 @@ interface IEventStore : AutoCloseable { * * Default impl runs each insert in its own transaction — correct * but loses the group-commit win — and cannot classify a throw from - * [insert], so it reports `Rejected`. Implementations that can tell - * a refusal from a write error should override and say which. + * [insert], so it reports `Failed`: re-offering a duplicate is + * idempotent, while dropping a good event on a transient store + * error is not. Implementations that can tell a refusal from a + * write error should override and say which. */ suspend fun batchInsert(events: List): List = events.map { event -> try { insert(event) InsertOutcome.Accepted + } catch (e: CancellationException) { + throw e } catch (e: Throwable) { - InsertOutcome.Rejected(e.message ?: e::class.simpleName ?: "insert failed") + InsertOutcome.Failed(e.message ?: e::class.simpleName ?: "insert failed") } } 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 c314c4e3b2..71a6e3c292 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 @@ -103,7 +103,7 @@ class ObservableEventStore( // order. Already-expired ephemerals are dropped (matching // [insert]). Accepted events are emitted on [_changes] only // after the inner batch returns, so a commit failure that - // converts everything to Rejected suppresses the emits. + // converts everything to Failed suppresses the emits. if (events.isEmpty()) return emptyList() val outcomes = arrayOfNulls(events.size) 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 93f3bb5bb6..3db879925d 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 @@ -510,7 +510,29 @@ class SQLiteEventStore( // ROLLBACK shouldn't mask the original cause. runCatching { db.execSQL("ROLLBACK TRANSACTION TO SAVEPOINT $sp") } runCatching { db.execSQL("RELEASE SAVEPOINT $sp") } - IEventStore.InsertOutcome.Rejected(e.message ?: e::class.simpleName ?: "insert failed") + classifyRowError(e) + } + } + + /** + * Which side failed decides whether the caller may drop the event. + * Policy refusals are recognizable — every schema trigger RAISEs with + * a `blocked:` prefix, the immutability guards say "not allowed", and + * a duplicate id is a constraint violation. Anything else (disk full, + * I/O error, schema drift) is the store failing to write an acceptable + * event: `Failed`, so a rising count is loud instead of blending into + * the duplicate tally. + */ + private fun classifyRowError(e: Throwable): IEventStore.InsertOutcome { + val message = e.message ?: e::class.simpleName ?: "insert failed" + val refusal = + message.contains("blocked:") || + message.contains("not allowed") || + message.contains("constraint", ignoreCase = true) + return if (refusal) { + IEventStore.InsertOutcome.Rejected(message) + } else { + IEventStore.InsertOutcome.Failed(message) } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/HostStrikes.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/HostStrikes.kt index 11a2c2d51d..1b846eac93 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/HostStrikes.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip66RelayMonitor/reachability/HostStrikes.kt @@ -82,7 +82,13 @@ class HostStrikes( val authority = authorityOf(url.url) if (authority in producedHosts || authority in deadHosts) return null if (strikes.merge(authority, 1) { old, new -> old + new } < strikeLimit) return null - deadHosts.add(authority) + // Concurrent strikers can cross the threshold together; add() is the + // atomic exactly-once gate on who publishes. Re-check produced after + // winning it: a delivery that landed while this strike was in flight + // outranks the verdict, and a verdict for a host that just answered + // would be a false public record. + if (!deadHosts.add(authority)) return null + if (authority in producedHosts) return null return Evicted(authority, strikeLimit) } diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt index 78abe61682..53120624a4 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt @@ -159,6 +159,22 @@ class SyncCoverageTest { assertEquals(2, walked.legs(relay, profiles).size, "a paged walk says nothing about what it never asked for") } + @Test + fun `a deeper floor re-opens history below a complete band`() { + // A reconcile only compared down to the window it ran against. When + // the operator raises the backfill window, the span below the band's + // recorded floor is ground nobody ever asked for. + val c = SyncCoverage() + c.record(relay, profiles, 1_700_000_000L, null, paged = false, reconciledThrough = 1_700_002_000L) + + assertEquals(1, c.legs(relay, profiles, floor = 1_700_000_000L).size, "same floor: nothing older to ask") + + val legs = c.legs(relay, profiles, floor = 1_600_000_000L) + assertEquals(2, legs.size, "a deeper floor re-opens the older span") + assertEquals(1_700_000_000L, legs[0].until, "up to the floor the reconcile actually compared") + assertEquals(1_700_002_000L, legs[1].since) + } + // ---- the periodic full re-walk ----------------------------------------- @Test @@ -224,6 +240,20 @@ class SyncCoverageTest { assertEquals(1_700_003_000L, c.coveringWindow(listOf(relay, other, third), profiles).since) } + @Test + fun `a fully covered relay does not widen the shared window`() { + // A complete band past a bounded filter's ceiling needs no legs at + // all. The best case must not force the snapshot back to the whole + // filter — that would make full coverage cost the most. + val window = Filter(kinds = listOf(0), since = 1_700_000_000L, until = 1_700_005_000L) + val c = SyncCoverage() + c.record(relay, window, 1_700_000_000L, null, paged = false, reconciledThrough = 1_700_009_000L) + c.record(other, window, 1_700_000_000L, null, paged = false, reconciledThrough = 1_700_003_000L) + + assertEquals(0, c.legs(relay, window).size, "covered past the ceiling: nothing to ask") + assertEquals(1_700_003_000L, c.coveringWindow(listOf(relay, other), window).since) + } + @Test fun `a relay with an older gap also widens the shared window`() { val c = SyncCoverage() diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt index 2c1c0ed433..35c5ecfdab 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt @@ -39,7 +39,7 @@ class PagingWindowProgressTest { val p = PagingWindowProgress() p.begin("a", top = 1_000L, bottom = 0L) - assertClose(0.0, p.fraction(), "nothing walked yet") + assertClose(0.0, p.fraction(), "nothing paged yet") p.mark("a", 750L) assertClose(0.25, p.fraction()) @@ -64,7 +64,7 @@ class PagingWindowProgressTest { @Test fun `windows average rather than sum`() { - // Two relays each walking their own window: one done and one untouched + // Two relays each paging their own window: one done and one untouched // is half way — not 100% as summing would give. val p = PagingWindowProgress() p.begin("a", top = 1_000L, bottom = 0L) @@ -84,14 +84,14 @@ class PagingWindowProgressTest { p.finish("a") - assertClose(0.5, p.fraction(), "only b is still walking") + assertClose(0.5, p.fraction(), "only b is still paging") p.finish("b") assertNull(p.fraction(), "nothing paging means no number to report") } @Test - fun `a group prefix scopes the numbers to its own walks`() { - // One instance serves many concurrent walks; without the scope two + fun `a group prefix scopes the numbers to its own paginations`() { + // One instance serves many concurrent paginations; without the scope two // streams would print each other's percentages. val p = PagingWindowProgress() p.begin("streamA|wss://r1", top = 1_000L, bottom = 0L) @@ -104,7 +104,7 @@ class PagingWindowProgressTest { } @Test - fun `an inverted or empty window is not a pagination`() { + fun `an inverted window is not a pagination`() { // A leg whose since is above its until asks for a range nothing can be // in. Dividing by that span would produce infinities on the status line. val p = PagingWindowProgress() @@ -114,6 +114,21 @@ class PagingWindowProgressTest { assertNull(p.fraction()) } + @Test + fun `a single-second window is a pagination`() { + // Coverage legs re-read a band's edge second: since == until is a + // real, one-second range, not an inverted one. + val p = PagingWindowProgress() + + p.begin("a", top = 500L, bottom = 500L) + + assertClose(0.0, p.fraction(), "tracked from its start") + p.mark("a", 500L) + assertClose(0.0, p.fraction(), "still at its only second") + p.finish("a") + assertNull(p.fraction()) + } + @Test fun `no ETA before the estimate means anything`() { val p = PagingWindowProgress() diff --git a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/ConcurrentIngestLossTest.kt b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/ConcurrentIngestLossTest.kt index 335d587cc6..3d87b2f457 100644 --- a/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/ConcurrentIngestLossTest.kt +++ b/quartz/src/jvmTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/prodbench/ConcurrentIngestLossTest.kt @@ -120,6 +120,7 @@ class ConcurrentIngestLossTest { val accepted = ConcurrentHashMap.newKeySet() val rejected = AtomicInteger() + val failed = AtomicInteger() val window = Semaphore(200) val done = CompletableDeferred() val remaining = AtomicInteger(events.size) @@ -130,7 +131,7 @@ class ConcurrentIngestLossTest { when (outcome) { is IEventStore.InsertOutcome.Accepted -> accepted.add(e.id) is IEventStore.InsertOutcome.Rejected -> rejected.incrementAndGet() - is IEventStore.InsertOutcome.Failed -> rejected.incrementAndGet() + is IEventStore.InsertOutcome.Failed -> failed.incrementAndGet() } window.release() if (remaining.decrementAndGet() == 0) done.complete(Unit) @@ -155,7 +156,7 @@ class ConcurrentIngestLossTest { } } val stored = store.count(Filter()) - println(" submitted=${events.size} accepted=${accepted.size} rejected=${rejected.get()} stored=$stored lostAcceptedRegular=$lost") + println(" submitted=${events.size} accepted=${accepted.size} rejected=${rejected.get()} failed=${failed.get()} stored=$stored lostAcceptedRegular=$lost") ingest.close() queueJob.cancel()