diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxCacheGateway.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxCacheGateway.kt new file mode 100644 index 0000000000..a9c451f05f --- /dev/null +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxCacheGateway.kt @@ -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, + ) +} diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxDispatcher.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxDispatcher.kt new file mode 100644 index 0000000000..8dbd7308b7 --- /dev/null +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxDispatcher.kt @@ -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, + 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() + + /** + * Pubkeys we've already successfully fetched kind-0 for this session. + */ + private val kind0Succeeded = mutableSetOf() + + /** + * 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() + private val kind0InFlight = mutableSetOf() + + /** + * 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): 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): 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): 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, + 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, + newForKind3: Set, + newForKind0: Set, + 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>() + val toDiscover = mutableSetOf() + 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>() + 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, + relays: Set, + ): Pair>>, 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>>() + val gate = BatchEoseGate(scope, target = relays.size) + + val listener = + object : SubscriptionListener { + override fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + if (event is AdvertisedRelayListEvent && event.pubKey in pubkeys) { + received + .getOrPut(event.pubKey) { mutableListOf() } + .add(event to relay) + } + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + 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, + kinds: List, + authors: Set, + 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?, + ) { + when (event.kind) { + MetadataEvent.KIND -> counters.kind0++ + ContactListEvent.KIND -> counters.kind3++ + } + gateway.onDiscoveredEvent(event, relay) + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + 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(Channel.UNLIMITED) + private val done = CompletableDeferred() + + @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() + 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 + } + } +} diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxDispatcherTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxDispatcherTest.kt new file mode 100644 index 0000000000..9238b5ce1d --- /dev/null +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/wot/OutboxDispatcherTest.kt @@ -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, + 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, + ) = 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() + val discoveredOutbox = mutableListOf>() + val discoveredEvents = mutableListOf>() + + 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, List>() + private val eoseNever = mutableSetOf() + val allSubscribeCalls = mutableListOf>>() + + fun scriptEvent( + kind: Int, + relay: NormalizedRelayUrl, + events: List, + ) { + script[kind to relay] = events + } + + fun neverEose(relay: NormalizedRelayUrl) { + eoseNever.add(relay) + } + + override fun subscribe( + subId: String, + filters: Map>, + 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. + } +} diff --git a/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/cache/DesktopLocalCache.kt b/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/cache/DesktopLocalCache.kt index cd63af0d83..2212833b67 100644 --- a/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/cache/DesktopLocalCache.kt +++ b/desktopApp/src/jvmMain/kotlin/com/vitorpamplona/amethyst/desktop/cache/DesktopLocalCache.kt @@ -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.