From 19d8e3069de5ca5144307bb3ab1ba153d24a1811 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 4 Aug 2026 04:24:19 +0000 Subject: [PATCH] Align the sync accessories with quartz vocabulary; geode catch-up resumes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Renames from review: SyncBands -> SyncCoverage ("sync" reads negentropy-ish in quartz, and coverage is the role — the bands are the records), and PagingProgress moves into relay.client.paging as PagingWindowProgress, beside RelayLoadingCursors and RelayPagingProgress, with its docs swept from "walk" dialect to quartz's pagination vocabulary and cross-references delineating the three: cursors are in-memory positions for demand-driven UI paging, the window progress is fraction/ETA for a bulk pagination over a known window, coverage is persistent intervals that license skipping work. geode adopts both halves of the new contract. MirrorWorker counts InsertOutcome.Failed in its own `failed` counter instead of folding it into `rejected`, and the down catch-up gains resume memory: SyncCoverageFile persists SyncCoverage next to the event database (admin state-file convention, temp-file + atomic move, daemon flush), and runCatchUpDown asks only for the legs outside the covered band. Bands are keyed on the stable scoped filter — never the boot window, whose since/until change every start — and clamped to the window, which only slides forward, so an old band can never license skipping a range an earlier boot could not ask about. A clean reconcile records completeness through its snapshot instant; a paged fallback earns only the span it saw. For an upstream without NIP-77 this turns the every-boot full re-download of the backfill window into a resumed walk. Off unless wired: MirrorWorker's coverage parameter defaults to null and in-memory stores keep no state file, so existing tests and setups are unchanged. Full :quartz:jvmTest and :geode:test pass. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01Y4Pi9YYMhdzTFxRiV2jF9R --- .../kotlin/com/vitorpamplona/geode/Main.kt | 16 ++ .../geode/config/StaticConfig.kt | 10 ++ .../geode/mirror/MirrorWorker.kt | 124 ++++++++++++--- .../geode/mirror/SyncCoverageFile.kt | 147 ++++++++++++++++++ .../geode/mirror/SyncCoverageFileTest.kt | 130 ++++++++++++++++ .../{SyncBands.kt => SyncCoverage.kt} | 7 +- .../PagingWindowProgress.kt} | 60 +++---- .../{SyncBandsTest.kt => SyncCoverageTest.kt} | 58 +++---- .../PagingWindowProgressTest.kt} | 32 ++-- 9 files changed, 489 insertions(+), 95 deletions(-) create mode 100644 geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt create mode 100644 geode/src/test/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFileTest.kt rename quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/{SyncBands.kt => SyncCoverage.kt} (97%) rename quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/{accessories/PagingProgress.kt => paging/PagingWindowProgress.kt} (62%) rename quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/{SyncBandsTest.kt => SyncCoverageTest.kt} (93%) rename quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/{accessories/PagingProgressTest.kt => paging/PagingWindowProgressTest.kt} (84%) diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt index 729b95b3a3..2026a7f3ba 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/Main.kt @@ -28,6 +28,7 @@ import com.vitorpamplona.geode.config.StaticConfig import com.vitorpamplona.geode.mirror.MirrorDirection import com.vitorpamplona.geode.mirror.MirrorUpstream import com.vitorpamplona.geode.mirror.MirrorWorker +import com.vitorpamplona.geode.mirror.SyncCoverageFile import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl @@ -411,6 +412,16 @@ private fun serve(args: Array) { require(upstreams.none { it.url.displayUrl() == advertisedIdentity }) { "[[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. + val syncCoverage = + if (upstreams.isEmpty()) { + null + } else { + (config.options.mirror_sync_state_file ?: config.database.file?.let { "$it.sync-coverage.json" }) + ?.let { SyncCoverageFile(File(it)) } + } val mirror = if (upstreams.isEmpty()) { null @@ -425,6 +436,9 @@ private fun serve(args: Array) { // for the historical window, then live REQ tail. Auto-falls back // to paged REQ for upstreams without NIP-77. negentropyBackfill = true, + // Without NIP-77 the catch-up is a paged re-download; the + // coverage bands remember what previous boots already walked. + coverage = syncCoverage?.coverage, ).also { it.start() } } @@ -462,6 +476,8 @@ private fun serve(args: Array) { // queue and store beneath them shut down. runCatching { maintenanceScope.cancel() } runCatching { mirror?.close() } + // After the mirror, so the final flush carries the last bands. + runCatching { syncCoverage?.close() } runCatching { server.stop() } runCatching { relay.close() } }, diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt index dffe4e50e3..8cfecb4e41 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/config/StaticConfig.kt @@ -185,6 +185,16 @@ data class StaticConfig( * that don't offer NIP-50 at all (strfry, for example). */ val full_text_search: Boolean = true, + /** + * Where the mirror catch-up's resume state lives: the per-upstream + * `created_at` coverage bands (quartz's `SyncCoverage`). Without it + * every restart re-syncs each upstream's whole backfill window — + * a full re-download for an upstream without NIP-77. Defaults to + * `.sync-coverage.json` when the store is + * file-backed; an in-memory store keeps no resume state (bands + * only pay off across restarts). + */ + val mirror_sync_state_file: String? = null, ) /** 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 942fe7769f..5916a5f794 100644 --- a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/MirrorWorker.kt @@ -23,6 +23,7 @@ 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.SyncCoverage 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.reqs.SubscriptionListener @@ -159,6 +160,13 @@ class MirrorWorker( * transparent and needs no separate toggle. */ private val negentropyBackfill: Boolean = false, + /** + * Resume memory for the down catch-up, shared across upstreams and — via + * [SyncCoverageFile] — across restarts. Null keeps the old behavior: + * every boot re-syncs the whole backfill window, which for an upstream + * without NIP-77 is a full re-download. + */ + private val coverage: SyncCoverage? = null, ) : AutoCloseable { private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) @@ -201,6 +209,14 @@ class MirrorWorker( /** Events the store rejected — mostly duplicate replays after a reconnect. */ val rejected = AtomicLong(0) + /** + * Good events the store could not write ([IEventStore.InsertOutcome.Failed]). + * The store's fault, not the event's, and nothing re-offers them — counted + * apart from [rejected] so a schema drift reads as store damage rather + * than as upstreams sending junk. + */ + val failed = AtomicLong(0) + /** * Deliveries dropped before ever reaching the store: events outside * the [MirrorUpstream.filter] scope (an upstream answering outside @@ -304,7 +320,7 @@ class MirrorWorker( // The store's error, not the event's: the event // was good and nothing will re-offer it. Louder // than a rejection on purpose. - rejected.incrementAndGet() + failed.incrementAndGet() Log.w("MirrorWorker") { "store failed ${msg.event.id}: ${outcome.reason}" } } } @@ -419,11 +435,21 @@ class MirrorWorker( initialSince: Long, until: Long, ) { - val catchUpFilter = scopedBase.copy(since = initialSince, until = until) - // Reconcile against what we already hold in this window → download only - // the diff (like `strfry sync`). No store wired → empty local set → the - // whole window is downloaded and the store's unique-id constraint dedups. - val localEntries = store?.snapshotIdsForNegentropy(listOf(catchUpFilter)) ?: emptyList() + // 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. + val legs = + coverage + ?.legs(up.url, scopedBase) + ?.mapNotNull { clampToWindow(it, initialSince, until) } + ?: listOf(scopedBase.copy(since = initialSince, until = until)) + if (legs.isEmpty()) { + Log.i("MirrorWorker") { "catch-up from ${up.url.url}: window already covered - nothing outside the synced band" } + return + } // Bounded hand-off → one ingest consumer. `onEvent` can't suspend, so it // blocks here when the sink falls behind; because negentropySyncOrFetch's @@ -439,7 +465,7 @@ class MirrorWorker( IEventStore.InsertOutcome.Accepted -> accepted.incrementAndGet() is IEventStore.InsertOutcome.Rejected -> rejected.incrementAndGet() is IEventStore.InsertOutcome.Failed -> { - rejected.incrementAndGet() + failed.incrementAndGet() Log.w("MirrorWorker") { "store failed ${event.id}: ${outcome.reason}" } } } @@ -454,24 +480,59 @@ class MirrorWorker( } try { - val result = - client.negentropySyncOrFetch( - relay = up.url, - filter = catchUpFilter, - 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)) { - handoff.trySendBlocking(event) - } else { - filtered.incrementAndGet() - } - }, + var downloaded = 0 + var paged = false + for (leg in legs) { + // Reconcile against what we already hold in this leg → download + // only the diff (like `strfry sync`). No store wired → empty + // local set → the whole leg is downloaded and the store's + // unique-id constraint dedups. + val localEntries = store?.snapshotIdsForNegentropy(listOf(leg)) ?: emptyList() + // Coverage is stamped from when the local ids were read — that + // is the state the relay is being compared against. + 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 + // 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, ) + } Log.i("MirrorWorker") { - val how = if (result.pagedFallback) "paged REQ (upstream has no NIP-77)" else "negentropy" - "catch-up from ${up.url.url}: ${result.downloaded} events via $how" + val how = if (paged) "paged REQ (upstream has no NIP-77)" else "negentropy" + val resumed = if (coverage != null && legs.size > 1) " [resumed: ${legs.size} legs outside the synced band]" else "" + "catch-up from ${up.url.url}: $downloaded events via $how$resumed" } } catch (e: CancellationException) { throw e @@ -731,3 +792,20 @@ class MirrorWorker( const val UP_SYNC_SETTLE_MS = 1_500L } } + +/** + * [leg] intersected with the boot window `[since, until]`, or null when the + * band already covers everything this window could ask. Coverage legs come + * off the stable scoped filter and are unbounded on one side; the catch-up + * only ever asks inside its own window. + */ +internal fun clampToWindow( + leg: Filter, + since: Long, + until: Long, +): Filter? { + val newSince = maxOf(leg.since ?: since, since) + val newUntil = minOf(leg.until ?: until, until) + if (newSince > newUntil) return null + return leg.copy(since = newSince, until = newUntil) +} diff --git a/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt new file mode 100644 index 0000000000..f02d8f56f5 --- /dev/null +++ b/geode/src/main/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFile.kt @@ -0,0 +1,147 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.mirror + +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.SyncCoverage +import com.vitorpamplona.quartz.utils.Log +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.boolean +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import kotlinx.serialization.json.long +import kotlinx.serialization.json.put +import java.io.File +import java.nio.file.Files +import java.nio.file.StandardCopyOption + +/** + * File persistence for [SyncCoverage], so the mirror's catch-up resumes + * across restarts instead of re-syncing each upstream's whole backfill + * window — which, for an upstream without NIP-77, is a full re-download + * every boot. + * + * Same shape as the admin state file: JSON next to the event database, + * written via a temp file and an atomic move so a reader never sees a half + * map. A daemon timer flushes changed state so progress survives a hard + * kill; [close] flushes the rest. A corrupt file starts fresh — the cost of + * losing it is one re-sync, the cost of refusing to start is the relay. + */ +class SyncCoverageFile( + private val file: File, + flushSeconds: Long = DEFAULT_FLUSH_SECONDS, +) : AutoCloseable { + @Volatile private var dirty = false + + val coverage = SyncCoverage(onChange = { dirty = true }) + + private val flusher: Thread + + init { + load() + // Loading marks every restored band dirty; the file already has them. + dirty = false + flusher = + Thread { + while (!Thread.currentThread().isInterrupted) { + try { + Thread.sleep(flushSeconds * 1000) + } catch (_: InterruptedException) { + return@Thread + } + flush() + } + }.apply { + isDaemon = true + name = "mirror-sync-coverage-flush" + start() + } + } + + /** Write the map if anything changed since the last write. */ + @Synchronized + fun flush() { + if (!dirty) return + dirty = false + save() + } + + override fun close() { + flusher.interrupt() + flush() + } + + private fun load() { + if (!file.isFile) return + runCatching { + val root = Json.parseToJsonElement(file.readText()).jsonObject + coverage.restore( + root.mapValues { (_, v) -> + val o = v.jsonObject + SyncCoverage.Band( + o.getValue("min").jsonPrimitive.long, + o.getValue("max").jsonPrimitive.long, + o["complete"]?.jsonPrimitive?.boolean ?: false, + o["fullAt"]?.jsonPrimitive?.long ?: 0L, + ) + }, + ) + }.onFailure { + Log.w("SyncCoverageFile") { "could not read ${file.path} (${it.message}); starting fresh" } + } + } + + @Synchronized + private fun save() { + runCatching { + val doc = + buildJsonObject { + coverage.export().forEach { (key, band) -> + put( + key, + buildJsonObject { + put("min", band.minCreatedAt) + put("max", band.maxCreatedAt) + put("complete", band.complete) + put("fullAt", band.fullAt) + }, + ) + } + } + 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) + }.onFailure { + Log.w("SyncCoverageFile") { "could not write ${file.path}: ${it.message}" } + } + } + + companion object { + // Pretty-printed: this file is read by a human debugging why an + // upstream re-synced. + private val json = Json { prettyPrint = true } + + // Often enough that a kill costs little, rare enough to be free. + private const val DEFAULT_FLUSH_SECONDS = 30L + } +} diff --git a/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFileTest.kt b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFileTest.kt new file mode 100644 index 0000000000..4a02507db8 --- /dev/null +++ b/geode/src/test/kotlin/com/vitorpamplona/geode/mirror/SyncCoverageFileTest.kt @@ -0,0 +1,130 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.geode.mirror + +import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer +import java.io.File +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * The catch-up's resume memory across restarts. Without it every boot + * re-syncs each upstream's whole backfill window — a full re-download for an + * upstream without NIP-77. These pin the restart round-trip and the window + * clamping that keys bands on the stable filter rather than the sliding + * boot window. + */ +class SyncCoverageFileTest { + private val relay = RelayUrlNormalizer.normalize("wss://relay.example") + private val profiles = Filter(kinds = listOf(0)) + + private fun tempFile(): File { + val f = File.createTempFile("sync-coverage", ".json") + f.delete() + return f + } + + @Test + fun `bands survive a restart`() { + val f = tempFile() + SyncCoverageFile(f).use { + it.coverage.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) + } + + // A fresh instance, as a restart would build. + SyncCoverageFile(f).use { reopened -> + val band = reopened.coverage.band(relay, profiles)!! + assertEquals(1_700_001_000L, band.minCreatedAt) + assertEquals(1_700_002_000L, band.maxCreatedAt) + assertFalse(band.complete) + } + } + + @Test + fun `a complete band survives with its completeness`() { + val f = tempFile() + SyncCoverageFile(f).use { + it.coverage.record(relay, profiles, null, null, paged = false, reconciledThrough = 1_700_005_000L) + } + SyncCoverageFile(f).use { reopened -> + assertTrue(reopened.coverage.band(relay, profiles)!!.complete) + } + } + + @Test + fun `a corrupt file starts fresh instead of refusing to start`() { + val f = tempFile() + f.writeText("{ not json") + SyncCoverageFile(f).use { + assertNull(it.coverage.band(relay, profiles)) + } + } + + @Test + fun `recording does not write but closing does`() { + val f = tempFile() + val store = SyncCoverageFile(f) + store.coverage.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) + assertFalse(f.isFile, "a record marks dirty; only a flush writes") + store.close() + assertTrue(f.isFile, "close flushes") + } + + @Test + fun `reopening without new records does not rewrite the file`() { + val f = tempFile() + SyncCoverageFile(f).use { + it.coverage.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) + } + val written = f.lastModified() + SyncCoverageFile(f).close() + assertEquals(written, f.lastModified(), "restoring bands must not mark the store dirty") + } + + // ---- the window clamp -------------------------------------------------- + + @Test + fun `an unbanded filter clamps to exactly the boot window`() { + val leg = clampToWindow(profiles, since = 1_000L, until = 2_000L)!! + assertEquals(1_000L, leg.since) + assertEquals(2_000L, leg.until) + } + + @Test + fun `a leg outside the window is dropped rather than inverted`() { + // The band covers past the window's floor: the older leg would ask + // [since..band.min] with since above until — a range nothing can be in. + val olderLeg = profiles.copy(until = 500L) + assertNull(clampToWindow(olderLeg, since = 1_000L, until = 2_000L)) + } + + @Test + fun `a leg inside the window keeps its own tighter bound`() { + val newerLeg = profiles.copy(since = 1_500L) + val clamped = clampToWindow(newerLeg, since = 1_000L, until = 2_000L)!! + assertEquals(1_500L, clamped.since, "the band's ceiling wins over the window floor") + assertEquals(2_000L, clamped.until) + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncBands.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt similarity index 97% rename from quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncBands.kt rename to quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt index 68bff2ca85..c6ddff4dce 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncBands.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverage.kt @@ -50,8 +50,13 @@ import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap * Persistence is the caller's: [export] the map on a schedule and [restore] * it at startup. [onChange] fires whenever a band changes, so a persistence * layer can mark itself dirty without polling. + * + * Not to be confused with the `relay.client.paging` package: its + * `RelayLoadingCursors` are in-memory POSITIONS for demand-driven UI paging + * within one session, while these are persistent INTERVALS — a claim about + * coverage that outlives the process and licenses skipping work. */ -class SyncBands( +class SyncCoverage( // How long a band may narrow work before the whole filter is walked // again. Everything a band claims is a claim about the past; this is how // long to trust it without re-testing. diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/PagingProgress.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt similarity index 62% rename from quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/PagingProgress.kt rename to quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt index 028fa34979..0f515d0ca5 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/PagingProgress.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgress.kt @@ -18,73 +18,81 @@ * 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.relay.client.accessories +package com.vitorpamplona.quartz.nip01Core.relay.client.paging import com.vitorpamplona.quartz.utils.TimeUtils import com.vitorpamplona.quartz.utils.concurrent.ConcurrentMap import kotlin.concurrent.Volatile /** - * How far a paged walk has got, measured on the time axis — the only axis - * whose end is known in advance. + * How far a bulk pagination has got, measured on the time axis — the only + * axis whose end is known in advance. * - * A paged fetch ([fetchAllPages]) has no event denominator: how many events + * A paged fetch (`fetchAllPages`) has no event denominator: how many events * exist is exactly what it is finding out, so every count-based percentage * degenerates to `downloaded/downloaded = 100%`. The time axis has both ends * before the first request — the filter's `until` (or now) down to its - * `since` (or [SyncBands.PLAUSIBLE_FLOOR]) — with each page's new `until` - * reporting the exact position between them. It needs no COUNT support. + * `since` (or the accessories' `SyncCoverage.PLAUSIBLE_FLOOR`) — with each + * page's new `until` cursor reporting the exact position between them. It + * needs no COUNT support. * * The estimate assumes events are spread evenly over time, which they are * not — so it errs pessimistic on the tail, and is a bound, not a promise. * - * One instance can serve many concurrent walks: keys are `"group|walk"`, and - * the group prefix scopes [fraction], [reached] and [etaMs] so two groups - * never report each other's numbers. + * One instance can serve many concurrent paginations: keys are + * `"group|name"`, and the group prefix scopes [fraction], [reached] and + * [etaMs] so two groups never report each other's numbers. + * + * Its siblings in this package track different things: [RelayLoadingCursors] + * is demand-driven `until`+`limit` paging for one scope (a feed pulling + * older pages on demand, no window), and [RelayPagingProgress] is the + * per-relay display state derived from it. This class is for a BULK + * pagination over a known `[since, until]` window, where "how far through + * the window, and when will it finish" is the question. */ -class PagingProgress( +class PagingWindowProgress( private val nowMillis: () -> Long = { TimeUtils.nowMillis() }, ) { - private class Walk( + private class Window( val top: Long, val bottom: Long, val startedMs: Long, @Volatile var current: Long, ) - private val walks = ConcurrentMap() + private val windows = ConcurrentMap() - /** Begin a walk over `[bottom, top]` seconds. An inverted window is not a walk. */ + /** Begin a pagination over `[bottom, top]` seconds. An inverted window is not one. */ fun begin( key: String, top: Long, bottom: Long, ) { - if (top > bottom) walks[key] = Walk(top, bottom, nowMillis(), top) + if (top > bottom) windows[key] = Window(top, bottom, nowMillis(), top) } - /** The walk 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. */ fun mark( key: String, until: Long, ) { - walks[key]?.let { - // Clamped to the walk's own floor: relays serve events stamped 0, - // and one of those would drag the position to the epoch. Below the - // floor means the walk is done, not time travel. + windows[key]?.let { + // Clamped to the window's own floor: relays serve events stamped + // 0, and one of those would drag the position to the epoch. Below + // the floor means the pagination is done, not time travel. val reached = until.coerceAtLeast(it.bottom) if (reached < it.current) it.current = reached } } fun finish(key: String) { - walks.remove(key) + windows.remove(key) } /** - * Fraction of the walk complete, averaged over every walk still going in + * Fraction complete, averaged over every pagination still going in * [group] (or all of them when null) — averaged rather than summed - * because each covers its own span, so "half the walks done and half at + * because each covers its own span, so "half of them done and half at * zero" is 50%. */ fun fraction(group: String? = null): Double? { @@ -96,18 +104,18 @@ class PagingProgress( } / live.size } - private fun live(group: String?): List = + private fun live(group: String?): List = if (group == null) { - walks.snapshot().values.toList() + windows.snapshot().values.toList() } else { - walks + windows .snapshot() .entries .filter { it.key.startsWith("$group|") } .map { it.value } } - /** The oldest second [group] has reached, or null when it is not walking. */ + /** The oldest second [group] has reached, or null when nothing is paging. */ fun reached(group: String? = null): Long? = live(group).minOfOrNull { it.current } /** Milliseconds left at the rate achieved so far, or null before it means anything. */ diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncBandsTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt similarity index 93% rename from quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncBandsTest.kt rename to quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt index d0ffe5290a..78abe61682 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncBandsTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/SyncCoverageTest.kt @@ -35,7 +35,7 @@ import kotlin.test.assertTrue * importantly, the cases where a band must NOT be used — a stale band silently * skips events, which is a worse failure than re-reading them. */ -class SyncBandsTest { +class SyncCoverageTest { private val relay = RelayUrlNormalizer.normalize("wss://relay.example") private val other = RelayUrlNormalizer.normalize("wss://other.example") private val profiles = Filter(kinds = listOf(0)) @@ -46,13 +46,13 @@ class SyncBandsTest { @Test fun `with nothing recorded the whole filter is fetched`() { - val c = SyncBands() + val c = SyncCoverage() assertEquals(listOf(profiles), c.legs(relay, profiles)) } @Test fun `a recorded band is fetched around rather than through`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, observedMin = 1_700_001_000L, observedMax = 1_700_002_000L, paged = true) val legs = c.legs(relay, profiles) @@ -68,7 +68,7 @@ class SyncBandsTest { // A paged relay cuts pages by count, so a boundary can fall inside a run // of events sharing one created_at. Excluding the edge would strand the // rest of that second in no leg at all, while the band called it covered. - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) val legs = c.legs(relay, profiles) @@ -85,7 +85,7 @@ class SyncBandsTest { @Test fun `successive runs widen the band rather than replacing it`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) // A later run reaches further back and picks up newer events. c.record(relay, profiles, 1_700_000_500L, 1_700_002_500L, paged = true) @@ -99,7 +99,7 @@ class SyncBandsTest { fun `a capped relay walks further back on each run`() { // The case that makes this worth having: a relay that only ever answers // with its newest N events. Each run starts below the last one's floor. - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_009_000L, 1_700_010_000L, paged = true) assertEquals(1_700_009_000L, c.legs(relay, profiles)[0].until) @@ -113,7 +113,7 @@ class SyncBandsTest { fun `a negentropy sync that reported no outcome records nothing`() { // Only a sync that says how far it reconciled earns a band; a bare // paged=false call carries no claim to record. - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = false) assertNull(c.band(relay, profiles)) assertEquals(listOf(profiles), c.legs(relay, profiles)) @@ -125,7 +125,7 @@ class SyncBandsTest { fun `a finished reconcile is in sync through the instant it started`() { // Not through the newest event it happened to see: "the relay had nothing // newer" and "we never asked" must not record the same thing. - val c = SyncBands() + val c = SyncCoverage() val startedAt = now() - 60 c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = false, reconciledThrough = startedAt) @@ -138,7 +138,7 @@ class SyncBandsTest { fun `a reconcile that downloaded nothing still records coverage`() { // The empty case is the WHOLE point: nothing came back because we already // have it, and that is exactly when the next run should ask for a sliver. - val c = SyncBands() + val c = SyncCoverage() val startedAt = now() - 60 c.record(relay, profiles, null, null, paged = false, reconciledThrough = startedAt) @@ -149,12 +149,12 @@ class SyncBandsTest { @Test fun `a complete band drops its older leg while a paged one keeps it`() { - val reconciled = SyncBands() + val reconciled = SyncCoverage() reconciled.record(relay, profiles, null, null, paged = false, reconciledThrough = 1_700_002_000L) val only = reconciled.legs(relay, profiles).single() assertEquals(1_700_002_000L, only.since) - val walked = SyncBands() + val walked = SyncCoverage() walked.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) assertEquals(2, walked.legs(relay, profiles).size, "a paged walk says nothing about what it never asked for") } @@ -163,13 +163,13 @@ class SyncBandsTest { @Test fun `a band stops narrowing once it is older than the resync period`() { - val c = SyncBands(fullResyncSeconds = 60) + val c = SyncCoverage(fullResyncSeconds = 60) c.record(relay, profiles, null, null, paged = false, reconciledThrough = now() - 3600) // Recorded 'now' whatever the created_at claim, so age it by rewriting. c.record(relay, profiles, null, null, paged = false, reconciledThrough = now()) assertEquals(1, c.legs(relay, profiles).size, "fresh band still narrows") - val stale = SyncBands(fullResyncSeconds = 0) + val stale = SyncCoverage(fullResyncSeconds = 0) stale.record(relay, profiles, null, null, paged = false, reconciledThrough = now()) assertSame(profiles, stale.legs(relay, profiles).single(), "a band past its period re-walks everything") } @@ -178,7 +178,7 @@ class SyncBandsTest { fun `the re-walk replaces the old claim instead of widening it`() { // Widening would carry the stale band's floor forward forever and the // periodic pass would never actually reset anything. - val c = SyncBands(fullResyncSeconds = 0) + val c = SyncCoverage(fullResyncSeconds = 0) c.record(relay, profiles, 1_700_000_000L, 1_700_001_000L, paged = true) c.record(relay, profiles, 1_700_005_000L, 1_700_006_000L, paged = true) @@ -190,7 +190,7 @@ class SyncBandsTest { @Test fun `covering window collapses to the oldest ceiling once everyone is caught up`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, null, null, paged = false, reconciledThrough = 1_700_009_000L) c.record(other, profiles, null, null, paged = false, reconciledThrough = 1_700_003_000L) @@ -201,7 +201,7 @@ class SyncBandsTest { fun `one relay that has never synced puts the window back to the whole filter`() { // It genuinely needs everything — narrowing the shared snapshot would // reconcile it against ids we never looked up. - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, null, null, paged = false, reconciledThrough = 1_700_009_000L) // The filter itself, unnarrowed — identity, since Filter has no equals. @@ -214,7 +214,7 @@ class SyncBandsTest { // Every url in a stream shares that stream's filter, so a backfill can // take ONE snapshot for all of them instead of walking the identical // range once per relay for byte-identical answers. - val c = SyncBands() + val c = SyncCoverage() val third = RelayUrlNormalizer.normalize("wss://third.example") c.record(relay, profiles, null, null, paged = false, reconciledThrough = 1_700_009_000L) c.record(other, profiles, null, null, paged = false, reconciledThrough = 1_700_003_000L) @@ -226,7 +226,7 @@ class SyncBandsTest { @Test fun `a relay with an older gap also widens the shared window`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, null, null, paged = false, reconciledThrough = 1_700_009_000L) c.record(other, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) @@ -237,7 +237,7 @@ class SyncBandsTest { fun `an empty fetch records nothing`() { // No events says nothing about what the relay holds, only that this // window was empty — recording it would fabricate coverage. - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, null, null, paged = true) assertNull(c.band(relay, profiles)) } @@ -246,11 +246,11 @@ class SyncBandsTest { fun `one misdated event does not cost a relay its whole band`() { // A single future-dated stamp among hundreds of thousands must not fail // a check applied to the aggregate. Screening per event keeps the rest. - val c = SyncBands() + val c = SyncCoverage() val far = now() + 400L * 86_400 val observed = listOf(1_700_001_000L, far, 1_700_002_000L, 0L) - val plausible = observed.filter { SyncBands.isPlausible(it) } + val plausible = observed.filter { SyncCoverage.isPlausible(it) } c.record(relay, profiles, plausible.min(), plausible.max(), paged = true) val band = c.band(relay, profiles)!! @@ -260,7 +260,7 @@ class SyncBandsTest { @Test fun `changing the filter starts over`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) // Widening the kinds means the old band skipped events it never fetched. @@ -273,7 +273,7 @@ class SyncBandsTest { @Test fun `each relay keeps its own band`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) assertEquals(listOf(profiles), c.legs(other, profiles)) } @@ -283,7 +283,7 @@ class SyncBandsTest { @Test fun `a bounded filter never widens past its own since and until`() { val bounded = Filter(kinds = listOf(0), since = 1_700_001_000L, until = 1_700_005_000L) - val c = SyncBands() + val c = SyncCoverage() c.record(relay, bounded, 1_700_002_000L, 1_700_003_000L, paged = true) val legs = c.legs(relay, bounded) @@ -300,7 +300,7 @@ class SyncBandsTest { // the two boundary seconds are always re-read, because that is the only // way to catch a run of same-second events a page boundary cut in half. val bounded = Filter(kinds = listOf(0), since = 1_700_001_000L, until = 1_700_005_000L) - val c = SyncBands() + val c = SyncCoverage() c.record(relay, bounded, 1_700_001_000L, 1_700_005_000L, paged = true) val legs = c.legs(relay, bounded) @@ -313,10 +313,10 @@ class SyncBandsTest { @Test fun `export and restore round-trip the bands`() { - val c = SyncBands() + val c = SyncCoverage() c.record(relay, profiles, 1_700_001_000L, 1_700_002_000L, paged = true) - val reopened = SyncBands() + val reopened = SyncCoverage() reopened.restore(c.export()) val band = reopened.band(relay, profiles)!! @@ -327,7 +327,7 @@ class SyncBandsTest { @Test fun `onChange fires when a band changes so persistence can mark dirty`() { var changes = 0 - val c = SyncBands(onChange = { changes++ }) + val c = SyncCoverage(onChange = { changes++ }) c.record(relay, profiles, null, null, paged = true) assertEquals(0, changes, "an empty fetch records nothing and must not dirty the store") @@ -341,7 +341,7 @@ class SyncBandsTest { // Filter.toJson() runs to tens of thousands of characters for an // author-scoped filter, and a fan-out keys once per relay per cycle. val big = Filter(kinds = listOf(30382), authors = (1..500).map { it.toString(16).padStart(64, '0') }) - val c = SyncBands() + val c = SyncCoverage() c.record(relay, big, 1_700_001_000L, 1_700_002_000L, paged = true) // Same instance, many lookups: still one band, and cheap. diff --git a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/PagingProgressTest.kt b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt similarity index 84% rename from quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/PagingProgressTest.kt rename to quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt index bf8a550c63..2c1c0ed433 100644 --- a/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/PagingProgressTest.kt +++ b/quartz/src/commonTest/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/paging/PagingWindowProgressTest.kt @@ -18,14 +18,14 @@ * 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.relay.client.accessories +package com.vitorpamplona.quartz.nip01Core.relay.client.paging import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertNull import kotlin.test.assertTrue -class PagingProgressTest { +class PagingWindowProgressTest { private fun assertClose( expected: Double, actual: Double?, @@ -35,8 +35,8 @@ class PagingProgressTest { } @Test - fun `progress is the walked share of the time window`() { - val p = PagingProgress() + fun `progress is the paged share of the time window`() { + val p = PagingWindowProgress() p.begin("a", top = 1_000L, bottom = 0L) assertClose(0.0, p.fraction(), "nothing walked yet") @@ -49,11 +49,11 @@ class PagingProgressTest { } @Test - fun `a page that jumps backwards cannot un-advance the walk`() { + fun `a page that jumps backwards cannot un-advance the pagination`() { // Pages arrive from one relay in order, but nothing in the protocol // guarantees it, and a percentage that goes DOWN is worse than one that // is slightly wrong — it reads as the sync having lost ground. - val p = PagingProgress() + val p = PagingWindowProgress() p.begin("a", top = 1_000L, bottom = 0L) p.mark("a", 200L) @@ -63,10 +63,10 @@ class PagingProgressTest { } @Test - fun `walks average rather than sum`() { + fun `windows average rather than sum`() { // Two relays each walking their own window: one done and one untouched // is half way — not 100% as summing would give. - val p = PagingProgress() + val p = PagingWindowProgress() p.begin("a", top = 1_000L, bottom = 0L) p.begin("b", top = 500L, bottom = 0L) @@ -76,8 +76,8 @@ class PagingProgressTest { } @Test - fun `a finished walk leaves the average`() { - val p = PagingProgress() + fun `a finished window leaves the average`() { + val p = PagingWindowProgress() p.begin("a", top = 1_000L, bottom = 0L) p.begin("b", top = 1_000L, bottom = 0L) p.mark("b", 500L) @@ -86,14 +86,14 @@ class PagingProgressTest { assertClose(0.5, p.fraction(), "only b is still walking") p.finish("b") - assertNull(p.fraction(), "nothing walking means no number to report") + 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 // streams would print each other's percentages. - val p = PagingProgress() + val p = PagingWindowProgress() p.begin("streamA|wss://r1", top = 1_000L, bottom = 0L) p.begin("streamB|wss://r2", top = 1_000L, bottom = 0L) p.mark("streamA|wss://r1", 0L) @@ -104,10 +104,10 @@ class PagingProgressTest { } @Test - fun `an inverted or empty window is not a walk`() { + fun `an inverted or empty 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 = PagingProgress() + val p = PagingWindowProgress() p.begin("a", top = 100L, bottom = 900L) @@ -116,7 +116,7 @@ class PagingProgressTest { @Test fun `no ETA before the estimate means anything`() { - val p = PagingProgress() + val p = PagingWindowProgress() p.begin("a", top = 1_000_000L, bottom = 0L) p.mark("a", 999_000L) @@ -129,7 +129,7 @@ class PagingProgressTest { @Test fun `ETA extrapolates from the rate achieved so far`() { var clock = 1_000_000L - val p = PagingProgress(nowMillis = { clock }) + val p = PagingWindowProgress(nowMillis = { clock }) p.begin("a", top = 1_000L, bottom = 0L) p.mark("a", 500L)