mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat(graperank): add reverse follower crawl (amy graperank followers)
The outbox model can't find an observer's followers — you don't know a follower exists until you've seen their kind:3, so you can't route to their outbox first. FollowerCrawler casts a wide net instead: it asks as many relays as possible for kind:3 lists that #p-tag the observer, paging each relay past its per-REQ cap via fetchAllPagesFromPool, verifies with ParallelEventVerifier, keeps only lists that genuinely tag the observer, dedups by id, and group-commits to the store. Each follower's list is a full contact list, so persisting it also enriches the graph a later `graperank score` builds — every follower becomes a FOLLOW edge into the observer. CLI: `amy graperank followers [OBSERVER]` assembles "all possible relays" from the reachability-cache live set + every kind:10002/30166 relay in the store + the index/aggregator relays, skipping proven-dead relays. Runs anonymously (no signing) when given an explicit observer. Tunable via --relay/--page-limit/--timeout/--relay-concurrency/--insert-batch. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Xc3Wm4qVCrAGvSAotTUVt4
This commit is contained in:
@@ -391,6 +391,7 @@ HTTP endpoint. Reuses quartz's `Nip86Client` and the shared `Nip86Retriever`
|
||||
| `amy follow USER` / `amy unfollow USER` | Add/remove USER from your kind:3 contact list (fetches the freshest list first). |
|
||||
| `amy graperank [OBSERVER] [--offline] [--min-rank N]` | Crawl + score: compute GrapeRank web-of-trust scores (0..1) over the follow/mute/report graph, then persist the result. Exhaustively crawls each user's kind:10002 outbox for their latest kind:3/10000/1984 until every discovered user is checked (no user cap), dropping reports the author retracted via NIP-09. **Every score run persists its result locally**: the cards (rank cutoff `--min-rank`, default 2) are reconciled into the shared store as NIP-85 kind:30382 cards signed by a per-observer **service key** — each carries `rank`, `followers` (trusted-follower count, cutoff `--followers-threshold`, default 0.02) and `hops` (follow distance from the observer); changed cards re-signed, unchanged skipped (no event-id churn), dropped targets retracted (kind:5). `--offline` skips the crawl. |
|
||||
| `amy graperank crawl [OBSERVER] [--max-hops N] [--no-preconnect]` | Pipeline stage 1 — network only: crawl the follow/mute/report graph (kind 3/10000/1984/10002) into the local store, no scoring. Idempotent and cumulative: run it a few times to load everything, then `score`. |
|
||||
| `amy graperank followers [OBSERVER] [--relay URL[,URL…]]` | The **reverse** crawl: find every user who *follows* the observer by asking as many relays as possible for kind:3 lists that `#p`-tag them (paged past each relay's cap). The outbox model can't find followers — you don't know one exists until you've seen their list — so this casts a wide net over the whole relay universe (reachability-cache live set + every kind:10002/30166 relay in the store + index/aggregator relays), skipping proven-dead relays. Persists each follower's list, enriching the graph a later `score` builds. Idempotent and cumulative. |
|
||||
| `amy graperank score [OBSERVER]` | Pipeline stage 2 — local only: score from the store and persist the cards (identical to bare `--offline`; same scoring flags). No network, so re-run with different `--rigor`/`--attenuation`/`--min-rank` without re-crawling. |
|
||||
| `amy graperank publish [OBSERVER] [--relay URL[,URL…]]` | Pipeline stage 3 — transport only: make the operator relay(s) converge to the locally persisted card set — one NIP-77 up-only reconcile per relay over the service key's kind:30382 + kind:5 (nothing is re-scored or re-signed; a relay that can't reconcile gets the full set published instead). Also refreshes the observer's kind:10040 pointer when we hold their key. |
|
||||
| `amy graperank rank USER [--provider PUBKEY] [--refresh]` | The consumer side: read the kind:30382 cards about USER — one rank per provider, newest card each. Local store first; `--refresh` (or a miss) drains the operator relays, the relays your kind:10040 declares, and the bootstrap set. |
|
||||
|
||||
@@ -26,6 +26,7 @@ import com.vitorpamplona.amethyst.cli.DataDir
|
||||
import com.vitorpamplona.amethyst.cli.Output
|
||||
import com.vitorpamplona.amethyst.commons.defaults.Constants
|
||||
import com.vitorpamplona.amethyst.commons.defaults.DefaultIndexerRelayList
|
||||
import com.vitorpamplona.quartz.experimental.graperank.FollowerCrawler
|
||||
import com.vitorpamplona.quartz.experimental.graperank.GrapeRank
|
||||
import com.vitorpamplona.quartz.experimental.graperank.GrapeRankCrawler
|
||||
import com.vitorpamplona.quartz.experimental.graperank.GrapeRankParams
|
||||
@@ -213,6 +214,7 @@ object GrapeRankCommand {
|
||||
"providers" -> providers(dataDir, tail.drop(1).toTypedArray())
|
||||
"operator" -> operator(dataDir, tail.drop(1).toTypedArray())
|
||||
"crawl" -> crawl(dataDir, tail.drop(1).toTypedArray())
|
||||
"followers" -> followers(dataDir, tail.drop(1).toTypedArray())
|
||||
"status" -> status(dataDir)
|
||||
// The relay census outgrew graperank (it feeds the shared NIP-66
|
||||
// reachability cache every command reads) and moved to `amy relay
|
||||
@@ -586,6 +588,123 @@ object GrapeRankCommand {
|
||||
return 0
|
||||
}
|
||||
|
||||
/**
|
||||
* `amy graperank followers [OBSERVER] [flags]` — the reverse crawl: find every
|
||||
* user who FOLLOWS the observer and persist their kind:3 contact lists.
|
||||
*
|
||||
* The outbox model (what `crawl` uses) walks follows *outward* and can't find
|
||||
* followers — you don't know a follower exists until you've seen their list, so
|
||||
* you can't route to their outbox first. This casts a wide net instead: it asks
|
||||
* as many relays as possible for kind:3 events that `#p`-tag the observer, paging
|
||||
* each relay past its per-REQ cap. The relay universe is "all possible relays":
|
||||
* the reachability-cache live set + every kind:10002 read/write relay and every
|
||||
* kind:30166 monitored relay in the store + the index/aggregator relays that hold
|
||||
* reverse-follow data for the network. Proven-dead relays are skipped
|
||||
* (`--no-reachability-cache` to query them anyway).
|
||||
*
|
||||
* Each follower's list lands in the store, so it also enriches the graph a later
|
||||
* `graperank score` builds — every follower becomes a FOLLOW edge into the
|
||||
* observer. Idempotent and cumulative; run it a few times for completeness.
|
||||
*
|
||||
* Flags: `--relay URL[,URL…]` (query only these instead of the whole universe),
|
||||
* `--page-limit N` (per-page REQ limit, default 500), `--timeout SECS` (per-page
|
||||
* EOSE watchdog, default 15), `--relay-concurrency N` (relays paged at once,
|
||||
* default 16), `--insert-batch N` (events per store commit, default 500).
|
||||
*/
|
||||
private suspend fun followers(
|
||||
dataDir: DataDir,
|
||||
rest: Array<String>,
|
||||
): Int {
|
||||
val args = Args(rest)
|
||||
val observerArg = args.positionalOrNull(0)
|
||||
val relayArg = args.flag("relay")
|
||||
val relayConcurrency = args.intFlag("relay-concurrency", args.intFlag("concurrency", 16))
|
||||
|
||||
// Read-only + never signs, so it runs anonymously — but then an OBSERVER
|
||||
// positional is required (there's no account to default to).
|
||||
Context.openOrAnonymous(dataDir).use { ctx ->
|
||||
ctx.prepare()
|
||||
if (ctx.anonymous && observerArg == null) {
|
||||
return Output.error("bad_args", "usage: amy graperank followers OBSERVER [flags] (no account — pass an observer)")
|
||||
}
|
||||
val observer = observerArg?.let { ctx.requireUserHex(it) } ?: ctx.identity.pubKeyHex
|
||||
|
||||
val relays =
|
||||
relayArg
|
||||
?.split(",")
|
||||
?.mapNotNull { RelayUrlNormalizer.normalizeOrNull(it.trim()) }
|
||||
?.toSet()
|
||||
?.takeIf { it.isNotEmpty() }
|
||||
?: allKnownRelays(ctx, args)
|
||||
|
||||
if (relays.isEmpty()) {
|
||||
return Output.error("no_relays", "no relays to query — run `amy relay probe` / `amy graperank crawl` first, or pass --relay")
|
||||
}
|
||||
|
||||
System.err.println("[followers] querying ${relays.size} relays for kind:3 #p=${observer.take(8)}…")
|
||||
val crawler =
|
||||
FollowerCrawler(
|
||||
client = ctx.client,
|
||||
store = ctx.store,
|
||||
config =
|
||||
FollowerCrawler.Config(
|
||||
relays = relays,
|
||||
pageLimit = args.intFlag("page-limit", 500),
|
||||
timeoutMs = args.longFlag("timeout", 15L) * 1000,
|
||||
maxConcurrentRelays = relayConcurrency,
|
||||
insertBatchSize = args.intFlag("insert-batch", 500),
|
||||
),
|
||||
log = { System.err.println(it) },
|
||||
)
|
||||
val stats = crawler.crawl(observer)
|
||||
|
||||
Output.emit(
|
||||
linkedMapOf<String, Any?>(
|
||||
"observer" to observer,
|
||||
"relays_queried" to stats.relaysQueried,
|
||||
"relays_answered" to stats.relaysAnswered,
|
||||
"followers_found" to stats.followersFound,
|
||||
"events_stored" to stats.eventsStored,
|
||||
"download_ms" to stats.downloadMs,
|
||||
),
|
||||
)
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
/**
|
||||
* "All possible relays" for the follower crawl: the reachability-cache live set,
|
||||
* every kind:10002 read/write relay and every kind:30166 monitored relay the
|
||||
* store knows, and the index/aggregator relays that hold reverse-follow data —
|
||||
* minus the ones a prior probe proved dead (unless `--no-reachability-cache`).
|
||||
*/
|
||||
private suspend fun allKnownRelays(
|
||||
ctx: Context,
|
||||
args: Args,
|
||||
): Set<NormalizedRelayUrl> {
|
||||
val reach = ctx.reachability.snapshot()
|
||||
val relays = LinkedHashSet<NormalizedRelayUrl>()
|
||||
relays.addAll(reach.live)
|
||||
// Every advertised outbox/inbox relay in the store — the homes of the users
|
||||
// whose followers we're hunting are exactly where those followers publish too.
|
||||
for (event in ctx.store.query<Event>(Filter(kinds = listOf(AdvertisedRelayListEvent.KIND)))) {
|
||||
if (event is AdvertisedRelayListEvent) relays.addAll(event.relaysNorm())
|
||||
}
|
||||
// Relays a NIP-66 monitor has ever seen (kind:30166) — the widest census.
|
||||
for (event in ctx.store.query<Event>(Filter(kinds = listOf(RelayDiscoveryEvent.KIND)))) {
|
||||
if (event is RelayDiscoveryEvent) event.relay()?.let { relays.add(it) }
|
||||
}
|
||||
// Aggregators/indexers hold reverse-follow data for users whose own outbox
|
||||
// we've never learned — the biggest completeness lever for followers.
|
||||
relays.addAll(CONTENT_AGGREGATOR_RELAYS)
|
||||
relays.addAll(EXTRA_DISCOVERY_RELAYS)
|
||||
relays.addAll(ctx.bootstrapRelays())
|
||||
relays.addAll(Constants.eventFinderRelays)
|
||||
|
||||
if (!args.bool(NO_REACHABILITY_CACHE_FLAG)) relays.removeAll(reach.dead)
|
||||
return relays
|
||||
}
|
||||
|
||||
/**
|
||||
* `amy graperank status` — read-only inventory of everything a GrapeRank run
|
||||
* depends on, straight from the local store: WoT record counts (the "do I
|
||||
|
||||
+189
@@ -0,0 +1,189 @@
|
||||
/*
|
||||
* 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.quartz.experimental.graperank
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.ParallelEventVerifier
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPagesFromPool
|
||||
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.tags.people.isTaggedUser
|
||||
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlin.time.TimeSource
|
||||
|
||||
/**
|
||||
* The reverse of [GrapeRankCrawler]: finds every user who **follows** the observer
|
||||
* by asking as many relays as possible for kind:3 contact lists that `#p`-tag the
|
||||
* observer, and persists each one to [store]. The outbox model can't find these —
|
||||
* you don't know a follower until you've seen their list, so you can't route to
|
||||
* their outbox first — so this casts a wide net across the whole relay universe
|
||||
* ([Config.relays], assembled by the caller from the reachability cache, the
|
||||
* kind:10002 relays in the store, and the index/aggregator relays that hold
|
||||
* reverse-follow data for the network).
|
||||
*
|
||||
* Each discovered kind:3 is a follower's full contact list, so persisting it also
|
||||
* enriches the graph a later `graperank score` builds from the store — the
|
||||
* follower becomes a FOLLOW edge into the observer (and into everyone else it
|
||||
* follows). A popular observer can have far more followers than a relay's per-REQ
|
||||
* cap, so every relay is paged on its own `created_at` cursor via
|
||||
* [fetchAllPagesFromPool].
|
||||
*
|
||||
* Transport-agnostic within quartz: it takes an [INostrClient] and an
|
||||
* [IEventStore]; relay *policy* (which relays make up "all possible relays") is the
|
||||
* caller's, injected through [Config]. Progress is emitted through [log].
|
||||
*/
|
||||
class FollowerCrawler(
|
||||
private val client: INostrClient,
|
||||
private val store: IEventStore,
|
||||
private val config: Config,
|
||||
private val log: (String) -> Unit = {},
|
||||
) {
|
||||
/**
|
||||
* @param relays the relay universe to query — "all possible relays" is the
|
||||
* caller's policy (reachability-cache live set + every kind:10002 relay in the
|
||||
* store + the index/aggregator relays). Empty means nothing to do.
|
||||
* @param pageLimit per-page `limit` handed to each relay's paged REQ; the cursor
|
||||
* walks past it, so this only bounds one page, not the total.
|
||||
* @param timeoutMs per-page EOSE timeout for a relay before its next page fires.
|
||||
* @param maxConcurrentRelays how many relays page at once (a global fan-out cap).
|
||||
* @param insertBatchSize verified events group-committed per [IEventStore.batchInsert].
|
||||
*/
|
||||
class Config(
|
||||
val relays: Set<NormalizedRelayUrl>,
|
||||
val pageLimit: Int = 500,
|
||||
val timeoutMs: Long = 15_000,
|
||||
val maxConcurrentRelays: Int = 16,
|
||||
val insertBatchSize: Int = 500,
|
||||
)
|
||||
|
||||
/** What the follower crawl found. */
|
||||
class Stats(
|
||||
val relaysQueried: Int,
|
||||
/** Relays that returned at least one valid follower list. */
|
||||
val relaysAnswered: Int,
|
||||
/** Distinct followers (kind:3 authors that `#p`-tag the observer). */
|
||||
val followersFound: Int,
|
||||
/** Verified kind:3 events handed to the store (superseded versions included). */
|
||||
val eventsStored: Long,
|
||||
val downloadMs: Long,
|
||||
)
|
||||
|
||||
/**
|
||||
* Find and persist every follower of [observer] across [Config.relays]. Verifies
|
||||
* each event (a relay can lie), keeps only kind:3 that actually `#p`-tag the
|
||||
* observer (a relay can over-return), dedups by event id, and group-commits to
|
||||
* the store. Returns [Stats].
|
||||
*/
|
||||
suspend fun crawl(observer: HexKey): Stats =
|
||||
withContext(Dispatchers.IO) {
|
||||
if (config.relays.isEmpty()) return@withContext Stats(0, 0, 0, 0L, 0L)
|
||||
|
||||
val mark = TimeSource.Monotonic.markNow()
|
||||
val filter = Filter(kinds = FOLLOW_KINDS, tags = mapOf("p" to listOf(observer)), limit = config.pageLimit)
|
||||
val perRelay = config.relays.associateWith { listOf(filter) }
|
||||
|
||||
// These three are read/written ONLY by the verifier's single drain
|
||||
// coroutine (ParallelEventVerifier dispatches onVerified in order from one
|
||||
// coroutine), so plain collections are safe — no concurrent access.
|
||||
val followers = HashSet<HexKey>()
|
||||
val seen = HashSet<HexKey>()
|
||||
val answered = HashSet<NormalizedRelayUrl>()
|
||||
var eventsStored = 0L
|
||||
|
||||
// Verified follower lists cross to a single persister coroutine that can
|
||||
// suspend on the store; onVerified itself must not (it's a plain callback).
|
||||
val toPersist = Channel<Event>(Channel.UNLIMITED)
|
||||
|
||||
coroutineScope {
|
||||
val persister =
|
||||
launch(Dispatchers.IO) {
|
||||
val flushAt = config.insertBatchSize.coerceAtLeast(1)
|
||||
val buffer = ArrayList<Event>(flushAt)
|
||||
|
||||
suspend fun flush() {
|
||||
if (buffer.isEmpty()) return
|
||||
store.batchInsert(buffer)
|
||||
eventsStored += buffer.size
|
||||
buffer.clear()
|
||||
}
|
||||
for (event in toPersist) {
|
||||
buffer.add(event)
|
||||
if (buffer.size >= flushAt) flush()
|
||||
}
|
||||
flush()
|
||||
}
|
||||
|
||||
val verifier =
|
||||
ParallelEventVerifier<NormalizedRelayUrl>(
|
||||
scope = this,
|
||||
onVerified = { event, relay ->
|
||||
// A relay can over-return; keep only genuine followers, and
|
||||
// dedup by id so a list mirrored across relays is stored once.
|
||||
if (event is ContactListEvent && event.isTaggedUser(observer) && seen.add(event.id)) {
|
||||
followers.add(event.pubKey)
|
||||
answered.add(relay)
|
||||
toPersist.trySend(event)
|
||||
}
|
||||
},
|
||||
)
|
||||
|
||||
client.fetchAllPagesFromPool(
|
||||
filters = perRelay,
|
||||
timeoutMs = config.timeoutMs,
|
||||
maxConcurrentRelays = config.maxConcurrentRelays,
|
||||
onRelayComplete = { relay, total ->
|
||||
if (total > 0) log("[followers] ${relay.url}: $total kind:3 pages drained")
|
||||
},
|
||||
) { event, relay ->
|
||||
// onEvent runs on the relay reader thread and must not suspend;
|
||||
// submit() is a non-suspending channel send that backpressures.
|
||||
if (event.kind == ContactListEvent.KIND) verifier.submit(event, relay)
|
||||
}
|
||||
|
||||
// No more events will arrive: drain the verifier, then the persister.
|
||||
verifier.close()
|
||||
verifier.join()
|
||||
toPersist.close()
|
||||
persister.join()
|
||||
}
|
||||
|
||||
Stats(
|
||||
relaysQueried = config.relays.size,
|
||||
relaysAnswered = answered.size,
|
||||
followersFound = followers.size,
|
||||
eventsStored = eventsStored,
|
||||
downloadMs = mark.elapsedNow().inWholeMilliseconds,
|
||||
)
|
||||
}
|
||||
|
||||
companion object {
|
||||
/** Reverse-follow lookup is kind:3 only — a follower IS the author of a kind:3 that `#p`-tags the observer. */
|
||||
private val FOLLOW_KINDS = listOf(ContactListEvent.KIND)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user