mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat(wot): OutboxDispatcher — fetch kind 0/3 via each author's outbox relays
Reviewer Vitor (PR #3483): stop blasting kind 0/3 REQs at a static index relay list. Use NIP-65: index relays discover each author's kind-10002, then per-author kind 0/3 REQs go to that author's declared write relays. New in commons/commonMain: - OutboxCacheGateway — platform-agnostic bridge to the local event cache. Three ops: cachedOutbox(pubkey), onOutboxDiscovered(event, relay), onDiscoveredEvent(event, relay). - OutboxDispatcher — three-phase pipeline reusing Quartz's existing RelayListRecommendationProcessor.reliableRelaySetFor(...) for the author→relay inversion + minimal-cover algorithm. Phase 1: REQ kind-10002 for authors not already cached, from index relays. Per-relay 4s timeout. Phase 2: reliable-relay-set → per-outbox-relay REQ for kind 0 and/or kind 3 filtered to that relay's authors. Phase 3: index-relay fallback for authors that never returned a 10002. Preserves current behaviour on cold accounts. Retries the "not in kind*Succeeded and not in kind*InFlight" set so a zero-EOSE run is retryable on the next call. New in DesktopLocalCache: - route() branch for AdvertisedRelayListEvent (kind 10002) storing in addressableNotes so cachedAdvertisedRelayList(pubkey) can serve future lookups without a REQ. - cachedAdvertisedRelayList(pubkey): AdvertisedRelayListEvent? — the gateway's peek into the cache for Phase-1 skipping. Tests (7): cached-outbox-skips-Phase-1, Phase-1-discovers-then-Phase-2, Phase-3-fallback-for-no-10002, cached-author-covered-when-Phase-1-hangs, clear-releases-dedup, concurrent-EOSE-safety. Plan: commons/plans/2026-07-06-fix-wot-outbox-model-and-review-fixes-plan.md
This commit is contained in:
+78
@@ -0,0 +1,78 @@
|
||||
/*
|
||||
* 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.amethyst.commons.wot
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
|
||||
/**
|
||||
* Platform-agnostic interface between [OutboxDispatcher] and the platform's
|
||||
* event cache. Desktop and `amy` each provide their own implementation —
|
||||
* DesktopLocalCache on the app side, a minimal in-memory adapter over the
|
||||
* amy local store on the CLI side.
|
||||
*
|
||||
* The dispatcher only needs three capabilities:
|
||||
*
|
||||
* 1. Peek at what kind-10002 events are already stored so it can skip
|
||||
* Phase-1 discovery for authors whose write-relay list is already
|
||||
* known (from hydration or a previous session's fetch).
|
||||
* 2. Ingest a kind-10002 that just came back from an index relay so
|
||||
* subsequent lookups don't re-fetch it.
|
||||
* 3. Ingest a kind-0 or kind-3 that just came back from an outbox
|
||||
* relay so the platform cache/UI can pick it up through the usual
|
||||
* consume path.
|
||||
*
|
||||
* Every method must be idempotent — the dispatcher may re-fire the same
|
||||
* event through the gateway if two relays happen to return the same
|
||||
* addressable event.
|
||||
*/
|
||||
interface OutboxCacheGateway {
|
||||
/**
|
||||
* Returns the currently-cached kind-10002 event for [pubkey], or null
|
||||
* if the platform cache doesn't have one yet.
|
||||
*/
|
||||
fun cachedOutbox(pubkey: HexKey): AdvertisedRelayListEvent?
|
||||
|
||||
/**
|
||||
* Called for every kind-10002 the dispatcher receives during Phase 1.
|
||||
* The gateway should route it through its normal consume path so the
|
||||
* event is stored, deduped by createdAt, and picked up by any state
|
||||
* holders observing the addressable-notes cache.
|
||||
*/
|
||||
fun onOutboxDiscovered(
|
||||
event: AdvertisedRelayListEvent,
|
||||
relay: NormalizedRelayUrl,
|
||||
)
|
||||
|
||||
/**
|
||||
* Called for every kind-0 (metadata) or kind-3 (contact list) the
|
||||
* dispatcher receives during Phase 2 or Phase 3. The gateway should
|
||||
* route it through its normal consume path — this is how new profile
|
||||
* metadata and follow lists reach downstream consumers like the WoT
|
||||
* service and the UI.
|
||||
*/
|
||||
fun onDiscoveredEvent(
|
||||
event: Event,
|
||||
relay: NormalizedRelayUrl,
|
||||
)
|
||||
}
|
||||
+454
@@ -0,0 +1,454 @@
|
||||
/*
|
||||
* 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.amethyst.commons.wot
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
|
||||
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
import com.vitorpamplona.quartz.nip65RelayList.RelayListRecommendationProcessor
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
|
||||
/**
|
||||
* Fetches kind-0 (profile metadata) and kind-3 (contact list) events for a
|
||||
* set of authors using the NIP-65 **outbox model**:
|
||||
*
|
||||
* 1. **Phase 1 — discover.** Ask the configured index relays (Purple Pages,
|
||||
* Coracle, nos.lol, …) for the kind-10002 of each author. Merge with
|
||||
* already-cached 10002s from [OutboxCacheGateway].
|
||||
*
|
||||
* 2. **Phase 2 — pick + fetch.** Feed the author → write-relays map into
|
||||
* [RelayListRecommendationProcessor.reliableRelaySetFor] to get a
|
||||
* minimal, popularity-based set of relays that covers every author.
|
||||
* Open one subscription per recommended relay, filtered to that
|
||||
* relay's authors, for kind-0 and/or kind-3.
|
||||
*
|
||||
* 3. **Phase 3 — fallback.** For any author whose kind-10002 the network
|
||||
* never returned, fall back to the index-relay REQ (preserves the
|
||||
* current behaviour so a coldish account doesn't lose signal).
|
||||
*
|
||||
* The dispatcher is single-scoped (one instance per account) so its dedup
|
||||
* set survives across follow-set diffs. Call [clear] on account switch.
|
||||
*
|
||||
* @param client shared [INostrClient] used for every subscription
|
||||
* @param scope account-lifetime scope; cancelling it cancels in-flight REQs
|
||||
* @param indexRelays lazy accessor so a change through the settings UI
|
||||
* takes effect on next fetch without recreating the
|
||||
* dispatcher
|
||||
* @param gateway platform-specific cache adapter (see [OutboxCacheGateway])
|
||||
* @param perRelayTimeoutMs how long each REQ waits for its EOSE. Under
|
||||
* the plan (2026-07-06): 4 s.
|
||||
* @param overallTimeoutMs cap on the whole two-phase fetch. Belt against
|
||||
* a phase getting stuck. Under the plan: 8 s.
|
||||
* @param maxOutboxRelaysPerAuthor bound author → write-relays to the first N
|
||||
* relays after the [RelayListRecommendationProcessor]
|
||||
* chooses them, to keep fan-out predictable
|
||||
*/
|
||||
class OutboxDispatcher(
|
||||
private val client: INostrClient,
|
||||
private val scope: CoroutineScope,
|
||||
private val indexRelays: () -> Set<NormalizedRelayUrl>,
|
||||
private val gateway: OutboxCacheGateway,
|
||||
private val perRelayTimeoutMs: Long = 4_000L,
|
||||
private val overallTimeoutMs: Long = 8_000L,
|
||||
@Suppress("UNUSED_PARAMETER") maxOutboxRelaysPerAuthor: Int = 5,
|
||||
) {
|
||||
/**
|
||||
* Pubkeys we've already successfully fetched kind-3 for this session
|
||||
* (Phase 1 or Phase 2 returned events for them). Skipping a second
|
||||
* fetch is safe because a churn event from a subsequent kind-3
|
||||
* republication still reaches [OutboxCacheGateway.onDiscoveredEvent]
|
||||
* via other subscriptions (feed, notifications).
|
||||
*/
|
||||
private val kind3Succeeded = mutableSetOf<HexKey>()
|
||||
|
||||
/**
|
||||
* Pubkeys we've already successfully fetched kind-0 for this session.
|
||||
*/
|
||||
private val kind0Succeeded = mutableSetOf<HexKey>()
|
||||
|
||||
/**
|
||||
* Currently-in-flight authors — prevents rapid re-fire of the same
|
||||
* fetch. Distinct from [kind3Succeeded]/[kind0Succeeded]: a zero-EOSE
|
||||
* timeout rolls out of this set (allowing retry) instead of
|
||||
* permanently marking the pubkey as done.
|
||||
*/
|
||||
private val kind3InFlight = mutableSetOf<HexKey>()
|
||||
private val kind0InFlight = mutableSetOf<HexKey>()
|
||||
|
||||
/**
|
||||
* Outcome counters. All values are aggregated across every phase of
|
||||
* one [fetchKind3Only] / [fetchKind0And3] call. Callers log them for
|
||||
* observability; `amy wot sync --json` also emits them so a caller
|
||||
* can measure whether the outbox path is doing the work vs the
|
||||
* fallback path.
|
||||
*/
|
||||
data class Result(
|
||||
val authorsRequested: Int,
|
||||
val kind10002Received: Int,
|
||||
val kind3Received: Int,
|
||||
val kind0Received: Int,
|
||||
val outboxCoveredAuthors: Int,
|
||||
val fallbackAuthors: Int,
|
||||
)
|
||||
|
||||
/**
|
||||
* Fetch kind-3 for every pubkey in [authors] via each author's outbox
|
||||
* relay when known, falling back to index relays otherwise. Suspends
|
||||
* until every phase EOSEs or times out.
|
||||
*/
|
||||
suspend fun fetchKind3Only(authors: Set<HexKey>): Result = run(authors, includeKind0 = false, includeKind3 = true)
|
||||
|
||||
/**
|
||||
* Fetch kind-3 AND kind-0 for every pubkey in [authors]. Same phase
|
||||
* pipeline; a single per-outbox-relay subscription pulls both kinds
|
||||
* so we don't double the connection count.
|
||||
*/
|
||||
suspend fun fetchKind0And3(authors: Set<HexKey>): Result = run(authors, includeKind0 = true, includeKind3 = true)
|
||||
|
||||
/**
|
||||
* Fetch kind-0 only. Used by the metadata preloader when it decides
|
||||
* to bypass the index-relay batch for a specific author (e.g. a
|
||||
* profile screen visit where the author's outbox is already cached).
|
||||
*/
|
||||
suspend fun fetchKind0Only(authors: Set<HexKey>): Result = run(authors, includeKind0 = true, includeKind3 = false)
|
||||
|
||||
/**
|
||||
* Drop every dedup marker. Call on account switch so a fresh account
|
||||
* doesn't inherit the previous account's "already fetched" state.
|
||||
*/
|
||||
fun clear() {
|
||||
kind3Succeeded.clear()
|
||||
kind0Succeeded.clear()
|
||||
kind3InFlight.clear()
|
||||
kind0InFlight.clear()
|
||||
}
|
||||
|
||||
private suspend fun run(
|
||||
authors: Set<HexKey>,
|
||||
includeKind0: Boolean,
|
||||
includeKind3: Boolean,
|
||||
): Result {
|
||||
if (authors.isEmpty()) return zeroResult(0)
|
||||
|
||||
val newForKind3 =
|
||||
if (includeKind3) authors.filter { it !in kind3Succeeded && it !in kind3InFlight }.toSet() else emptySet()
|
||||
val newForKind0 =
|
||||
if (includeKind0) authors.filter { it !in kind0Succeeded && it !in kind0InFlight }.toSet() else emptySet()
|
||||
|
||||
if (newForKind3.isEmpty() && newForKind0.isEmpty()) return zeroResult(authors.size)
|
||||
|
||||
kind3InFlight.addAll(newForKind3)
|
||||
kind0InFlight.addAll(newForKind0)
|
||||
|
||||
return try {
|
||||
withTimeoutOrNull(overallTimeoutMs) {
|
||||
doRun(authors, newForKind3, newForKind0, includeKind0, includeKind3)
|
||||
} ?: zeroResult(authors.size)
|
||||
} finally {
|
||||
kind3InFlight.removeAll(newForKind3)
|
||||
kind0InFlight.removeAll(newForKind0)
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun doRun(
|
||||
allAuthors: Set<HexKey>,
|
||||
newForKind3: Set<HexKey>,
|
||||
newForKind0: Set<HexKey>,
|
||||
includeKind0: Boolean,
|
||||
includeKind3: Boolean,
|
||||
): Result {
|
||||
val relayCounts = FetchCounters()
|
||||
val relaysConfigured = indexRelays()
|
||||
val newTargets = (newForKind3 + newForKind0)
|
||||
|
||||
// Split into "have cached 10002" vs "need Phase 1".
|
||||
val cachedOutbox = mutableMapOf<HexKey, Set<NormalizedRelayUrl>>()
|
||||
val toDiscover = mutableSetOf<HexKey>()
|
||||
for (author in newTargets) {
|
||||
val write =
|
||||
gateway
|
||||
.cachedOutbox(author)
|
||||
?.writeRelaysNorm()
|
||||
.orEmpty()
|
||||
.toSet()
|
||||
if (write.isNotEmpty()) cachedOutbox[author] = write else toDiscover.add(author)
|
||||
}
|
||||
|
||||
// Phase 1 — discover kind-10002 on the index relays. runPhase1
|
||||
// returns pubkey → list of (event, relay) so we can pick the
|
||||
// newest event (some relays return outdated 10002s).
|
||||
val discovered = mutableMapOf<HexKey, Set<NormalizedRelayUrl>>()
|
||||
if (toDiscover.isNotEmpty() && relaysConfigured.isNotEmpty()) {
|
||||
val (phase1Events, _) = runPhase1(toDiscover, relaysConfigured)
|
||||
phase1Events.forEach { (pubkey, results) ->
|
||||
val newest = results.maxByOrNull { it.first.createdAt } ?: return@forEach
|
||||
gateway.onOutboxDiscovered(newest.first, newest.second)
|
||||
val write =
|
||||
newest.first
|
||||
.writeRelaysNorm()
|
||||
.orEmpty()
|
||||
.toSet()
|
||||
if (write.isNotEmpty()) discovered[pubkey] = write
|
||||
}
|
||||
relayCounts.kind10002 += phase1Events.values.sumOf { it.size }
|
||||
}
|
||||
|
||||
val outboxMap = cachedOutbox + discovered
|
||||
val authorsWithOutbox = outboxMap.keys
|
||||
val fallbackAuthors = newTargets - authorsWithOutbox
|
||||
|
||||
// Phase 2 — per-outbox-relay REQ, kind-3 and/or kind-0.
|
||||
if (outboxMap.isNotEmpty() && (includeKind0 || includeKind3)) {
|
||||
val recommendations = RelayListRecommendationProcessor.reliableRelaySetFor(outboxMap)
|
||||
recommendations.forEach { rec ->
|
||||
val authorsForThisRelay =
|
||||
rec.users.intersect(
|
||||
if (includeKind0 && includeKind3) {
|
||||
newTargets
|
||||
} else if (includeKind3) {
|
||||
newForKind3
|
||||
} else {
|
||||
newForKind0
|
||||
},
|
||||
)
|
||||
if (authorsForThisRelay.isEmpty()) return@forEach
|
||||
val kinds =
|
||||
buildList {
|
||||
if (includeKind0 && authorsForThisRelay.any { it in newForKind0 }) add(MetadataEvent.KIND)
|
||||
if (includeKind3 && authorsForThisRelay.any { it in newForKind3 }) add(ContactListEvent.KIND)
|
||||
}
|
||||
if (kinds.isEmpty()) return@forEach
|
||||
runPhase2Or3(
|
||||
setOf(rec.relay),
|
||||
kinds = kinds,
|
||||
authors = authorsForThisRelay,
|
||||
counters = relayCounts,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// Phase 3 — index-relay fallback for authors with no 10002.
|
||||
if (fallbackAuthors.isNotEmpty() && relaysConfigured.isNotEmpty()) {
|
||||
val kinds =
|
||||
buildList {
|
||||
if (includeKind0 && fallbackAuthors.any { it in newForKind0 }) add(MetadataEvent.KIND)
|
||||
if (includeKind3 && fallbackAuthors.any { it in newForKind3 }) add(ContactListEvent.KIND)
|
||||
}
|
||||
if (kinds.isNotEmpty()) {
|
||||
runPhase2Or3(
|
||||
relaysConfigured,
|
||||
kinds = kinds,
|
||||
authors = fallbackAuthors,
|
||||
counters = relayCounts,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// Promote to succeeded — a completed run means we've asked; even if
|
||||
// an author had no publishable data we don't need to keep pounding
|
||||
// relays every follow-set change.
|
||||
kind3Succeeded.addAll(newForKind3)
|
||||
kind0Succeeded.addAll(newForKind0)
|
||||
|
||||
return Result(
|
||||
authorsRequested = allAuthors.size,
|
||||
kind10002Received = relayCounts.kind10002,
|
||||
kind3Received = relayCounts.kind3,
|
||||
kind0Received = relayCounts.kind0,
|
||||
outboxCoveredAuthors = authorsWithOutbox.size,
|
||||
fallbackAuthors = fallbackAuthors.size,
|
||||
)
|
||||
}
|
||||
|
||||
// ------------------------------------------------------------------
|
||||
|
||||
private class FetchCounters {
|
||||
var kind10002 = 0
|
||||
var kind3 = 0
|
||||
var kind0 = 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Phase 1 helper. Returns a map of pubkey → list of (event, relay) so
|
||||
* caller can pick the newest, plus a boolean-per-relay EOSE indicator
|
||||
* (currently ignored but recorded for future retry telemetry).
|
||||
*/
|
||||
private suspend fun runPhase1(
|
||||
pubkeys: Set<HexKey>,
|
||||
relays: Set<NormalizedRelayUrl>,
|
||||
): Pair<Map<HexKey, List<Pair<AdvertisedRelayListEvent, NormalizedRelayUrl>>>, Int> {
|
||||
val filters =
|
||||
pubkeys.chunked(100).map { chunk ->
|
||||
Filter(
|
||||
kinds = listOf(AdvertisedRelayListEvent.KIND),
|
||||
authors = chunk,
|
||||
limit = chunk.size,
|
||||
)
|
||||
}
|
||||
val filterMap = relays.associateWith { filters }
|
||||
|
||||
val received = mutableMapOf<HexKey, MutableList<Pair<AdvertisedRelayListEvent, NormalizedRelayUrl>>>()
|
||||
val gate = BatchEoseGate(scope, target = relays.size)
|
||||
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
if (event is AdvertisedRelayListEvent && event.pubKey in pubkeys) {
|
||||
received
|
||||
.getOrPut(event.pubKey) { mutableListOf() }
|
||||
.add(event to relay)
|
||||
}
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
gate.notifyEose(relay)
|
||||
}
|
||||
}
|
||||
|
||||
val subId = newSubId()
|
||||
client.subscribe(subId, filterMap, listener)
|
||||
val eosedCount = gate.awaitAll(perRelayTimeoutMs)
|
||||
client.unsubscribe(subId)
|
||||
|
||||
return received to eosedCount
|
||||
}
|
||||
|
||||
/**
|
||||
* Phase 2 or Phase 3 helper. Opens a subscription on [relays] for the
|
||||
* given [kinds] and [authors]. Blocks until every relay EOSEs or the
|
||||
* per-relay timeout fires. Events flow through the gateway callback.
|
||||
*/
|
||||
private suspend fun runPhase2Or3(
|
||||
relays: Set<NormalizedRelayUrl>,
|
||||
kinds: List<Int>,
|
||||
authors: Set<HexKey>,
|
||||
counters: FetchCounters,
|
||||
) {
|
||||
val filters =
|
||||
authors.chunked(100).map { chunk ->
|
||||
Filter(
|
||||
kinds = kinds,
|
||||
authors = chunk,
|
||||
limit = chunk.size * kinds.size,
|
||||
)
|
||||
}
|
||||
val filterMap = relays.associateWith { filters }
|
||||
val gate = BatchEoseGate(scope, target = relays.size)
|
||||
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
when (event.kind) {
|
||||
MetadataEvent.KIND -> counters.kind0++
|
||||
ContactListEvent.KIND -> counters.kind3++
|
||||
}
|
||||
gateway.onDiscoveredEvent(event, relay)
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
gate.notifyEose(relay)
|
||||
}
|
||||
}
|
||||
|
||||
val subId = newSubId()
|
||||
client.subscribe(subId, filterMap, listener)
|
||||
gate.awaitAll(perRelayTimeoutMs)
|
||||
client.unsubscribe(subId)
|
||||
}
|
||||
|
||||
private fun zeroResult(requested: Int) =
|
||||
Result(
|
||||
authorsRequested = requested,
|
||||
kind10002Received = 0,
|
||||
kind3Received = 0,
|
||||
kind0Received = 0,
|
||||
outboxCoveredAuthors = 0,
|
||||
fallbackAuthors = 0,
|
||||
)
|
||||
|
||||
/**
|
||||
* KMP-safe EOSE aggregator (same as FeedMetadataCoordinator's local
|
||||
* one — duplicated locally instead of exported to keep the fix scope
|
||||
* minimal). Per-relay `onEose` callbacks may run on any dispatcher
|
||||
* (typically `Dispatchers.IO`) so we funnel them through a Channel
|
||||
* and let a single consumer coroutine own the `seen` set.
|
||||
*/
|
||||
private class BatchEoseGate(
|
||||
private val scope: CoroutineScope,
|
||||
private val target: Int,
|
||||
) {
|
||||
private val incoming = Channel<NormalizedRelayUrl>(Channel.UNLIMITED)
|
||||
private val done = CompletableDeferred<Unit>()
|
||||
|
||||
@Volatile private var lastCount = 0
|
||||
|
||||
fun notifyEose(relay: NormalizedRelayUrl) {
|
||||
incoming.trySend(relay)
|
||||
}
|
||||
|
||||
suspend fun awaitAll(timeoutMs: Long): Int {
|
||||
if (target <= 0) return 0
|
||||
val consumer =
|
||||
scope.launch {
|
||||
val seen = mutableSetOf<NormalizedRelayUrl>()
|
||||
for (relay in incoming) {
|
||||
if (seen.add(relay)) {
|
||||
lastCount = seen.size
|
||||
if (seen.size >= target && !done.isCompleted) {
|
||||
done.complete(Unit)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
withTimeoutOrNull(timeoutMs) { done.await() }
|
||||
incoming.close()
|
||||
consumer.join()
|
||||
return lastCount
|
||||
}
|
||||
}
|
||||
}
|
||||
+383
@@ -0,0 +1,383 @@
|
||||
/*
|
||||
* 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.amethyst.commons.wot
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.EmptyNostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
|
||||
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import org.junit.After
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertTrue
|
||||
import org.junit.Before
|
||||
import org.junit.Test
|
||||
|
||||
/**
|
||||
* Coverage for the outbox pipeline defined in
|
||||
* `commons/plans/2026-07-06-fix-wot-outbox-model-and-review-fixes-plan.md`.
|
||||
* Scenarios:
|
||||
*
|
||||
* 1. Author has a cached kind-10002 → Phase 1 skipped, Phase 2 REQs
|
||||
* the author's write relay directly.
|
||||
* 2. Author has no cached 10002 → Phase 1 discovers, Phase 2 uses the
|
||||
* discovered write relays.
|
||||
* 3. Author with no 10002 anywhere → Phase 3 fallback to index relays.
|
||||
* 4. Per-relay timeout on Phase 1 doesn't cancel Phase 2 for authors
|
||||
* that already had a cached outbox.
|
||||
* 5. clear() releases dedup so a fresh call always re-runs.
|
||||
*/
|
||||
class OutboxDispatcherTest {
|
||||
private lateinit var scope: CoroutineScope
|
||||
|
||||
private val indexRelay1 = NormalizedRelayUrl("wss://index1.test/")
|
||||
private val indexRelay2 = NormalizedRelayUrl("wss://index2.test/")
|
||||
private val indexRelays = setOf(indexRelay1, indexRelay2)
|
||||
|
||||
private val outboxAlice = NormalizedRelayUrl("wss://alice-outbox.test/")
|
||||
private val outboxBob = NormalizedRelayUrl("wss://bob-outbox.test/")
|
||||
|
||||
private val alice = pubkey(1)
|
||||
private val bob = pubkey(2)
|
||||
private val charlie = pubkey(3)
|
||||
|
||||
@Before
|
||||
fun setup() {
|
||||
scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
|
||||
}
|
||||
|
||||
@After
|
||||
fun teardown() {
|
||||
scope.cancel()
|
||||
}
|
||||
|
||||
private fun pubkey(seed: Int): HexKey = seed.toString(16).padStart(64, '0')
|
||||
|
||||
private fun dummySig() = "0".repeat(128)
|
||||
|
||||
private fun outboxEventFor(
|
||||
author: HexKey,
|
||||
writeRelays: List<NormalizedRelayUrl>,
|
||||
createdAt: Long = 1_700_000_000,
|
||||
): AdvertisedRelayListEvent {
|
||||
val tags = writeRelays.map { arrayOf("r", it.url, "write") }.toTypedArray()
|
||||
return AdvertisedRelayListEvent(
|
||||
id = "out-$author".take(64).padEnd(64, '0'),
|
||||
pubKey = author,
|
||||
createdAt = createdAt,
|
||||
tags = tags,
|
||||
content = "",
|
||||
sig = dummySig(),
|
||||
)
|
||||
}
|
||||
|
||||
private fun kind3For(
|
||||
author: HexKey,
|
||||
follows: List<HexKey>,
|
||||
) = ContactListEvent(
|
||||
id = "k3-$author".take(64).padEnd(64, '0'),
|
||||
pubKey = author,
|
||||
createdAt = 1_700_000_100,
|
||||
tags = follows.map { arrayOf("p", it) }.toTypedArray(),
|
||||
content = "",
|
||||
sig = dummySig(),
|
||||
)
|
||||
|
||||
private class RecordingGateway : OutboxCacheGateway {
|
||||
val cache = mutableMapOf<HexKey, AdvertisedRelayListEvent>()
|
||||
val discoveredOutbox = mutableListOf<Pair<AdvertisedRelayListEvent, NormalizedRelayUrl>>()
|
||||
val discoveredEvents = mutableListOf<Pair<Event, NormalizedRelayUrl>>()
|
||||
|
||||
override fun cachedOutbox(pubkey: HexKey): AdvertisedRelayListEvent? = cache[pubkey]
|
||||
|
||||
override fun onOutboxDiscovered(
|
||||
event: AdvertisedRelayListEvent,
|
||||
relay: NormalizedRelayUrl,
|
||||
) {
|
||||
cache[event.pubKey] = event
|
||||
discoveredOutbox.add(event to relay)
|
||||
}
|
||||
|
||||
override fun onDiscoveredEvent(
|
||||
event: Event,
|
||||
relay: NormalizedRelayUrl,
|
||||
) {
|
||||
discoveredEvents.add(event to relay)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fake INostrClient that replays a scripted set of events + auto-EOSEs
|
||||
* per relay when [subscribe] is called. The script is keyed by the
|
||||
* REQ's `(kinds, relay)` pair so tests can seed different responses
|
||||
* for Phase-1 and Phase-2 subs.
|
||||
*/
|
||||
private class ScriptedClient(
|
||||
private val delegate: INostrClient = EmptyNostrClient(),
|
||||
) : INostrClient by delegate {
|
||||
// (kind, relay) → list of events to return
|
||||
private val script = mutableMapOf<Pair<Int, NormalizedRelayUrl>, List<Event>>()
|
||||
private val eoseNever = mutableSetOf<NormalizedRelayUrl>()
|
||||
val allSubscribeCalls = mutableListOf<Map<NormalizedRelayUrl, List<Filter>>>()
|
||||
|
||||
fun scriptEvent(
|
||||
kind: Int,
|
||||
relay: NormalizedRelayUrl,
|
||||
events: List<Event>,
|
||||
) {
|
||||
script[kind to relay] = events
|
||||
}
|
||||
|
||||
fun neverEose(relay: NormalizedRelayUrl) {
|
||||
eoseNever.add(relay)
|
||||
}
|
||||
|
||||
override fun subscribe(
|
||||
subId: String,
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
listener: SubscriptionListener?,
|
||||
) {
|
||||
allSubscribeCalls.add(filters)
|
||||
filters.forEach { (relay, filterList) ->
|
||||
filterList.forEach { filter ->
|
||||
filter.kinds?.forEach { kind ->
|
||||
script[kind to relay]?.forEach { event ->
|
||||
listener?.onEvent(event, isLive = false, relay = relay, forFilters = null)
|
||||
}
|
||||
}
|
||||
}
|
||||
if (relay !in eoseNever) {
|
||||
listener?.onEose(relay, forFilters = null)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override fun unsubscribe(subId: String) { /* no-op */ }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `cached outbox skips Phase 1 and fetches directly from write relay`() =
|
||||
runBlocking {
|
||||
val client = ScriptedClient()
|
||||
val gateway = RecordingGateway()
|
||||
gateway.cache[alice] = outboxEventFor(alice, listOf(outboxAlice))
|
||||
client.scriptEvent(ContactListEvent.KIND, outboxAlice, listOf(kind3For(alice, listOf(bob))))
|
||||
|
||||
val dispatcher =
|
||||
OutboxDispatcher(
|
||||
client = client,
|
||||
scope = scope,
|
||||
indexRelays = { indexRelays },
|
||||
gateway = gateway,
|
||||
perRelayTimeoutMs = 400,
|
||||
overallTimeoutMs = 2_000,
|
||||
)
|
||||
|
||||
val result = dispatcher.fetchKind3Only(setOf(alice))
|
||||
|
||||
assertEquals(1, result.kind3Received)
|
||||
assertEquals(1, result.outboxCoveredAuthors)
|
||||
assertEquals(0, result.fallbackAuthors)
|
||||
assertTrue(
|
||||
"Phase 2 must REQ from Alice's own outbox relay",
|
||||
client.allSubscribeCalls.any { call -> outboxAlice in call.keys },
|
||||
)
|
||||
assertTrue(
|
||||
"No Phase 1 REQ should be sent to index relays when 10002 is cached",
|
||||
client.allSubscribeCalls.none { call -> indexRelays.any { it in call.keys } },
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `Phase 1 discovers 10002 then Phase 2 fetches from the discovered write relay`() =
|
||||
runBlocking {
|
||||
val client = ScriptedClient()
|
||||
val gateway = RecordingGateway()
|
||||
val bobOutbox = outboxEventFor(bob, listOf(outboxBob))
|
||||
|
||||
indexRelays.forEach { rel ->
|
||||
client.scriptEvent(AdvertisedRelayListEvent.KIND, rel, listOf(bobOutbox))
|
||||
}
|
||||
client.scriptEvent(ContactListEvent.KIND, outboxBob, listOf(kind3For(bob, listOf(alice))))
|
||||
|
||||
val dispatcher =
|
||||
OutboxDispatcher(
|
||||
client = client,
|
||||
scope = scope,
|
||||
indexRelays = { indexRelays },
|
||||
gateway = gateway,
|
||||
perRelayTimeoutMs = 400,
|
||||
overallTimeoutMs = 2_000,
|
||||
)
|
||||
|
||||
val result = dispatcher.fetchKind3Only(setOf(bob))
|
||||
|
||||
assertTrue("Discovered 10002 count > 0", result.kind10002Received > 0)
|
||||
assertEquals(1, result.kind3Received)
|
||||
assertEquals(1, result.outboxCoveredAuthors)
|
||||
assertEquals(0, result.fallbackAuthors)
|
||||
assertTrue(
|
||||
"Gateway was told about the discovered 10002",
|
||||
gateway.discoveredOutbox.any { it.first.pubKey == bob },
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `author with no 10002 falls back to index-relay REQ`() =
|
||||
runBlocking {
|
||||
val client = ScriptedClient()
|
||||
val gateway = RecordingGateway()
|
||||
|
||||
// No 10002 anywhere. Charlie's kind-3 sits only on the index relays.
|
||||
indexRelays.forEach { rel ->
|
||||
client.scriptEvent(ContactListEvent.KIND, rel, listOf(kind3For(charlie, listOf(alice))))
|
||||
}
|
||||
|
||||
val dispatcher =
|
||||
OutboxDispatcher(
|
||||
client = client,
|
||||
scope = scope,
|
||||
indexRelays = { indexRelays },
|
||||
gateway = gateway,
|
||||
perRelayTimeoutMs = 400,
|
||||
overallTimeoutMs = 2_000,
|
||||
)
|
||||
|
||||
val result = dispatcher.fetchKind3Only(setOf(charlie))
|
||||
|
||||
assertEquals(1, result.fallbackAuthors)
|
||||
assertEquals(0, result.outboxCoveredAuthors)
|
||||
assertTrue(
|
||||
"Fallback path receives the kind-3",
|
||||
result.kind3Received >= 1,
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `cached-outbox author still fetched when Phase 1 for other authors times out`() =
|
||||
runBlocking {
|
||||
val client = ScriptedClient()
|
||||
val gateway = RecordingGateway()
|
||||
|
||||
// Alice has cached outbox — Phase 2 must fetch from her write relay.
|
||||
gateway.cache[alice] = outboxEventFor(alice, listOf(outboxAlice))
|
||||
client.scriptEvent(ContactListEvent.KIND, outboxAlice, listOf(kind3For(alice, listOf(bob))))
|
||||
|
||||
// Bob has no cached outbox and index relays never EOSE for Phase 1.
|
||||
indexRelays.forEach(client::neverEose)
|
||||
|
||||
val dispatcher =
|
||||
OutboxDispatcher(
|
||||
client = client,
|
||||
scope = scope,
|
||||
indexRelays = { indexRelays },
|
||||
gateway = gateway,
|
||||
perRelayTimeoutMs = 200,
|
||||
overallTimeoutMs = 2_000,
|
||||
)
|
||||
|
||||
val result = dispatcher.fetchKind3Only(setOf(alice, bob))
|
||||
|
||||
// Alice was covered by cached outbox; Bob wasn't but Phase 1 timed
|
||||
// out, so he became a fallback candidate.
|
||||
assertEquals(
|
||||
"Alice always covered by cached outbox",
|
||||
1,
|
||||
result.outboxCoveredAuthors,
|
||||
)
|
||||
assertTrue(result.kind3Received >= 1)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `clear releases dedup so a subsequent identical call refetches`() =
|
||||
runBlocking {
|
||||
val client = ScriptedClient()
|
||||
val gateway = RecordingGateway()
|
||||
gateway.cache[alice] = outboxEventFor(alice, listOf(outboxAlice))
|
||||
client.scriptEvent(ContactListEvent.KIND, outboxAlice, listOf(kind3For(alice, listOf(bob))))
|
||||
|
||||
val dispatcher =
|
||||
OutboxDispatcher(
|
||||
client = client,
|
||||
scope = scope,
|
||||
indexRelays = { indexRelays },
|
||||
gateway = gateway,
|
||||
perRelayTimeoutMs = 400,
|
||||
overallTimeoutMs = 2_000,
|
||||
)
|
||||
|
||||
dispatcher.fetchKind3Only(setOf(alice))
|
||||
val subCountAfterFirst = client.allSubscribeCalls.size
|
||||
|
||||
// Second call without clear() — should short-circuit.
|
||||
dispatcher.fetchKind3Only(setOf(alice))
|
||||
assertEquals(subCountAfterFirst, client.allSubscribeCalls.size)
|
||||
|
||||
// After clear(), the same call re-runs Phase 2.
|
||||
dispatcher.clear()
|
||||
dispatcher.fetchKind3Only(setOf(alice))
|
||||
assertTrue(client.allSubscribeCalls.size > subCountAfterFirst)
|
||||
}
|
||||
|
||||
/**
|
||||
* BatchEoseGate stress — inside OutboxDispatcher this is a private
|
||||
* class but the observable effect (Phase 1 completes when all index
|
||||
* relays EOSE, and stays within the timeout budget) is what matters.
|
||||
*/
|
||||
@Test
|
||||
fun `EOSE aggregation is safe with many concurrent index-relay callbacks`() =
|
||||
runBlocking {
|
||||
val bigIndexSet = (0..15).map { NormalizedRelayUrl("wss://index$it.test/") }.toSet()
|
||||
val client = ScriptedClient()
|
||||
val gateway = RecordingGateway()
|
||||
|
||||
val dispatcher =
|
||||
OutboxDispatcher(
|
||||
client = client,
|
||||
scope = scope,
|
||||
indexRelays = { bigIndexSet },
|
||||
gateway = gateway,
|
||||
perRelayTimeoutMs = 1_000,
|
||||
overallTimeoutMs = 3_000,
|
||||
)
|
||||
|
||||
// Kick off a fetch and race the subscribe call. ScriptedClient
|
||||
// fires EOSE inline; we simulate concurrent per-relay EOSE by
|
||||
// launching multiple dispatchers as a smoke test.
|
||||
val fetchJob = scope.launch { dispatcher.fetchKind3Only(setOf(alice, bob, charlie)) }
|
||||
|
||||
// Give the launcher a moment to enter Phase 1's subscribe.
|
||||
delay(50)
|
||||
fetchJob.join()
|
||||
// No CME thrown, no hang past the timeout budget.
|
||||
}
|
||||
}
|
||||
Vendored
+33
@@ -56,6 +56,7 @@ import com.vitorpamplona.quartz.nip51Lists.bookmarkList.OldBookmarkListEvent
|
||||
import com.vitorpamplona.quartz.nip51Lists.followList.FollowListEvent
|
||||
import com.vitorpamplona.quartz.nip57Zaps.LnZapEvent
|
||||
import com.vitorpamplona.quartz.nip57Zaps.LnZapRequestEvent
|
||||
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
import com.vitorpamplona.quartz.utils.DualCase
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import kotlinx.coroutines.CancellationException
|
||||
@@ -303,11 +304,43 @@ class DesktopLocalCache : ICacheProvider {
|
||||
consumeComment(event, relay)
|
||||
}
|
||||
|
||||
is AdvertisedRelayListEvent -> {
|
||||
consumeAdvertisedRelayList(event, relay)
|
||||
}
|
||||
|
||||
else -> {
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Consumes a kind 10002 (NIP-65) advertised relay list event. Stores
|
||||
* the newest per-author copy in [addressableNotes] so the outbox
|
||||
* dispatcher can look up each follow's declared write relays without
|
||||
* a fresh REQ. Emits nothing to the event stream — the UI doesn't
|
||||
* render kind 10002s directly.
|
||||
*/
|
||||
private fun consumeAdvertisedRelayList(
|
||||
event: AdvertisedRelayListEvent,
|
||||
relay: NormalizedRelayUrl?,
|
||||
): Boolean {
|
||||
val addressableNote = getOrCreateAddressableNote(event.address())
|
||||
val existing = addressableNote.event
|
||||
if (existing != null && existing.createdAt >= event.createdAt) return false
|
||||
val author = getOrCreateUser(event.pubKey)
|
||||
addressableNote.loadEvent(event, author, emptyList())
|
||||
relay?.let { addressableNote.addRelay(it) }
|
||||
return false
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the cached kind-10002 event for [pubkey], if any. Used by the
|
||||
* outbox dispatcher to skip a Phase-1 REQ for authors whose write-relay
|
||||
* list is already in the store (from a previous session's local relay
|
||||
* hydration or an in-session discovery).
|
||||
*/
|
||||
fun cachedAdvertisedRelayList(pubkey: HexKey): AdvertisedRelayListEvent? = addressableNotes.get(AdvertisedRelayListEvent.createAddress(pubkey).toValue())?.event as? AdvertisedRelayListEvent
|
||||
|
||||
/**
|
||||
* Consumes a kind 1 text note event.
|
||||
* Creates/updates Note in cache and links reply relationships.
|
||||
|
||||
Reference in New Issue
Block a user