mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 19:53:08 +00:00
feat(cli): dedup drainAllPages via SeenIds; unbounded amy fetch --paginate
Two changes to the paginated fetch path: - Cross-relay dedup before verify. drainAllPages' single consumer now runs a SeenIds filter: the same widely-mirrored event arrives once per relay, and the repeats are dropped BEFORE the expensive Schnorr verify + store instead of after (they were only trimmed by FetchCommand's distinctBy). An id is marked seen only once it verifies, so a forged copy (valid id, bad sig) delivered first can't suppress the genuine one from another relay. Adds SeenIds.contains (peek without recording) for that check-then-add. - `amy fetch --paginate` no longer forces a --limit. With --limit N it still pages up to N per relay; WITHOUT --limit it drains the whole filter unbounded (the filter's null limit flows straight through). Plain (non-paginate) fetch still trims to the default 100. Verified live: unbounded --paginate over a ~20-min nos.lol firehose window returns 406 (all unique, 3s) vs the old 100 cap; --limit 50 caps at 50; default caps at 100; cross-relay fetch stays count==uniq. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015YEbdqCRPkszkGCoi89RMt
This commit is contained in:
@@ -70,6 +70,7 @@ import com.vitorpamplona.quartz.nip61Nutzaps.info.NutzapInfoEvent
|
||||
import com.vitorpamplona.quartz.nip61Nutzaps.nutzap.NutzapEvent
|
||||
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
import com.vitorpamplona.quartz.nip87Ecash.recommendation.MintRecommendationEvent
|
||||
import com.vitorpamplona.quartz.utils.SeenIds
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
|
||||
@@ -487,10 +488,13 @@ class Context(
|
||||
* [fetchAllPagesFromPool] instead of stopping at the first EOSE — so a query
|
||||
* larger than a relay's per-`REQ` cap (strfry's `limit`, ~500) is fully
|
||||
* retrieved instead of silently truncated. Each relay is walked on its own
|
||||
* `until` cursor, up to [maxConcurrentRelays] at once, and every event still
|
||||
* funnels through [verifyAndStore]; the result is tagged by relay exactly like
|
||||
* [drain] (and, like [drain], is NOT deduped across relays — callers dedup by
|
||||
* id).
|
||||
* `until` cursor, up to [maxConcurrentRelays] at once, and every event funnels
|
||||
* through [verifyAndStore]; the result is tagged by the relay that first
|
||||
* delivered it. Unlike [drain], it IS deduped across relays: the same
|
||||
* widely-mirrored event arrives once per relay, and the repeats are dropped by a
|
||||
* [SeenIds] filter BEFORE the expensive verify+store — an id is marked seen only
|
||||
* after it verifies, so a forged copy (valid id, bad signature) delivered first
|
||||
* can't suppress the genuine one from another relay.
|
||||
*
|
||||
* Bound the work with the filters' `limit`: each relay pages until it reaches
|
||||
* the limit, so an unbounded filter pages that relay's entire matching history.
|
||||
@@ -511,8 +515,16 @@ class Context(
|
||||
coroutineScope {
|
||||
val consumer =
|
||||
launch {
|
||||
// One writer → SeenIds' single-writer contract holds. Skip a
|
||||
// cross-relay duplicate before verifying it; mark it seen only once
|
||||
// it verifies so a bad-sig copy can't pre-empt a good one.
|
||||
val seen = SeenIds()
|
||||
for ((relay, event) in eventChannel) {
|
||||
if (verifyAndStore(event)) collected.add(relay to event)
|
||||
if (seen.contains(event.id)) continue
|
||||
if (verifyAndStore(event)) {
|
||||
seen.add(event.id)
|
||||
collected.add(relay to event)
|
||||
}
|
||||
}
|
||||
}
|
||||
try {
|
||||
|
||||
@@ -52,35 +52,42 @@ import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
* NIP-65 write relays, exactly how the app downloads an event/profile from
|
||||
* a shared link. This is nak's `fetch` (nip19-hint resolution).
|
||||
*
|
||||
* Results are deduplicated by id, sorted newest-first, capped at `--limit`
|
||||
* (default 100), and emitted as full event JSON under an `events` array.
|
||||
* Results are deduplicated by id, sorted newest-first, and emitted as full event
|
||||
* JSON under an `events` array.
|
||||
*
|
||||
* By default (filter mode) a single `REQ` is drained to EOSE, so a relay that
|
||||
* caps its response (strfry's per-`REQ` `limit`, ~500) truncates the result.
|
||||
* `--paginate` (alias `--all`) instead walks each relay page-by-page on `until`
|
||||
* cursors up to `--limit`, fully draining sets larger than one `REQ` — the
|
||||
* multi-relay [Context.drainAllPages] path. Code mode is always single-shot.
|
||||
* By default (filter mode) a single `REQ` is drained to EOSE — so a relay that
|
||||
* caps its response (strfry's per-`REQ` `limit`, ~500) truncates the result — and
|
||||
* the output is trimmed to the newest `--limit` (default 100). `--paginate` (alias
|
||||
* `--all`) instead walks each relay page-by-page on `until` cursors via the
|
||||
* multi-relay [Context.drainAllPages] path: with `--limit N` it pages up to N per
|
||||
* relay, and WITHOUT `--limit` it drains the whole filter unbounded (mind broad
|
||||
* filters — that can be a lot). Code mode is always single-shot.
|
||||
*/
|
||||
object FetchCommand {
|
||||
/** Output cap for a plain (non-`--paginate`) fetch when `--limit` is omitted. */
|
||||
private const val DEFAULT_LIMIT = 100
|
||||
|
||||
suspend fun run(
|
||||
dataDir: DataDir,
|
||||
rest: Array<String>,
|
||||
): Int {
|
||||
val args = Args(rest)
|
||||
val limit = args.flag("limit")?.toIntOrNull() ?: 100
|
||||
if (limit <= 0) return Output.error("bad_args", "--limit must be > 0")
|
||||
// `--limit` is optional. When absent, `--paginate` drains the whole filter
|
||||
// (unbounded) while a plain fetch still trims to DEFAULT_LIMIT.
|
||||
val explicitLimit = args.flag("limit")?.toIntOrNull()
|
||||
if (explicitLimit != null && explicitLimit <= 0) return Output.error("bad_args", "--limit must be > 0")
|
||||
val timeoutMs = (args.flag("timeout")?.toLongOrNull() ?: 8L) * 1000
|
||||
|
||||
// Code mode: a nip19/nip05 positional resolves its own relays via the
|
||||
// outbox model rather than using a hand-built filter.
|
||||
args.positionalOrNull(0)?.takeIf { looksLikeCode(it) }?.let {
|
||||
return fetchByCode(dataDir, it, limit, timeoutMs)
|
||||
return fetchByCode(dataDir, it, explicitLimit ?: DEFAULT_LIMIT, timeoutMs)
|
||||
}
|
||||
|
||||
// buildFilter already carries `--limit` (or null) as the filter's limit, so
|
||||
// the paginate path uses it verbatim: bounded per relay with --limit, or a
|
||||
// full drain of the filter without one.
|
||||
val filter = RawEventSupport.buildFilter(args)
|
||||
// --paginate/--all walks each relay past its per-REQ cap (strfry's ~500)
|
||||
// by following `until` cursors, bounded by `limit`; default stops at the
|
||||
// first EOSE like nak's `req`.
|
||||
val paginate = args.bool("paginate") || args.bool("all")
|
||||
|
||||
Context.open(dataDir).use { ctx ->
|
||||
@@ -90,17 +97,26 @@ object FetchCommand {
|
||||
|
||||
val received =
|
||||
if (paginate) {
|
||||
ctx.drainAllPages(relays.associateWith { listOf(filter.copy(limit = limit)) }, timeoutMs)
|
||||
ctx.drainAllPages(relays.associateWith { listOf(filter) }, timeoutMs)
|
||||
} else {
|
||||
ctx.drain(relays.associateWith { listOf(filter) }, timeoutMs)
|
||||
}
|
||||
val events =
|
||||
val ordered =
|
||||
received
|
||||
.asSequence()
|
||||
.map { it.second }
|
||||
.distinctBy { it.id }
|
||||
.sortedByDescending { it.createdAt }
|
||||
.take(limit)
|
||||
// An explicit --limit always caps the output. Without it, a single-REQ
|
||||
// fetch trims to DEFAULT_LIMIT; a --paginate drain returns everything.
|
||||
val capped =
|
||||
when {
|
||||
explicitLimit != null -> ordered.take(explicitLimit)
|
||||
paginate -> ordered
|
||||
else -> ordered.take(DEFAULT_LIMIT)
|
||||
}
|
||||
val events =
|
||||
capped
|
||||
.map { Output.mapper.readTree(it.toJson()) }
|
||||
.toList()
|
||||
|
||||
|
||||
@@ -84,6 +84,29 @@ class SeenIds(
|
||||
return addKey(Hex.readLong(idHex, 0), Hex.readLong(idHex, 16))
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether [idHex] has already been [add]ed this pass, WITHOUT recording it. A
|
||||
* too-short/malformed id is never "seen" (returns false) — the mirror of [add]
|
||||
* letting it through. Use this to check-then-conditionally-add, e.g. to mark an
|
||||
* id seen only after it verifies (so a forged copy sharing a valid id can't
|
||||
* pre-empt the genuine one).
|
||||
*/
|
||||
fun contains(idHex: String): Boolean {
|
||||
if (idHex.length < 32) return false
|
||||
val hi = Hex.readLong(idHex, 0)
|
||||
val lo = Hex.readLong(idHex, 16)
|
||||
if (hi == 0L && lo == 0L) return zeroSeen
|
||||
var i = (mix(hi, lo).toInt() and mask)
|
||||
while (true) {
|
||||
val s = i * 2
|
||||
val h = table[s]
|
||||
val l = table[s + 1]
|
||||
if (h == 0L && l == 0L) return false
|
||||
if (h == hi && l == lo) return true
|
||||
i = (i + 1) and mask
|
||||
}
|
||||
}
|
||||
|
||||
private fun addKey(
|
||||
hi: Long,
|
||||
lo: Long,
|
||||
|
||||
@@ -85,4 +85,16 @@ class SeenIdsTest {
|
||||
assertTrue(seen.add("not-hex"), "malformed -> flows to verify")
|
||||
assertTrue(seen.add("short"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `contains peeks without recording`() {
|
||||
val seen = SeenIds()
|
||||
assertFalse(seen.contains(id(9)), "never added -> not contained")
|
||||
assertFalse(seen.contains(id(9)), "peeking does not add it")
|
||||
assertEquals(0, seen.size(), "contains must not grow the set")
|
||||
assertTrue(seen.add(id(9)))
|
||||
assertTrue(seen.contains(id(9)), "now contained")
|
||||
assertFalse(seen.contains("short"), "malformed is never contained")
|
||||
assertEquals(1, seen.size())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user