Align the sync accessories with quartz vocabulary; geode catch-up resumes

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Y4Pi9YYMhdzTFxRiV2jF9R
This commit is contained in:
Claude
2026-08-04 04:24:19 +00:00
parent 564cafedc4
commit 19d8e3069d
9 changed files with 489 additions and 95 deletions
@@ -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<String>) {
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<String>) {
// 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<String>) {
// 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() }
},
@@ -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
* `<database file>.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,
)
/**
@@ -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)
}
@@ -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
}
}
@@ -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)
}
}
@@ -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.
@@ -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<String, Walk>()
private val windows = ConcurrentMap<String, Window>()
/** 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<Walk> =
private fun live(group: String?): List<Window> =
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. */
@@ -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.
@@ -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)