mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
refactor: use quartz INostrClient.fetchAll instead of a bespoke Context.fetchRaw
fetchAll already does exactly what the need-metadata fetch needs — subscribe, collect (deduped by id), return on EOSE/timeout, no verify, no store — so drop the duplicated Context.fetchRaw and call the existing extension. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JgL1WTV4Hkp2uuXcUHCHGt
This commit is contained in:
@@ -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<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 8_000,
|
||||
): List<Event> {
|
||||
if (filters.isEmpty()) return emptyList()
|
||||
val eventChannel = Channel<Event>(UNLIMITED)
|
||||
val doneChannel = Channel<NormalizedRelayUrl>(UNLIMITED)
|
||||
val remaining = filters.keys.toMutableSet()
|
||||
val subId = newSubId()
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
eventChannel.trySend(event)
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay)
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
message: String,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay)
|
||||
}
|
||||
|
||||
override fun onCannotConnect(
|
||||
relay: NormalizedRelayUrl,
|
||||
message: String,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay)
|
||||
}
|
||||
}
|
||||
val collected = mutableListOf<Event>()
|
||||
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
|
||||
|
||||
@@ -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)) {
|
||||
|
||||
Reference in New Issue
Block a user