mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor(cli): align fetch default limit across paths; harden paging
Audit follow-ups before merge: - amy fetch default limit is now the same on both paths: absent --limit → 100 for plain AND --paginate (previously --paginate silently meant "unbounded"). `--limit 0` is the explicit opt-in to drain everything (unbounded); negative is rejected. The effective limit is carried on the filter so both paths agree. - drainAllPages sizes its SeenIds for CLI-scale fetches (initialSlotsPow2 = 12, ~64 KB) instead of the large-walk default (~16 MB eagerly allocated per fetch); it grows if an unbounded drain needs it. - fetchAllPages clamps the inclusive advance to `min(pageMinTs, boundary)` so a misbehaving relay that answers with an event past the requested `until` can't push the cursor upward — the boundary dedup and termination rely on `until` never increasing. No-op for honest relays (they only return events ≤ until). Verified live: default and --paginate both cap at 100; --limit 50 → 50; --limit 0 --paginate drains the full window (>100); paging tests + SeenIds tests still pass. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015YEbdqCRPkszkGCoi89RMt
This commit is contained in:
@@ -517,8 +517,11 @@ class Context(
|
||||
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()
|
||||
// it verifies so a bad-sig copy can't pre-empt a good one. Start
|
||||
// small (CLI fetches are typically hundreds of events); it grows if
|
||||
// an unbounded drain needs it, rather than eagerly taking the
|
||||
// large-walk default table.
|
||||
val seen = SeenIds(initialSlotsPow2 = 12)
|
||||
for ((relay, event) in eventChannel) {
|
||||
if (seen.contains(event.id)) continue
|
||||
if (verifyAndStore(event)) {
|
||||
|
||||
@@ -52,19 +52,20 @@ 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, and emitted as full event
|
||||
* JSON under an `events` array.
|
||||
* Results are deduplicated by id, sorted newest-first, capped at `--limit` (default
|
||||
* 100 — the SAME on both paths), and emitted as full event JSON under an `events`
|
||||
* array. `--limit 0` removes the cap entirely.
|
||||
*
|
||||
* 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.
|
||||
* 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 via the
|
||||
* multi-relay [Context.drainAllPages] path, fully draining sets larger than one
|
||||
* `REQ`. Both honor the same limit: `--limit N` returns the newest N, absent is 100,
|
||||
* and `--limit 0` is unbounded — combined with `--paginate` that drains the entire
|
||||
* filter, so mind broad filters. Code mode is always single-shot.
|
||||
*/
|
||||
object FetchCommand {
|
||||
/** Output cap for a plain (non-`--paginate`) fetch when `--limit` is omitted. */
|
||||
/** Output/paging cap for a fetch (either path) when `--limit` is omitted. */
|
||||
private const val DEFAULT_LIMIT = 100
|
||||
|
||||
suspend fun run(
|
||||
@@ -72,22 +73,25 @@ object FetchCommand {
|
||||
rest: Array<String>,
|
||||
): Int {
|
||||
val args = Args(rest)
|
||||
// `--limit` is optional. When absent, `--paginate` drains the whole filter
|
||||
// (unbounded) while a plain fetch still trims to DEFAULT_LIMIT.
|
||||
// `--limit`: omitted → DEFAULT_LIMIT on BOTH the plain and --paginate paths;
|
||||
// `0` → unbounded (drain everything — only useful with --paginate); negative
|
||||
// → error. `effectiveLimit == null` means "no cap".
|
||||
val explicitLimit = args.flag("limit")?.toIntOrNull()
|
||||
if (explicitLimit != null && explicitLimit <= 0) return Output.error("bad_args", "--limit must be > 0")
|
||||
if (explicitLimit != null && explicitLimit < 0) return Output.error("bad_args", "--limit must be >= 0 (0 = unbounded)")
|
||||
val effectiveLimit: Int? = if (explicitLimit == 0) null else (explicitLimit ?: DEFAULT_LIMIT)
|
||||
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.
|
||||
// outbox model rather than using a hand-built filter. It fetches a single
|
||||
// entity, so it just uses the (positive) default cap.
|
||||
args.positionalOrNull(0)?.takeIf { looksLikeCode(it) }?.let {
|
||||
return fetchByCode(dataDir, it, explicitLimit ?: DEFAULT_LIMIT, timeoutMs)
|
||||
return fetchByCode(dataDir, it, effectiveLimit ?: 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)
|
||||
// Carry the effective limit on the filter so both paths agree: --paginate
|
||||
// pages each relay up to it (or fully drains the filter when unbounded) and a
|
||||
// plain fetch asks the relay for that many.
|
||||
val filter = RawEventSupport.buildFilter(args).copy(limit = effectiveLimit)
|
||||
val paginate = args.bool("paginate") || args.bool("all")
|
||||
|
||||
Context.open(dataDir).use { ctx ->
|
||||
@@ -107,14 +111,8 @@ object FetchCommand {
|
||||
.map { it.second }
|
||||
.distinctBy { it.id }
|
||||
.sortedByDescending { it.createdAt }
|
||||
// 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)
|
||||
}
|
||||
// effectiveLimit caps the output on both paths; null (--limit 0) is uncapped.
|
||||
val capped = if (effectiveLimit != null) ordered.take(effectiveLimit) else ordered
|
||||
val events =
|
||||
capped
|
||||
.map { Output.mapper.readTree(it.toJson()) }
|
||||
|
||||
+7
-2
@@ -257,12 +257,17 @@ suspend fun INostrClient.fetchAllPages(
|
||||
|
||||
// Advance inclusively to the oldest second seen, carrying its dedup set:
|
||||
// still the same boundary → accumulate; a genuinely older one → replace.
|
||||
if (boundary != null && pageMinTs == boundary) {
|
||||
// Clamp to `boundary` so a misbehaving relay that answers with an event past
|
||||
// the requested `until` can't push the cursor UPWARD — the boundary dedup and
|
||||
// termination both rely on `until` never increasing. Honest relays only
|
||||
// return events at-or-below `until`, so this is a no-op for them.
|
||||
val nextUntil = if (boundary != null) minOf(pageMinTs, boundary) else pageMinTs
|
||||
if (boundary != null && nextUntil == boundary) {
|
||||
seenAtBoundary.addAll(idsAtPageMin)
|
||||
} else {
|
||||
seenAtBoundary = idsAtPageMin
|
||||
}
|
||||
until = pageMinTs
|
||||
until = nextUntil
|
||||
}
|
||||
|
||||
return totalEvents
|
||||
|
||||
Reference in New Issue
Block a user