mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 19:53:08 +00:00
refactor(cli): extract the drain loop to quartz accessories; split Context per domain
Context.drain/drainAllPages/requestResponse re-implemented the subscription state machine quartz already ships in relay/client/accessories — the CLI-specific needs (per-event verify-and-store hook, dead-relay collection, pending-on-auth) now live in an option-rich fetchAll variant there, and Context keeps thin adapters. The per-domain sections bolted onto Context (Cashu seed warming/snapshot/restore counters; the Concord stream-key AUTH registry) move to CashuContext/ConcordAuth, with the NUT-09 restore counter rule shared via commons CashuWalletOps so the CLI and Android can't drift. Context.kt: 1246 -> 950 lines; behavior unchanged (cli tests green). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CP4kfLCa3wWtE8Khy21Pkj
This commit is contained in:
@@ -0,0 +1,152 @@
|
||||
/*
|
||||
* 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.cli
|
||||
|
||||
import com.vitorpamplona.amethyst.cli.stores.FileCashuKeysetCounterStore
|
||||
import com.vitorpamplona.amethyst.commons.cashu.CashuWalletReader
|
||||
import com.vitorpamplona.amethyst.commons.cashu.ops.CashuWalletOps
|
||||
import com.vitorpamplona.amethyst.commons.cashu.ops.RestoreOutcome
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip60Cashu.history.CashuSpendingHistoryEvent
|
||||
import com.vitorpamplona.quartz.nip60Cashu.mintApi.DeterministicSecretFactory
|
||||
import com.vitorpamplona.quartz.nip60Cashu.quote.CashuMintQuoteEvent
|
||||
import com.vitorpamplona.quartz.nip60Cashu.seed.CashuDeterministic
|
||||
import com.vitorpamplona.quartz.nip60Cashu.token.CashuTokenEvent
|
||||
import com.vitorpamplona.quartz.nip60Cashu.wallet.CashuWalletEvent
|
||||
import com.vitorpamplona.quartz.nip61Nutzaps.info.NutzapInfoEvent
|
||||
import com.vitorpamplona.quartz.nip61Nutzaps.nutzap.NutzapEvent
|
||||
import com.vitorpamplona.quartz.nip87Ecash.recommendation.MintRecommendationEvent
|
||||
|
||||
/**
|
||||
* Cashu (NIP-60 / NIP-61) wiring for the CLI, split out of [Context] — shared
|
||||
* wallet code from commons, driven by the CLI's file-backed NUT-13 counter
|
||||
* store and the run's [Context.store] snapshot projection. Instantiated
|
||||
* lazily by [Context.cashu], so a run that never touches the wallet pays
|
||||
* nothing.
|
||||
*/
|
||||
class CashuContext(
|
||||
private val ctx: Context,
|
||||
) {
|
||||
/** Durable NUT-13 counter store at `<data-dir>/cashu.json`. */
|
||||
private val counters by lazy { FileCashuKeysetCounterStore(ctx.dataDir.cashuFile) }
|
||||
|
||||
@Volatile private var cachedSeed: ByteArray? = null
|
||||
|
||||
/**
|
||||
* Decrypt the wallet's NUT-13 seed once per run and cache it. The
|
||||
* [DeterministicSecretFactory] thunk reads this synchronously, so any
|
||||
* mint/swap op must warm it first (CashuWalletOps' seedWarmer does).
|
||||
*/
|
||||
private suspend fun warmSeed() {
|
||||
if (cachedSeed != null) return
|
||||
val priv = snapshot().walletEvent?.let { runCatching { it.privkey(ctx.signer) }.getOrNull() } ?: return
|
||||
cachedSeed = CashuDeterministic.deriveWalletSeed(priv.hexToByteArray())
|
||||
}
|
||||
|
||||
/**
|
||||
* Wallet operations driven by the exact same `commons` [CashuWalletOps]
|
||||
* the Android app uses. Wired to publish on the account's outbox relays,
|
||||
* the shared OkHttp instance for mint HTTP, and the file-backed NUT-13
|
||||
* counter store.
|
||||
*/
|
||||
fun ops(): CashuWalletOps =
|
||||
CashuWalletOps(
|
||||
signer = ctx.signer,
|
||||
// Amethyst publishes cashu events via sendLiterallyEverywhere
|
||||
// (all the user's relays). anyRelays() — outbox + inbox +
|
||||
// keypackage — is the CLI's closest analog, so the wallet lands
|
||||
// on the same broad relay set the app would use, not just outbox.
|
||||
publish = { event -> ctx.publish(event, ctx.anyRelays()) },
|
||||
okHttpClient = { ctx.okhttp },
|
||||
secretFactory =
|
||||
DeterministicSecretFactory(
|
||||
seedProvider = { cachedSeed },
|
||||
reserveCounters = { keysetId, count -> counters.reserve(keysetId, count) },
|
||||
),
|
||||
seedWarmer = { warmSeed() },
|
||||
seedForRestore = {
|
||||
warmSeed()
|
||||
cachedSeed
|
||||
},
|
||||
peekCashuCounter = { keysetId -> counters.peek(keysetId) },
|
||||
reserveCashuCounters = { keysetId, count -> counters.reserve(keysetId, count) },
|
||||
)
|
||||
|
||||
/**
|
||||
* Project this account's locally-stored NIP-60/61/87 events into a wallet
|
||||
* snapshot via the shared [CashuWalletReader] — the same decrypt +
|
||||
* del-rollover + pending-quote logic the Android holder runs. Reads the
|
||||
* cache only; commands that need fresh state should [Context.drain] first.
|
||||
*/
|
||||
suspend fun snapshot(): CashuWalletReader.WalletSnapshot {
|
||||
val pk = ctx.identity.pubKeyHex
|
||||
// Mirror commons' CashuWalletFilterAssembler exactly: authored wallet
|
||||
// kinds by authors=[pk], inbound nutzaps by #p — so amy projects the
|
||||
// same event set the Android app subscribes to.
|
||||
val authored =
|
||||
ctx.store.query<Event>(
|
||||
Filter(
|
||||
authors = listOf(pk),
|
||||
kinds =
|
||||
listOf(
|
||||
CashuWalletEvent.KIND,
|
||||
CashuTokenEvent.KIND,
|
||||
CashuSpendingHistoryEvent.KIND,
|
||||
CashuMintQuoteEvent.KIND,
|
||||
NutzapInfoEvent.KIND,
|
||||
MintRecommendationEvent.KIND,
|
||||
),
|
||||
),
|
||||
)
|
||||
val inboundNutzaps =
|
||||
ctx.store.query<Event>(
|
||||
Filter(kinds = listOf(NutzapEvent.KIND), tags = mapOf("p" to listOf(pk))),
|
||||
)
|
||||
return CashuWalletReader(ctx.signer).project(authored + inboundNutzaps)
|
||||
}
|
||||
|
||||
/** Warm and return the wallet's NUT-13 seed, or null if no wallet. */
|
||||
suspend fun seed(): ByteArray? {
|
||||
warmSeed()
|
||||
return cachedSeed
|
||||
}
|
||||
|
||||
/**
|
||||
* NUT-09 restore for one mint, mirroring Android's
|
||||
* `CashuWalletState.restoreFromMint`: re-derive proofs from the seed
|
||||
* (skipping secrets we already hold), then bump the persisted NUT-13
|
||||
* counter past every slot the scan confirmed in use so later mints can't
|
||||
* reuse one — the advance rule is the shared
|
||||
* [CashuWalletOps.restoreFromMintAdvancingCounter]. Returns null when the
|
||||
* wallet has no seed yet.
|
||||
*/
|
||||
suspend fun restore(mintUrl: String): RestoreOutcome? {
|
||||
val seed = seed() ?: return null
|
||||
val existingSecrets =
|
||||
snapshot()
|
||||
.tokenEntries
|
||||
.flatMap { it.content.proofs }
|
||||
.mapTo(HashSet()) { it.secret }
|
||||
return ops().restoreFromMintAdvancingCounter(mintUrl = mintUrl, seed = seed, startCounter = 0L, existingSecrets = existingSecrets)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,69 @@
|
||||
/*
|
||||
* 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.cli
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.EventTemplate
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
* Concord plane stream-key AUTH registry (CORD-01 §4b). Concord relays gate a
|
||||
* plane's kind-1059 wraps behind NIP-42 and serve them only to a connection
|
||||
* authenticated AS the plane's derived *stream key* — the member is neither the
|
||||
* wrap's author (the stream key) nor its recipient (a throwaway ephemeral key),
|
||||
* so an account AUTH is refused. `amy concord` verbs register their control +
|
||||
* channel stream secrets here (scoped to the community's relays, the same scope
|
||||
* the plane REQ uses) before draining, and the [Context]'s NIP-42 responder
|
||||
* answers a challenge from one of those relays with one kind-22242 per stream
|
||||
* key — signed locally from the raw derived key, never the account, so no user
|
||||
* identity is exposed.
|
||||
*/
|
||||
class ConcordAuth {
|
||||
private val streamSecrets = ConcurrentHashMap<NormalizedRelayUrl, MutableSet<HexKey>>()
|
||||
private val streamSigners = ConcurrentHashMap<HexKey, NostrSignerSync>()
|
||||
|
||||
/** Registers raw 32-byte Concord stream [secrets] to answer NIP-42 challenges from [relays]. */
|
||||
fun register(
|
||||
relays: Set<NormalizedRelayUrl>,
|
||||
secrets: List<ByteArray>,
|
||||
) {
|
||||
if (relays.isEmpty() || secrets.isEmpty()) return
|
||||
val hexes = secrets.map { it.toHexKey() }
|
||||
for (relay in relays) streamSecrets.getOrPut(relay) { ConcurrentHashMap.newKeySet() }.addAll(hexes)
|
||||
}
|
||||
|
||||
/** Signs one kind-22242 per Concord stream key registered for [relay] (empty if none). */
|
||||
fun signAuths(
|
||||
relay: NormalizedRelayUrl,
|
||||
template: EventTemplate<RelayAuthEvent>,
|
||||
): List<RelayAuthEvent> =
|
||||
streamSecrets[relay].orEmpty().mapNotNull { hex ->
|
||||
runCatching {
|
||||
streamSigners.getOrPut(hex) { NostrSignerSync(KeyPair(privKey = hex.hexToByteArray())) }.sign(template)
|
||||
}.getOrNull()
|
||||
}
|
||||
}
|
||||
@@ -21,7 +21,6 @@
|
||||
package com.vitorpamplona.amethyst.cli
|
||||
|
||||
import com.sun.management.UnixOperatingSystemMXBean
|
||||
import com.vitorpamplona.amethyst.cli.stores.FileCashuKeysetCounterStore
|
||||
import com.vitorpamplona.amethyst.cli.stores.FileKeyPackageBundleStore
|
||||
import com.vitorpamplona.amethyst.cli.stores.FileMarmotMessageStore
|
||||
import com.vitorpamplona.amethyst.cli.stores.FileMlsGroupStateStore
|
||||
@@ -36,21 +35,17 @@ import com.vitorpamplona.quartz.marmot.RecipientRelayFetcher
|
||||
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageRelayListEvent
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.AdaptiveRelayLimiter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.DrainFailure
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.classifyDrainFailure
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPagesFromPool
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllPagesFromPoolWithHooks
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.fetchAllWithHooks
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.publishAndConfirmDetailed
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.auth.RelayAuthenticator
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.SubscriptionListener
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.single.newSubId
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CachingEventDecoder
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.MachineReadablePrefix
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
@@ -59,44 +54,25 @@ import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSoc
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDnsStore
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.EventTemplate
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
|
||||
import com.vitorpamplona.quartz.nip01Core.store.IEventStore
|
||||
import com.vitorpamplona.quartz.nip01Core.store.verifyAndInsert
|
||||
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
|
||||
import com.vitorpamplona.quartz.nip05DnsIdentifiers.resolveUserHexOrNull
|
||||
import com.vitorpamplona.quartz.nip11RelayInfo.Nip11RelayInformation
|
||||
import com.vitorpamplona.quartz.nip17Dm.settings.ChatMessageRelayListEvent
|
||||
import com.vitorpamplona.quartz.nip42RelayAuth.RelayAuthEvent
|
||||
import com.vitorpamplona.quartz.nip46RemoteSigner.signer.NostrSignerRemote
|
||||
import com.vitorpamplona.quartz.nip60Cashu.history.CashuSpendingHistoryEvent
|
||||
import com.vitorpamplona.quartz.nip60Cashu.mintApi.DeterministicSecretFactory
|
||||
import com.vitorpamplona.quartz.nip60Cashu.quote.CashuMintQuoteEvent
|
||||
import com.vitorpamplona.quartz.nip60Cashu.seed.CashuDeterministic
|
||||
import com.vitorpamplona.quartz.nip60Cashu.token.CashuTokenEvent
|
||||
import com.vitorpamplona.quartz.nip60Cashu.wallet.CashuWalletEvent
|
||||
import com.vitorpamplona.quartz.nip61Nutzaps.info.NutzapInfoEvent
|
||||
import com.vitorpamplona.quartz.nip61Nutzaps.nutzap.NutzapEvent
|
||||
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
|
||||
import com.vitorpamplona.quartz.nip66RelayMonitor.reachability.RelayReachabilityStore
|
||||
import com.vitorpamplona.quartz.nip87Ecash.recommendation.MintRecommendationEvent
|
||||
import com.vitorpamplona.quartz.utils.SeenIds
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.selects.select
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
import okhttp3.Dispatcher
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import java.lang.management.ManagementFactory
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
/**
|
||||
@@ -156,7 +132,9 @@ class Context(
|
||||
runCatching { it.load() }
|
||||
}
|
||||
|
||||
private val okhttp =
|
||||
// Internal (not private) so [CashuContext] can reuse the shared instance
|
||||
// for mint HTTP.
|
||||
internal val okhttp =
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
.socketFactory(TcpNoDelaySocketFactory)
|
||||
@@ -257,39 +235,18 @@ class Context(
|
||||
).also { client.addConnectionListener(it) }
|
||||
|
||||
/**
|
||||
* Concord plane stream-key AUTH (CORD-01 §4b). Concord relays gate a plane's
|
||||
* kind-1059 wraps behind NIP-42 and serve them only to a connection authenticated
|
||||
* AS the plane's derived *stream key* — the member is neither the wrap's author
|
||||
* (the stream key) nor its recipient (a throwaway ephemeral key), so an account
|
||||
* AUTH is refused. `amy concord` verbs register their control + channel stream
|
||||
* secrets here (scoped to the community's relays, the same scope the plane REQ
|
||||
* uses) before draining, and [relayAuth] answers a challenge from one of those
|
||||
* relays with one kind-22242 per stream key — signed locally from the raw derived
|
||||
* key, never the account, so no user identity is exposed.
|
||||
* Concord plane stream-key AUTH registry (CORD-01 §4b) — see [ConcordAuth].
|
||||
* `amy concord` verbs register their stream secrets via
|
||||
* [registerConcordStreamKeys] before draining, and [relayAuth] answers a
|
||||
* challenge from one of those relays with one kind-22242 per stream key.
|
||||
*/
|
||||
private val concordStreamSecrets = ConcurrentHashMap<NormalizedRelayUrl, MutableSet<HexKey>>()
|
||||
private val concordStreamSigners = ConcurrentHashMap<HexKey, NostrSignerSync>()
|
||||
private val concordAuth = ConcordAuth()
|
||||
|
||||
/** Registers raw 32-byte Concord stream [secrets] to answer NIP-42 challenges from [relays]. */
|
||||
fun registerConcordStreamKeys(
|
||||
relays: Set<NormalizedRelayUrl>,
|
||||
secrets: List<ByteArray>,
|
||||
) {
|
||||
if (relays.isEmpty() || secrets.isEmpty()) return
|
||||
val hexes = secrets.map { it.toHexKey() }
|
||||
for (relay in relays) concordStreamSecrets.getOrPut(relay) { ConcurrentHashMap.newKeySet() }.addAll(hexes)
|
||||
}
|
||||
|
||||
/** Signs one kind-22242 per Concord stream key registered for [relay] (empty if none). */
|
||||
private fun signConcordStreamAuths(
|
||||
relay: NormalizedRelayUrl,
|
||||
template: EventTemplate<RelayAuthEvent>,
|
||||
): List<RelayAuthEvent> =
|
||||
concordStreamSecrets[relay].orEmpty().mapNotNull { hex ->
|
||||
runCatching {
|
||||
concordStreamSigners.getOrPut(hex) { NostrSignerSync(KeyPair(privKey = hex.hexToByteArray())) }.sign(template)
|
||||
}.getOrNull()
|
||||
}
|
||||
) = concordAuth.register(relays, secrets)
|
||||
|
||||
/**
|
||||
* NIP-42 responder: answers a relay's AUTH challenge by signing with the
|
||||
@@ -311,7 +268,7 @@ class Context(
|
||||
} else {
|
||||
emptyList()
|
||||
}
|
||||
accountAuth + signConcordStreamAuths(relay, template)
|
||||
accountAuth + concordAuth.signAuths(relay, template)
|
||||
},
|
||||
)
|
||||
|
||||
@@ -386,109 +343,25 @@ class Context(
|
||||
// Cashu (NIP-60 / NIP-61) — shared wallet code from commons
|
||||
// ------------------------------------------------------------------
|
||||
|
||||
/** Durable NUT-13 counter store at `<data-dir>/cashu.json`. */
|
||||
private val cashuCounters by lazy { FileCashuKeysetCounterStore(dataDir.cashuFile) }
|
||||
|
||||
@Volatile private var cachedCashuSeed: ByteArray? = null
|
||||
|
||||
/**
|
||||
* Decrypt the wallet's NUT-13 seed once per run and cache it. The
|
||||
* [DeterministicSecretFactory] thunk reads this synchronously, so any
|
||||
* mint/swap op must warm it first (CashuWalletOps' seedWarmer does).
|
||||
* Cashu wiring (seed warming, snapshot projection, NUT-09 restore) —
|
||||
* see [CashuContext]. Lazy so a run that never touches the wallet pays
|
||||
* nothing. The `cashuOps` / `cashuSnapshot` / `cashuSeed` / `cashuRestore`
|
||||
* members below forward to it, keeping the command surface unchanged.
|
||||
*/
|
||||
private suspend fun warmCashuSeed() {
|
||||
if (cachedCashuSeed != null) return
|
||||
val priv = cashuSnapshot().walletEvent?.let { runCatching { it.privkey(signer) }.getOrNull() } ?: return
|
||||
cachedCashuSeed = CashuDeterministic.deriveWalletSeed(priv.hexToByteArray())
|
||||
}
|
||||
val cashu: CashuContext by lazy { CashuContext(this) }
|
||||
|
||||
/**
|
||||
* Wallet operations driven by the exact same `commons` [CashuWalletOps]
|
||||
* the Android app uses. Wired to publish on the account's outbox relays,
|
||||
* the shared OkHttp instance for mint HTTP, and the file-backed NUT-13
|
||||
* counter store.
|
||||
*/
|
||||
fun cashuOps(): CashuWalletOps =
|
||||
CashuWalletOps(
|
||||
signer = signer,
|
||||
// Amethyst publishes cashu events via sendLiterallyEverywhere
|
||||
// (all the user's relays). anyRelays() — outbox + inbox +
|
||||
// keypackage — is the CLI's closest analog, so the wallet lands
|
||||
// on the same broad relay set the app would use, not just outbox.
|
||||
publish = { event -> publish(event, anyRelays()) },
|
||||
okHttpClient = { okhttp },
|
||||
secretFactory =
|
||||
DeterministicSecretFactory(
|
||||
seedProvider = { cachedCashuSeed },
|
||||
reserveCounters = { keysetId, count -> cashuCounters.reserve(keysetId, count) },
|
||||
),
|
||||
seedWarmer = { warmCashuSeed() },
|
||||
seedForRestore = {
|
||||
warmCashuSeed()
|
||||
cachedCashuSeed
|
||||
},
|
||||
peekCashuCounter = { keysetId -> cashuCounters.peek(keysetId) },
|
||||
reserveCashuCounters = { keysetId, count -> cashuCounters.reserve(keysetId, count) },
|
||||
)
|
||||
/** See [CashuContext.ops]. */
|
||||
fun cashuOps(): CashuWalletOps = cashu.ops()
|
||||
|
||||
/**
|
||||
* Project this account's locally-stored NIP-60/61/87 events into a wallet
|
||||
* snapshot via the shared [CashuWalletReader] — the same decrypt +
|
||||
* del-rollover + pending-quote logic the Android holder runs. Reads the
|
||||
* cache only; commands that need fresh state should [drain] first.
|
||||
*/
|
||||
suspend fun cashuSnapshot(): CashuWalletReader.WalletSnapshot {
|
||||
val pk = identity.pubKeyHex
|
||||
// Mirror commons' CashuWalletFilterAssembler exactly: authored wallet
|
||||
// kinds by authors=[pk], inbound nutzaps by #p — so amy projects the
|
||||
// same event set the Android app subscribes to.
|
||||
val authored =
|
||||
store.query<Event>(
|
||||
Filter(
|
||||
authors = listOf(pk),
|
||||
kinds =
|
||||
listOf(
|
||||
CashuWalletEvent.KIND,
|
||||
CashuTokenEvent.KIND,
|
||||
CashuSpendingHistoryEvent.KIND,
|
||||
CashuMintQuoteEvent.KIND,
|
||||
NutzapInfoEvent.KIND,
|
||||
MintRecommendationEvent.KIND,
|
||||
),
|
||||
),
|
||||
)
|
||||
val inboundNutzaps =
|
||||
store.query<Event>(
|
||||
Filter(kinds = listOf(NutzapEvent.KIND), tags = mapOf("p" to listOf(pk))),
|
||||
)
|
||||
return CashuWalletReader(signer).project(authored + inboundNutzaps)
|
||||
}
|
||||
/** See [CashuContext.snapshot]. */
|
||||
suspend fun cashuSnapshot(): CashuWalletReader.WalletSnapshot = cashu.snapshot()
|
||||
|
||||
/** Warm and return the wallet's NUT-13 seed, or null if no wallet. */
|
||||
suspend fun cashuSeed(): ByteArray? {
|
||||
warmCashuSeed()
|
||||
return cachedCashuSeed
|
||||
}
|
||||
/** See [CashuContext.seed]. */
|
||||
suspend fun cashuSeed(): ByteArray? = cashu.seed()
|
||||
|
||||
/**
|
||||
* NUT-09 restore for one mint, mirroring Android's
|
||||
* `CashuWalletState.restoreFromMint`: re-derive proofs from the seed
|
||||
* (skipping secrets we already hold), then bump the persisted NUT-13
|
||||
* counter past every slot the scan confirmed in use so later mints can't
|
||||
* reuse one. Returns null when the wallet has no seed yet.
|
||||
*/
|
||||
suspend fun cashuRestore(mintUrl: String): RestoreOutcome? {
|
||||
val seed = cashuSeed() ?: return null
|
||||
val existingSecrets =
|
||||
cashuSnapshot()
|
||||
.tokenEntries
|
||||
.flatMap { it.content.proofs }
|
||||
.mapTo(HashSet()) { it.secret }
|
||||
val outcome = cashuOps().restoreFromMint(mintUrl = mintUrl, seed = seed, startCounter = 0L, existingSecrets = existingSecrets)
|
||||
val delta = (outcome.nextCounterAfterScan - cashuCounters.peek(outcome.keysetId)).coerceAtLeast(0L)
|
||||
if (delta > 0) cashuCounters.reserve(outcome.keysetId, delta.toInt())
|
||||
return outcome
|
||||
}
|
||||
/** See [CashuContext.restore]. */
|
||||
suspend fun cashuRestore(mintUrl: String): RestoreOutcome? = cashu.restore(mintUrl)
|
||||
|
||||
private var prepared = false
|
||||
|
||||
@@ -640,6 +513,8 @@ class Context(
|
||||
* Subscribe to the given filters across the given relays, drain all events
|
||||
* until either every relay has sent EOSE or the timeout elapses, and
|
||||
* return them. Used for one-shot catch-up queries — not live subscriptions.
|
||||
* Thin adapter over the shared [fetchAllWithHooks] accessory: every arriving
|
||||
* event is verified + persisted via [verifyAndStore] before it is surfaced.
|
||||
*
|
||||
* When [deadOut] is provided, every relay that reported it could not be
|
||||
* connected to (`onCannotConnect`) is added to it, so callers can prune
|
||||
@@ -661,92 +536,19 @@ class Context(
|
||||
diagnoseSlow: Boolean = false,
|
||||
deadOut: MutableMap<NormalizedRelayUrl, DrainFailure>? = null,
|
||||
pendingOnAuthRequired: Boolean = false,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
if (filters.isEmpty()) return emptyList()
|
||||
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(UNLIMITED)
|
||||
// Carries the terminal reason per relay so a timeout can distinguish a slow
|
||||
// relay (never terminal) from a connect failure / CLOSED.
|
||||
val doneChannel = Channel<Pair<NormalizedRelayUrl, String>>(UNLIMITED)
|
||||
val remaining = filters.keys.toMutableSet()
|
||||
val doneReasons = HashMap<NormalizedRelayUrl, String>()
|
||||
val subId = newSubId()
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
eventChannel.trySend(relay to event)
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay to "eose")
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
message: String,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
// Keep the relay pending on an auth-required refusal: the authenticator answers the
|
||||
// challenge and re-fires this subscription, so the post-auth events still arrive.
|
||||
if (pendingOnAuthRequired && MachineReadablePrefix.parse(message) == MachineReadablePrefix.AUTH_REQUIRED) return
|
||||
doneChannel.trySend(relay to "closed:$message")
|
||||
}
|
||||
|
||||
override fun onCannotConnect(
|
||||
relay: NormalizedRelayUrl,
|
||||
message: String,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay to "cannot:$message")
|
||||
}
|
||||
}
|
||||
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
|
||||
try {
|
||||
client.subscribe(subId, filters, listener)
|
||||
val completed =
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
while (remaining.isNotEmpty()) {
|
||||
select {
|
||||
eventChannel.onReceive { pair ->
|
||||
if (verifyAndStore(pair.second)) collected.add(pair)
|
||||
}
|
||||
doneChannel.onReceive { (relay, reason) ->
|
||||
remaining.remove(relay)
|
||||
doneReasons[relay] = reason
|
||||
}
|
||||
}
|
||||
}
|
||||
// Drain any events that landed after EOSE but before cancel
|
||||
while (true) {
|
||||
val r = eventChannel.tryReceive()
|
||||
if (!r.isSuccess) break
|
||||
val pair = r.getOrThrow()
|
||||
if (verifyAndStore(pair.second)) collected.add(pair)
|
||||
}
|
||||
true
|
||||
}
|
||||
if (diagnoseSlow && completed == null && remaining.isNotEmpty()) {
|
||||
logSlowDrain(timeoutMs, remaining, doneReasons, collected)
|
||||
}
|
||||
} finally {
|
||||
client.unsubscribe(subId)
|
||||
eventChannel.close()
|
||||
doneChannel.close()
|
||||
}
|
||||
deadOut?.let { out ->
|
||||
for ((relay, reason) in doneReasons) {
|
||||
classifyDrainFailure(reason)?.let { out[relay] = it }
|
||||
}
|
||||
}
|
||||
return collected
|
||||
}
|
||||
): List<Pair<NormalizedRelayUrl, Event>> =
|
||||
client.fetchAllWithHooks(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
pendingOnAuthRequired = pendingOnAuthRequired,
|
||||
deadOut = deadOut,
|
||||
onTimeout =
|
||||
if (diagnoseSlow) {
|
||||
{ stalled, doneReasons, collected -> logSlowDrain(timeoutMs, stalled, doneReasons, collected) }
|
||||
} else {
|
||||
null
|
||||
},
|
||||
) { _, event -> verifyAndStore(event) }
|
||||
|
||||
/**
|
||||
* On a [drain] timeout, report which relays stalled and why — a relay that
|
||||
@@ -775,16 +577,16 @@ class Context(
|
||||
|
||||
/**
|
||||
* Like [drain], but paginates every relay to completion via
|
||||
* [fetchAllPagesFromPool] instead of stopping at the first EOSE — so a query
|
||||
* larger than a relay's per-`REQ` cap (strfry's `limit`, ~500) is fully
|
||||
* [fetchAllPagesFromPoolWithHooks] instead of stopping at the first EOSE — so a
|
||||
* query larger than a relay's per-`REQ` cap (strfry's `limit`, ~500) is fully
|
||||
* retrieved instead of silently truncated. Each relay is walked on its own
|
||||
* `until` cursor, up to [maxConcurrentRelays] at once, and every event funnels
|
||||
* through [verifyAndStore]; the result is tagged by the relay that first
|
||||
* delivered it. Unlike [drain], it IS deduped across relays: the same
|
||||
* widely-mirrored event arrives once per relay, and the repeats are dropped by a
|
||||
* [SeenIds] filter BEFORE the expensive verify+store — an id is marked seen only
|
||||
* after it verifies, so a forged copy (valid id, bad signature) delivered first
|
||||
* can't suppress the genuine one from another relay.
|
||||
* widely-mirrored event arrives once per relay, and the repeats are dropped by
|
||||
* the accessory's `SeenIds` filter BEFORE the expensive verify+store — an id is
|
||||
* marked seen only after it verifies, so a forged copy (valid id, bad signature)
|
||||
* delivered first can't suppress the genuine one from another relay.
|
||||
*
|
||||
* Bound the work with the filters' `limit`: each relay pages until it reaches
|
||||
* the limit, so an unbounded filter pages that relay's entire matching history.
|
||||
@@ -795,44 +597,12 @@ class Context(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
if (filters.isEmpty()) return emptyList()
|
||||
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
|
||||
// fetchAllPages' onEvent can't suspend, but verifyAndStore does — bridge
|
||||
// through a channel and verify+store single-threaded in one consumer so the
|
||||
// store writes stay serialized (same shape as `drain`).
|
||||
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(UNLIMITED)
|
||||
coroutineScope {
|
||||
val consumer =
|
||||
launch {
|
||||
// One writer → SeenIds' single-writer contract holds. Skip a
|
||||
// cross-relay duplicate before verifying it; mark it seen only once
|
||||
// it verifies so a bad-sig copy can't pre-empt a good one. Start
|
||||
// small (CLI fetches are typically hundreds of events); it grows if
|
||||
// an unbounded drain needs it, rather than eagerly taking the
|
||||
// large-walk default table.
|
||||
val seen = SeenIds(initialSlotsPow2 = 12)
|
||||
for ((relay, event) in eventChannel) {
|
||||
if (seen.contains(event.id)) continue
|
||||
if (verifyAndStore(event)) {
|
||||
seen.add(event.id)
|
||||
collected.add(relay to event)
|
||||
}
|
||||
}
|
||||
}
|
||||
try {
|
||||
client.fetchAllPagesFromPool(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
maxConcurrentRelays = maxConcurrentRelays,
|
||||
) { event, relay -> eventChannel.trySend(relay to event) }
|
||||
} finally {
|
||||
eventChannel.close()
|
||||
}
|
||||
consumer.join()
|
||||
}
|
||||
return collected
|
||||
}
|
||||
): List<Pair<NormalizedRelayUrl, Event>> =
|
||||
client.fetchAllPagesFromPoolWithHooks(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
maxConcurrentRelays = maxConcurrentRelays,
|
||||
) { _, event -> verifyAndStore(event) }
|
||||
|
||||
/**
|
||||
* Publish [request] to [relays], then wait for the FIRST event matching [responseFilter]
|
||||
|
||||
+21
@@ -1204,6 +1204,27 @@ class CashuWalletOps(
|
||||
existingSecrets: Set<String> = emptySet(),
|
||||
): RestoreOutcome = publishRecoveredProofs(scanRecoverableProofs(mintUrl, seed, startCounter, existingSecrets))
|
||||
|
||||
/**
|
||||
* [restoreFromMint], then bump the persisted NUT-13 counter past every
|
||||
* slot the scan confirmed in use — via the [peekCashuCounter] /
|
||||
* [reserveCashuCounters] hooks — so a later mint/swap can't reuse one of
|
||||
* the recovered secrets. Same advance rule Android's
|
||||
* `CashuWalletState.restoreFromMint` applies after its restore, shared
|
||||
* here so every front end bumps identically. `reserveCashuCounters` is
|
||||
* atomic, so two restores running concurrently can't collide.
|
||||
*/
|
||||
suspend fun restoreFromMintAdvancingCounter(
|
||||
mintUrl: String,
|
||||
seed: ByteArray,
|
||||
startCounter: Long = 0L,
|
||||
existingSecrets: Set<String> = emptySet(),
|
||||
): RestoreOutcome {
|
||||
val outcome = restoreFromMint(mintUrl, seed, startCounter, existingSecrets)
|
||||
val delta = (outcome.nextCounterAfterScan - peekCashuCounter(outcome.keysetId)).coerceAtLeast(0L)
|
||||
if (delta > 0) reserveCashuCounters(outcome.keysetId, delta.toInt())
|
||||
return outcome
|
||||
}
|
||||
|
||||
companion object {
|
||||
/**
|
||||
* How far to rewind the persisted counter when looking for proofs
|
||||
|
||||
+215
@@ -0,0 +1,215 @@
|
||||
/*
|
||||
* 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.nip01Core.relay.client.accessories
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
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.commands.toClient.MachineReadablePrefix
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.utils.SeenIds
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.selects.select
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
|
||||
/**
|
||||
* Option-rich sibling of [fetchAll]: subscribe [filters] across their relays,
|
||||
* funnel every arriving event through the suspending [onEvent] hook (verify /
|
||||
* persist / filter — return `true` to keep it in the result), and return the
|
||||
* accepted `(relay, event)` pairs once every relay reached a terminal state
|
||||
* (EOSE, CLOSED, or cannot-connect) or [timeoutMs] elapsed.
|
||||
*
|
||||
* Extras over [fetchAll]:
|
||||
* - **[onEvent] hook** — suspending per-event callback, invoked single-threaded
|
||||
* in arrival order, so callers can serialize verify+store work. Only events it
|
||||
* accepts (`true`) are collected. No cross-relay dedup is applied here — the
|
||||
* hook sees every copy.
|
||||
* - **[deadOut]** — when provided, every relay whose terminal reason classifies
|
||||
* as a hard failure via [classifyDrainFailure] (connect refused / DNS / TLS /
|
||||
* dead HTTP upgrade — NOT slow relays or 429s) is recorded, so callers can
|
||||
* prune proven-dead relays from future routing instead of paying the full
|
||||
* [timeoutMs] on them again.
|
||||
* - **[pendingOnAuthRequired]** — a relay that refuses the REQ with an
|
||||
* `auth-required:` CLOSED is kept pending rather than treated as terminal:
|
||||
* the caller's NIP-42 responder answers the challenge and the client re-fires
|
||||
* this same subscription, so the post-auth events are collected instead of
|
||||
* returning empty. If auth never satisfies it, the relay simply falls through
|
||||
* to the [timeoutMs].
|
||||
* - **[onTimeout]** — diagnostic hook fired when the deadline elapsed with
|
||||
* relays still pending: receives the stalled set, the terminal reasons seen so
|
||||
* far (`"eose"` / `"closed:<msg>"` / `"cannot:<msg>"`), and what was collected.
|
||||
*/
|
||||
suspend fun INostrClient.fetchAllWithHooks(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 8_000L,
|
||||
subscriptionId: String = newSubId(),
|
||||
pendingOnAuthRequired: Boolean = false,
|
||||
deadOut: MutableMap<NormalizedRelayUrl, DrainFailure>? = null,
|
||||
onTimeout: ((stalled: Set<NormalizedRelayUrl>, doneReasons: Map<NormalizedRelayUrl, String>, collected: List<Pair<NormalizedRelayUrl, Event>>) -> Unit)? = null,
|
||||
onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
if (filters.isEmpty()) return emptyList()
|
||||
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(UNLIMITED)
|
||||
// Carries the terminal reason per relay so a timeout can distinguish a slow
|
||||
// relay (never terminal) from a connect failure / CLOSED.
|
||||
val doneChannel = Channel<Pair<NormalizedRelayUrl, String>>(UNLIMITED)
|
||||
val remaining = filters.keys.toMutableSet()
|
||||
val doneReasons = HashMap<NormalizedRelayUrl, String>()
|
||||
val listener =
|
||||
object : SubscriptionListener {
|
||||
override fun onEvent(
|
||||
event: Event,
|
||||
isLive: Boolean,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
eventChannel.trySend(relay to event)
|
||||
}
|
||||
|
||||
override fun onEose(
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay to "eose")
|
||||
}
|
||||
|
||||
override fun onClosed(
|
||||
message: String,
|
||||
relay: NormalizedRelayUrl,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
// Keep the relay pending on an auth-required refusal: the authenticator answers the
|
||||
// challenge and re-fires this subscription, so the post-auth events still arrive.
|
||||
if (pendingOnAuthRequired && MachineReadablePrefix.parse(message) == MachineReadablePrefix.AUTH_REQUIRED) return
|
||||
doneChannel.trySend(relay to "closed:$message")
|
||||
}
|
||||
|
||||
override fun onCannotConnect(
|
||||
relay: NormalizedRelayUrl,
|
||||
message: String,
|
||||
forFilters: List<Filter>?,
|
||||
) {
|
||||
doneChannel.trySend(relay to "cannot:$message")
|
||||
}
|
||||
}
|
||||
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
|
||||
try {
|
||||
subscribe(subscriptionId, filters, listener)
|
||||
val completed =
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
while (remaining.isNotEmpty()) {
|
||||
select {
|
||||
eventChannel.onReceive { pair ->
|
||||
if (onEvent(pair.first, pair.second)) collected.add(pair)
|
||||
}
|
||||
doneChannel.onReceive { (relay, reason) ->
|
||||
remaining.remove(relay)
|
||||
doneReasons[relay] = reason
|
||||
}
|
||||
}
|
||||
}
|
||||
// Drain any events that landed after EOSE but before cancel
|
||||
while (true) {
|
||||
val r = eventChannel.tryReceive()
|
||||
if (!r.isSuccess) break
|
||||
val pair = r.getOrThrow()
|
||||
if (onEvent(pair.first, pair.second)) collected.add(pair)
|
||||
}
|
||||
true
|
||||
}
|
||||
if (completed == null && remaining.isNotEmpty()) {
|
||||
onTimeout?.invoke(remaining, doneReasons, collected)
|
||||
}
|
||||
} finally {
|
||||
unsubscribe(subscriptionId)
|
||||
eventChannel.close()
|
||||
doneChannel.close()
|
||||
}
|
||||
deadOut?.let { out ->
|
||||
for ((relay, reason) in doneReasons) {
|
||||
classifyDrainFailure(reason)?.let { out[relay] = it }
|
||||
}
|
||||
}
|
||||
return collected
|
||||
}
|
||||
|
||||
/**
|
||||
* [fetchAllPagesFromPool] with a suspending per-event hook: paginates every relay
|
||||
* to completion (each on its own `until` cursor, up to [maxConcurrentRelays] at
|
||||
* once) and funnels every event through [onEvent] — invoked single-threaded in
|
||||
* one consumer coroutine, so suspending verify/persist work stays serialized.
|
||||
* Returns the accepted `(relay, event)` pairs, tagged by the relay that first
|
||||
* delivered each.
|
||||
*
|
||||
* Unlike [fetchAllWithHooks], results ARE deduped across relays: the same
|
||||
* widely-mirrored event arrives once per relay, and the repeats are dropped by a
|
||||
* [SeenIds] filter BEFORE the (potentially expensive) [onEvent] — an id is marked
|
||||
* seen only after the hook accepts it, so a forged copy (valid id, bad signature)
|
||||
* delivered first can't suppress the genuine one from another relay.
|
||||
*/
|
||||
suspend fun INostrClient.fetchAllPagesFromPoolWithHooks(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
if (filters.isEmpty()) return emptyList()
|
||||
val collected = mutableListOf<Pair<NormalizedRelayUrl, Event>>()
|
||||
// fetchAllPagesFromPool's onEvent can't suspend, but the hook does — bridge
|
||||
// through a channel and run the hook single-threaded in one consumer so its
|
||||
// side effects (e.g. store writes) stay serialized.
|
||||
val eventChannel = Channel<Pair<NormalizedRelayUrl, Event>>(UNLIMITED)
|
||||
coroutineScope {
|
||||
val consumer =
|
||||
launch {
|
||||
// One writer → SeenIds' single-writer contract holds. Skip a
|
||||
// cross-relay duplicate before running the hook on it; mark it seen
|
||||
// only once the hook accepts it so a bad-sig copy can't pre-empt a
|
||||
// good one. Start small (one-shot fetches are typically hundreds of
|
||||
// events); it grows if an unbounded drain needs it, rather than
|
||||
// eagerly taking the large-walk default table.
|
||||
val seen = SeenIds(initialSlotsPow2 = 12)
|
||||
for ((relay, event) in eventChannel) {
|
||||
if (seen.contains(event.id)) continue
|
||||
if (onEvent(relay, event)) {
|
||||
seen.add(event.id)
|
||||
collected.add(relay to event)
|
||||
}
|
||||
}
|
||||
}
|
||||
try {
|
||||
fetchAllPagesFromPool(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
maxConcurrentRelays = maxConcurrentRelays,
|
||||
) { event, relay -> eventChannel.trySend(relay to event) }
|
||||
} finally {
|
||||
eventChannel.close()
|
||||
}
|
||||
consumer.join()
|
||||
}
|
||||
return collected
|
||||
}
|
||||
+2
@@ -20,6 +20,8 @@ Import as `com.vitorpamplona.quartz.nip01Core.relay.client.accessories.<name>` (
|
||||
| `fetchFirst(relay, filter, timeoutMs)` | `NostrClientFetchFirstExt` | Get the first matching event and stop (returns `null` on none/timeout). |
|
||||
| `fetchAllPages(relay, filters, timeoutMs)` | `NostrClientFetchAllPagesExt` | Fully retrieve a result set larger than the relay's per-REQ cap (strfry `limit`, ~500) by walking a `created_at` cursor. Bound it with the filter's `limit`. |
|
||||
| `fetchAllPagesFromPool(filters, ...)` | `NostrClientFetchAllPagesPoolExt` | Same paging, across several relays at once, deduped across them. |
|
||||
| `fetchAllWithHooks(filters, ...)` | `NostrClientFetchAllWithHooksExt` | `fetchAll` with a suspending per-`(relay, event)` accept hook (verify+store as events arrive), per-relay terminal-reason tracking, optional dead-relay collection (`deadOut` + `classifyDrainFailure`), keep-pending-on-`auth-required` CLOSED (NIP-42 re-fire), and a timeout diagnostic hook. |
|
||||
| `fetchAllPagesFromPoolWithHooks(filters, ...)` | `NostrClientFetchAllWithHooksExt` | `fetchAllPagesFromPool` with the same suspending accept hook, run single-threaded in one consumer; deduped across relays by `SeenIds` before the hook. |
|
||||
|
||||
## Streaming (`Flow`)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user