mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-08-09 16:14:40 +00:00
negentropy: audit fixes over the windowing change
A read-back over the two commits before this, rather than a failure — which is the only way these would have turned up, since every one of them lives on a path that runs when something has already gone wrong. **accept() is no longer single-threaded, and its comment said it was.** "Both phases run sequentially, so no concurrent access" was true right up until a paged window started running on a reconciler coroutine while the sync's own delivery consumer was still calling accept(). An unguarded HashSet between two coroutines can corrupt, and the delivered counter can lose updates. Now behind a Mutex — with onEvent kept INSIDE it, because callers are promised it never runs concurrently with itself and some of them keep unsynchronised state in that callback. pagedWindows becomes an AtomicInt for the same reason. **The kotlinx cap parse could take down the whole frame.** `.jsonPrimitive` throws on an object or array, so a relay putting something structured in the fourth element would have failed the NEG-ERR and lost the reason with it — where before that element existed, anything extra was simply ignored. `as?` restores that. Both mappers are now tested against a structured fourth element as well as a string one. **Int overflow in the split fan-out.** `mine + ceiling - 1` wraps when a window holds close to Int.MAX events, which is reachable on exactly the corpora this targets; done in Long now. **The count-driven split cuts N ways, not two.** The work queue is FIFO, so halving means every internal node's count() runs before the first NEG-OPEN goes out: on a corpus ~30,000 windows wide that is ~30,000 store counts of dead time with nothing downloading. Cutting into ceil(count/budget) pieces (capped at 32) reaches the same corpus in about three levels instead of fifteen, and pieces that guess wrong are re-split by the same rule. **The budget moves by CAS.** With reconcileConcurrency > 1 two reconcilers adjust it at once, and a lost SHRINK is the one that costs something real: the next window is then asked at a size the relay has already refused.
This commit is contained in:
+4
-2
@@ -191,8 +191,10 @@ object MessageKSerializer : KSerializer<Message> {
|
||||
reason = if (array.size > 2) array[2].jsonPrimitive.content else "",
|
||||
// Optional, and only a number: a relay that puts something
|
||||
// else there is telling us nothing rather than breaking the
|
||||
// frame.
|
||||
cap = if (array.size > 3) array[3].jsonPrimitive.longOrNull else null,
|
||||
// frame. `as?` rather than `.jsonPrimitive`, which THROWS on
|
||||
// an object or array — that would fail the whole message and
|
||||
// lose the reason, where before this element was ignored.
|
||||
cap = if (array.size > 3) (array[3] as? JsonPrimitive)?.longOrNull else null,
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
+46
-1
@@ -26,6 +26,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IdAndTime
|
||||
import com.vitorpamplona.quartz.nip01Core.store.verifyAndInsert
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.awaitAll
|
||||
@@ -92,6 +93,15 @@ class NegentropyStoreSync(
|
||||
* @param concurrency relays synced at once by [sync] (a relay's own filters stay sequential).
|
||||
* @param idleTimeoutMs idle watchdog for reconciles / fetches / pages.
|
||||
* @param publishTimeoutSecs OK-confirmation wait per uploaded event.
|
||||
* @param targetWindow events per reconcile window, or `0` to snapshot the
|
||||
* whole filter up front (the default, and what this class always did).
|
||||
*
|
||||
* Above zero, the store is read one `created_at` window at a time through
|
||||
* a [NegentropyLocalIndex] instead: the id snapshot stops being O(matched
|
||||
* set) — it is the largest thing this class holds — at the price of an
|
||||
* indexed count + range read per window. Worth turning on exactly when the
|
||||
* filter matches more than fits comfortably in memory; pointless below
|
||||
* that, where one snapshot shared by the whole group is cheaper.
|
||||
*/
|
||||
class Config(
|
||||
val down: Boolean = true,
|
||||
@@ -105,6 +115,7 @@ class NegentropyStoreSync(
|
||||
val concurrency: Int = 4,
|
||||
val idleTimeoutMs: Long = 30_000L,
|
||||
val publishTimeoutSecs: Long = 15,
|
||||
val targetWindow: Int = 0,
|
||||
)
|
||||
|
||||
/** Outcome of one `(relay, filter)` group. `error` is null on success. */
|
||||
@@ -166,7 +177,14 @@ class NegentropyStoreSync(
|
||||
// events (~40 B/entry vs ~1 KB), which matters when a relay hosts a large
|
||||
// matched set. The events the reconcile decides to UP-publish (the small
|
||||
// residual haves) are fetched by id on demand in the uploader below.
|
||||
val localEntries = store.snapshotIdsForNegentropy(listOf(filter))
|
||||
//
|
||||
// With a targetWindow, even those 40 B/entry are read per window rather
|
||||
// than for the whole filter — on a large store that snapshot is the
|
||||
// biggest thing this class allocates, and it is allocated before the
|
||||
// first frame goes out.
|
||||
val windowed = config.targetWindow > 0
|
||||
val localIndex = if (windowed) StoreWindowIndex(store) else null
|
||||
val localEntries = if (windowed) emptyList() else store.snapshotIdsForNegentropy(listOf(filter))
|
||||
|
||||
val downloaded = AtomicInt(0)
|
||||
val uploaded = AtomicInt(0)
|
||||
@@ -210,6 +228,8 @@ class NegentropyStoreSync(
|
||||
relay = relay,
|
||||
filter = filter,
|
||||
localEntries = localEntries,
|
||||
localIndex = localIndex,
|
||||
targetWindow = config.targetWindow,
|
||||
batchSize = config.idChunk,
|
||||
idleTimeoutMs = config.idleTimeoutMs,
|
||||
reconcileConcurrency = config.reconcileConcurrency,
|
||||
@@ -311,3 +331,28 @@ class NegentropyStoreSync(
|
||||
return stored.load()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* [NegentropyLocalIndex] over an [IEventStore]: the window engine's per-window
|
||||
* reads answered straight from the store's `created_at` index.
|
||||
*
|
||||
* The windows handed here are the caller's own filter with `since`/`until`
|
||||
* narrowed, so they can go to the store as-is. A count the store cannot answer
|
||||
* comes back null rather than throwing — the engine then simply stops
|
||||
* pre-splitting that window and lets the relay's refusal decide, which is the
|
||||
* behaviour without an index at all.
|
||||
*/
|
||||
private class StoreWindowIndex(
|
||||
private val store: IEventStore,
|
||||
) : NegentropyLocalIndex {
|
||||
override suspend fun count(window: Filter): Int? =
|
||||
try {
|
||||
store.count(window)
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (_: Exception) {
|
||||
null
|
||||
}
|
||||
|
||||
override suspend fun entriesFor(window: Filter): List<IdAndTime> = store.snapshotIdsForNegentropy(listOf(window))
|
||||
}
|
||||
|
||||
+105
-38
@@ -329,6 +329,7 @@ class NegentropyOrFetchResult(
|
||||
* Use [negentropySync] directly if you want to decide the fallback yourself (try
|
||||
* another relay, narrow the filter, abort, …) instead of always paging.
|
||||
*/
|
||||
@OptIn(ExperimentalAtomicApi::class)
|
||||
suspend fun INostrClient.negentropySyncOrFetch(
|
||||
relay: NormalizedRelayUrl,
|
||||
filter: Filter,
|
||||
@@ -346,18 +347,29 @@ suspend fun INostrClient.negentropySyncOrFetch(
|
||||
): NegentropyOrFetchResult {
|
||||
val seen = HashSet<HexKey>()
|
||||
var delivered = 0
|
||||
var pagedWindows = 0
|
||||
val pagedWindows = AtomicInt(0)
|
||||
|
||||
// Shared dedup + cap across both phases. Returns true if the event was new and
|
||||
// delivered. Both phases run sequentially, so no concurrent access.
|
||||
suspend fun accept(event: Event): Boolean {
|
||||
if ((maxEvents <= 0 || delivered < maxEvents) && seen.add(event.id)) {
|
||||
delivered++
|
||||
onEvent(event)
|
||||
return true
|
||||
// Shared dedup + cap across every path that delivers.
|
||||
//
|
||||
// The lock is not optional. The two phases used to run strictly one after
|
||||
// the other, but a paged window now runs DURING the negentropy phase, on a
|
||||
// reconciler coroutine, while the sync's own delivery consumer is calling
|
||||
// this too — an unguarded HashSet between them can corrupt, and the count
|
||||
// can lose updates. onEvent stays INSIDE the lock deliberately: callers are
|
||||
// promised it never runs concurrently with itself, and some of them keep
|
||||
// unsynchronised state in it.
|
||||
val gate = Mutex()
|
||||
|
||||
suspend fun accept(event: Event): Boolean =
|
||||
gate.withLock {
|
||||
if ((maxEvents <= 0 || delivered < maxEvents) && seen.add(event.id)) {
|
||||
delivered++
|
||||
onEvent(event)
|
||||
true
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
return try {
|
||||
val result =
|
||||
@@ -379,7 +391,7 @@ suspend fun INostrClient.negentropySyncOrFetch(
|
||||
// already reconciled cleanly walked again over REQ, which on a
|
||||
// large corpus is the entire cost negentropy was there to save.
|
||||
onUnreconcilableWindow = { window ->
|
||||
pagedWindows++
|
||||
pagedWindows.incrementAndFetch()
|
||||
val pageTimeoutMs = if (idleTimeoutMs > 0) idleTimeoutMs else DEFAULT_DOWNLOAD_IDLE_MS
|
||||
fetchAllPages(relay, listOf(window), pageTimeoutMs) { event ->
|
||||
if (accept(event)) onProgress?.invoke(delivered, delivered)
|
||||
@@ -392,10 +404,10 @@ suspend fun INostrClient.negentropySyncOrFetch(
|
||||
// Any paged window makes this not a clean reconcile — see the
|
||||
// property doc: under-reporting it would let a caller record
|
||||
// coverage it never compared.
|
||||
pagedFallback = pagedWindows > 0,
|
||||
pagedFallback = pagedWindows.load() > 0,
|
||||
negentropy = result,
|
||||
fallbackCause = null,
|
||||
pagedWindows = pagedWindows,
|
||||
pagedWindows = pagedWindows.load(),
|
||||
)
|
||||
} catch (e: NegentropySyncException) {
|
||||
// Negentropy couldn't enumerate the set — page the whole filter instead,
|
||||
@@ -411,7 +423,7 @@ suspend fun INostrClient.negentropySyncOrFetch(
|
||||
pagedFallback = true,
|
||||
negentropy = null,
|
||||
fallbackCause = e,
|
||||
pagedWindows = pagedWindows,
|
||||
pagedWindows = pagedWindows.load(),
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -567,7 +579,9 @@ internal suspend fun reconcileWindows(
|
||||
onPeerCap: ((Long) -> Unit)? = null,
|
||||
// Given a minimal window the relay will not reconcile at any size, instead
|
||||
// of throwing. The caller drains it however it can (paging it over REQ) and
|
||||
// the sweep carries on with the rest of the filter.
|
||||
// the sweep carries on with the rest of the filter. It runs ON the reconciler
|
||||
// that hit the window, so a slow drain holds that reconciler — with
|
||||
// reconcileConcurrency = 1 the rest of the sweep waits for it.
|
||||
onUnreconcilableWindow: (suspend (Filter) -> Unit)? = null,
|
||||
sendNeedBatch: suspend (List<HexKey>) -> Unit,
|
||||
sendHaveBatch: (suspend (List<HexKey>) -> Unit)?,
|
||||
@@ -597,25 +611,57 @@ internal suspend fun reconcileWindows(
|
||||
// targetWindow at 0 nothing below this line does anything.
|
||||
val budget = AtomicInt(targetWindow)
|
||||
|
||||
// Splits a window in two and queues both halves. Returns false when the
|
||||
// window is already minimal — `created_at` is in seconds, so that is the
|
||||
// floor, not a tuning choice.
|
||||
fun splitInto(
|
||||
// Every budget move goes through here. With reconcileConcurrency > 1 two
|
||||
// reconcilers adjust it at once, and read-then-store can drop one of them —
|
||||
// a lost SHRINK being the one that costs something real, since the next
|
||||
// window is then asked at a size the relay has already refused.
|
||||
fun budgetTo(next: (Int) -> Int) {
|
||||
while (true) {
|
||||
val now = budget.load()
|
||||
val want = next(now)
|
||||
if (want == now || budget.compareAndSet(now, want)) return
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Cuts `[lo, hi]` into [pieces] equal spans of time and queues them all. A
|
||||
* window already at the floor is left alone — `created_at` is in seconds, so
|
||||
* that is where splitting ends, not a tuning choice. Both callers check that
|
||||
* themselves; the guard here is so a third one cannot silently lose a window.
|
||||
*
|
||||
* [pieces] > 2 exists for the count-driven split, where we know HOW FAR over
|
||||
* the budget a window is and can land near the right size in one step.
|
||||
* Halving instead costs a store count per level of a tree that can be ~15
|
||||
* deep on a large corpus, and — since the queue is FIFO — every one of those
|
||||
* counts happens before the first window is reconciled at all.
|
||||
*/
|
||||
suspend fun splitInto(
|
||||
pendingWindow: Filter,
|
||||
lo: Long,
|
||||
hi: Long,
|
||||
): Boolean {
|
||||
if (hi - lo <= MIN_WINDOW_SECONDS) return false
|
||||
val mid = lo + (hi - lo) / 2
|
||||
remaining.incrementAndFetch()
|
||||
// The lower child gets the finite midpoint; the upper child KEEPS this
|
||||
// window's original `until` (which may be null = unbounded). Replacing
|
||||
// null with `now()` here would drop every event dated after now()
|
||||
// (clock skew) once any split happens, while the un-split path would
|
||||
// have included them.
|
||||
pending.trySend(pendingWindow.copy(since = lo, until = mid))
|
||||
pending.trySend(pendingWindow.copy(since = mid + 1, until = pendingWindow.until))
|
||||
return true
|
||||
pieces: Int = 2,
|
||||
) {
|
||||
if (hi - lo <= MIN_WINDOW_SECONDS) return
|
||||
val span = hi - lo + 1
|
||||
// Never more pieces than there are seconds to give them.
|
||||
val n = pieces.toLong().coerceIn(2L, minOf(span, MAX_SPLIT_FANOUT.toLong())).toInt()
|
||||
val step = span / n
|
||||
remaining.addAndFetch(n - 1)
|
||||
var start = lo
|
||||
repeat(n) { i ->
|
||||
val last = i == n - 1
|
||||
// The top piece KEEPS this window's original `until` (which may be
|
||||
// null = unbounded). Replacing null with `now()` here would drop
|
||||
// every event dated after now() (clock skew) once any split happens,
|
||||
// while the un-split path would have included them.
|
||||
if (last) {
|
||||
pending.send(pendingWindow.copy(since = start, until = pendingWindow.until))
|
||||
} else {
|
||||
val end = start + step - 1
|
||||
pending.send(pendingWindow.copy(since = start, until = end))
|
||||
start = end + 1
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
val reconcilers =
|
||||
@@ -636,7 +682,14 @@ internal suspend fun reconcileWindows(
|
||||
if (ceiling > 0 && hi - lo > MIN_WINDOW_SECONDS) {
|
||||
val mine = local.count(window)
|
||||
if (mine != null && mine > ceiling) {
|
||||
splitInto(window, lo, hi)
|
||||
// How many windows this one is worth, not just "two":
|
||||
// the count says how far over budget we are, and
|
||||
// uneven density is corrected by the same check on
|
||||
// each piece.
|
||||
// Long arithmetic: `mine` can be near Int.MAX on a
|
||||
// corpus this size, and the +ceiling would wrap.
|
||||
val over = (mine.toLong() + ceiling - 1) / ceiling
|
||||
splitInto(window, lo, hi, pieces = over.coerceAtMost(MAX_SPLIT_FANOUT.toLong()).toInt())
|
||||
continue
|
||||
}
|
||||
}
|
||||
@@ -662,9 +715,12 @@ internal suspend fun reconcileWindows(
|
||||
// asked for, so a sync that met one dense stretch
|
||||
// does not stay small for the rest of the timeline.
|
||||
if (targetWindow > 0) {
|
||||
val now = budget.load()
|
||||
if (now < targetWindow) {
|
||||
budget.store(minOf(targetWindow, (now * BUDGET_GROWTH).toInt().coerceAtLeast(now + 1)))
|
||||
budgetTo { now ->
|
||||
if (now >= targetWindow) {
|
||||
now
|
||||
} else {
|
||||
minOf(targetWindow, (now * BUDGET_GROWTH).toInt().coerceAtLeast(now + 1))
|
||||
}
|
||||
}
|
||||
}
|
||||
if (remaining.decrementAndFetch() == 0) pending.close()
|
||||
@@ -677,13 +733,16 @@ internal suspend fun reconcileWindows(
|
||||
outcome.cap?.let { cap ->
|
||||
onPeerCap?.invoke(cap)
|
||||
if (targetWindow > 0) {
|
||||
val fitted = (cap * CAP_MARGIN).toInt().coerceAtLeast(1)
|
||||
if (fitted < budget.load()) budget.store(fitted)
|
||||
val fitted =
|
||||
(cap * CAP_MARGIN)
|
||||
.coerceIn(1.0, Int.MAX_VALUE.toDouble())
|
||||
.toInt()
|
||||
budgetTo { now -> minOf(now, fitted) }
|
||||
}
|
||||
}
|
||||
if (outcome.cap == null && targetWindow > 0) {
|
||||
// No number to go on: halve and find out.
|
||||
budget.store((budget.load() / 2).coerceAtLeast(1))
|
||||
budgetTo { now -> (now / 2).coerceAtLeast(1) }
|
||||
}
|
||||
if (hi - lo <= MIN_WINDOW_SECONDS) {
|
||||
// A minimal window that still overflows:
|
||||
@@ -1295,6 +1354,14 @@ private const val MIN_WINDOW_SECONDS = 1L
|
||||
*/
|
||||
private const val MAX_WINDOWS = 100_000
|
||||
|
||||
/**
|
||||
* Most pieces one count-driven split may cut a window into. Bounds both the
|
||||
* queue and the depth: with 32, a corpus 30,000 windows wide is reached in
|
||||
* three levels instead of fifteen, and the pieces that guessed wrong are
|
||||
* re-split by the same rule.
|
||||
*/
|
||||
private const val MAX_SPLIT_FANOUT = 32
|
||||
|
||||
/**
|
||||
* How much of a relay's stated `max_sync_events` a window actually aims for.
|
||||
* The margin absorbs what the relay gains between stating that number and
|
||||
|
||||
+18
@@ -175,6 +175,24 @@ class Nip77SerializationTest {
|
||||
assertEquals(null, kotlin.cap)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun deserializeNegErrMessageWithStructuredFourthElement_bothMappers() {
|
||||
// A fourth element that is an object or array must degrade to no cap,
|
||||
// NOT fail the frame — the reason is the part that matters, and before
|
||||
// this element existed any extra was simply ignored.
|
||||
val json = """["NEG-ERR","neg-sub1","blocked: too many query results",{"max":10}]"""
|
||||
|
||||
val jackson = JacksonMapper.fromJsonToMessage(json)
|
||||
assertTrue(jackson is NegErrMessage)
|
||||
assertEquals("blocked: too many query results", jackson.reason)
|
||||
assertEquals(null, jackson.cap)
|
||||
|
||||
val kotlin = KotlinSerializationMapper.fromJsonToMessage(json)
|
||||
assertTrue(kotlin is NegErrMessage)
|
||||
assertEquals("blocked: too many query results", kotlin.reason)
|
||||
assertEquals(null, kotlin.cap)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun negErrMessageWithCap_crossDeserialization() {
|
||||
val msg = NegErrMessage("neg-sub1", "blocked: too many records", 500_000L)
|
||||
|
||||
Reference in New Issue
Block a user