mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
Fix audit findings across the sync accessories and their consumers
quartz: - SQLiteEventStore: classify per-row savepoint errors — policy refusals (blocked:/constraint/not allowed) stay Rejected, everything else is now Failed, so disk-full no longer masquerades as 2M duplicate rejections - IEventStore.batchInsert default: rethrow CancellationException and map unknown throws to Failed (re-offering a duplicate is idempotent; dropping a good event on a transient store error is not) - IngestQueue: rethrow CancellationException instead of stamping a cancelled batch Failed and continuing - HostStrikes: make the eviction verdict exactly-once under concurrency (deadHosts.add is the atomic gate) and re-check produced before publishing - SyncCoverage: bound the identity fingerprint cache (a caller minting fresh Filter instances per cycle could grow it forever); legs() gains a floor parameter so a complete band re-opens its older span when the caller's window deepens; coveringWindow no longer treats a fully covered relay as needing the whole filter - PagingWindowProgress: accept single-second windows (a band's re-read edge leg is exactly that shape) geode: - MirrorWorker: cap reconciledThrough at the leg's own ceiling — the older leg of a resumed catch-up no longer stamps the band complete through 'now' before the newer leg has run (silent event loss for up to fullResyncSeconds if that leg failed) - MirrorWorker: run negentropy and the paged fallback by hand instead of negentropySyncOrFetch: drops the O(delivered-ids) dedup set from the mirror path, and a fallback resets the observed span so a band never claims interior ranges only a half-finished reconcile scattered over - MirrorWorker: clamp a paged band's ceiling to the snapshot instant so one future-dated event cannot suppress the next boot's newer leg - MirrorWorker.close(): join the workers (bounded) so the final coverage flush carries the last records - Main: gate the coverage file on the store actually being persistent — database.file with in_memory=true (the default) persisted bands over a volatile store, and the next boot skipped the backfill over an empty database; honor --db overrides - SyncCoverageFile: request ATOMIC_MOVE explicitly; fix the restore/dirty comment - Import summary now prints the failed count; document mirror_sync_state_file in config.example.toml
This commit is contained in:
@@ -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 "<database file>.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.
|
||||
|
||||
@@ -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<String>) {
|
||||
}
|
||||
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<String>) {
|
||||
"[[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()) {
|
||||
|
||||
@@ -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<Event>(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.
|
||||
*/
|
||||
|
||||
@@ -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}" }
|
||||
}
|
||||
|
||||
+39
-6
@@ -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<Filter, String>()
|
||||
|
||||
/**
|
||||
@@ -106,6 +108,7 @@ class SyncCoverage(
|
||||
fun legs(
|
||||
url: NormalizedRelayUrl,
|
||||
filter: Filter,
|
||||
floor: Long? = null,
|
||||
): List<Filter> {
|
||||
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<Filter>()
|
||||
|
||||
// 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
|
||||
|
||||
+12
-3
@@ -62,16 +62,25 @@ class PagingWindowProgress(
|
||||
|
||||
private val windows = ConcurrentMap<String, Window>()
|
||||
|
||||
/** 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,
|
||||
|
||||
+6
@@ -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"
|
||||
|
||||
+8
-3
@@ -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<Event>): List<InsertOutcome> =
|
||||
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")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -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<IEventStore.InsertOutcome>(events.size)
|
||||
|
||||
+23
-1
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+7
-1
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
+30
@@ -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()
|
||||
|
||||
+21
-6
@@ -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()
|
||||
|
||||
+3
-2
@@ -120,6 +120,7 @@ class ConcurrentIngestLossTest {
|
||||
|
||||
val accepted = ConcurrentHashMap.newKeySet<String>()
|
||||
val rejected = AtomicInteger()
|
||||
val failed = AtomicInteger()
|
||||
val window = Semaphore(200)
|
||||
val done = CompletableDeferred<Unit>()
|
||||
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()
|
||||
|
||||
Reference in New Issue
Block a user