diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt index 928dd3a7ae..dd0bac444e 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -483,79 +483,6 @@ class Context( return collected } - /** - * Like [drain] but does NOT verify or store — collects the raw events until every - * relay EOSEs or the timeout elapses and returns them. Used when the caller needs - * an event's metadata to make a decision (e.g. which local deletion covers it) - * rather than to keep it. Verification/storage, if wanted, is the caller's job. - */ - suspend fun fetchRaw( - filters: Map>, - timeoutMs: Long = 8_000, - ): List { - if (filters.isEmpty()) return emptyList() - val eventChannel = Channel(UNLIMITED) - val doneChannel = Channel(UNLIMITED) - val remaining = filters.keys.toMutableSet() - val subId = newSubId() - val listener = - object : SubscriptionListener { - override fun onEvent( - event: Event, - isLive: Boolean, - relay: NormalizedRelayUrl, - forFilters: List?, - ) { - eventChannel.trySend(event) - } - - override fun onEose( - relay: NormalizedRelayUrl, - forFilters: List?, - ) { - doneChannel.trySend(relay) - } - - override fun onClosed( - message: String, - relay: NormalizedRelayUrl, - forFilters: List?, - ) { - doneChannel.trySend(relay) - } - - override fun onCannotConnect( - relay: NormalizedRelayUrl, - message: String, - forFilters: List?, - ) { - doneChannel.trySend(relay) - } - } - val collected = mutableListOf() - try { - client.subscribe(subId, filters, listener) - withTimeoutOrNull(timeoutMs) { - while (remaining.isNotEmpty()) { - select { - eventChannel.onReceive { collected.add(it) } - doneChannel.onReceive { r -> remaining.remove(r) } - } - } - while (true) { - val r = eventChannel.tryReceive() - if (!r.isSuccess) break - collected.add(r.getOrThrow()) - } - } - } finally { - client.unsubscribe(subId) - eventChannel.close() - doneChannel.close() - } - return collected - } - /** * Like [drain], but paginates every relay to completion via * [fetchAllPagesFromPool] instead of stopping at the first EOSE — so a query diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt index c1e8294fc0..333c8a2677 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/SyncCommand.kt @@ -27,6 +27,7 @@ import com.vitorpamplona.amethyst.cli.Output import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.NegentropySyncException +import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAll import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.negentropyReconcile import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer @@ -110,7 +111,7 @@ object SyncCommand { // publish the local deletions that would make the relay remove them — an id- or // address-based kind-5, or a kind-62 vanish that targets this relay. Only those // deletions, nothing else (not other deletions by the same author). We fetch the - // need events (not to keep — [Context.fetchRaw] neither verifies nor stores) only + // need events (not to keep — `fetchAll` neither verifies nor stores) only // to learn their author/address/created_at so [deletionsCovering] can tell which // of our deletions actually apply. Nothing is pulled down or applied locally, so // this can never over-delete the local store. @@ -145,8 +146,9 @@ object SyncCommand { List(DOWNLOAD_WORKERS) { launch { for (batch in needBatches) { - // Fetch the need events once (raw — no verify/store). - val events = ctx.fetchRaw(mapOf(relay to listOf(Filter(ids = batch))), timeoutMs) + // Fetch the need events once (no verify/store — we only + // need their metadata to decide which deletions apply). + val events = ctx.client.fetchAll(relay, Filter(ids = batch), timeoutMs) // Push up the deletions that would remove them from the relay. if (syncDeletions) { for (del in ctx.store.deletionsCovering(events, relay)) {