From 6e664b2db66574770ff5617f786673eb58db3e14 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 18 Jul 2026 22:55:56 +0000 Subject: [PATCH] refactor(cli): extract the drain loop to quartz accessories; split Context per domain MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 Claude-Session: https://claude.ai/code/session_01CP4kfLCa3wWtE8Khy21Pkj --- .../amethyst/cli/CashuContext.kt | 152 ++++++++ .../vitorpamplona/amethyst/cli/ConcordAuth.kt | 69 ++++ .../com/vitorpamplona/amethyst/cli/Context.kt | 334 +++--------------- .../commons/cashu/ops/CashuWalletOps.kt | 21 ++ .../NostrClientFetchAllWithHooksExt.kt | 215 +++++++++++ .../relay/client/accessories/README.md | 2 + 6 files changed, 511 insertions(+), 282 deletions(-) create mode 100644 cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/CashuContext.kt create mode 100644 cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/ConcordAuth.kt create mode 100644 quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllWithHooksExt.kt diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/CashuContext.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/CashuContext.kt new file mode 100644 index 0000000000..9542c53a8b --- /dev/null +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/CashuContext.kt @@ -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 `/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( + Filter( + authors = listOf(pk), + kinds = + listOf( + CashuWalletEvent.KIND, + CashuTokenEvent.KIND, + CashuSpendingHistoryEvent.KIND, + CashuMintQuoteEvent.KIND, + NutzapInfoEvent.KIND, + MintRecommendationEvent.KIND, + ), + ), + ) + val inboundNutzaps = + ctx.store.query( + 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) + } +} diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/ConcordAuth.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/ConcordAuth.kt new file mode 100644 index 0000000000..ffeb44831e --- /dev/null +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/ConcordAuth.kt @@ -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>() + private val streamSigners = ConcurrentHashMap() + + /** Registers raw 32-byte Concord stream [secrets] to answer NIP-42 challenges from [relays]. */ + fun register( + relays: Set, + secrets: List, + ) { + 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, + ): List = + streamSecrets[relay].orEmpty().mapNotNull { hex -> + runCatching { + streamSigners.getOrPut(hex) { NostrSignerSync(KeyPair(privKey = hex.hexToByteArray())) }.sign(template) + }.getOrNull() + } +} 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 bc50f62720..1766c08fd9 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -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>() - private val concordStreamSigners = ConcurrentHashMap() + private val concordAuth = ConcordAuth() /** Registers raw 32-byte Concord stream [secrets] to answer NIP-42 challenges from [relays]. */ fun registerConcordStreamKeys( relays: Set, secrets: List, - ) { - 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, - ): List = - 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 `/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( - Filter( - authors = listOf(pk), - kinds = - listOf( - CashuWalletEvent.KIND, - CashuTokenEvent.KIND, - CashuSpendingHistoryEvent.KIND, - CashuMintQuoteEvent.KIND, - NutzapInfoEvent.KIND, - MintRecommendationEvent.KIND, - ), - ), - ) - val inboundNutzaps = - store.query( - 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? = null, pendingOnAuthRequired: Boolean = false, - ): List> { - if (filters.isEmpty()) return emptyList() - val eventChannel = Channel>(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>(UNLIMITED) - val remaining = filters.keys.toMutableSet() - val doneReasons = HashMap() - val subId = newSubId() - val listener = - object : SubscriptionListener { - override fun onEvent( - event: Event, - isLive: Boolean, - relay: NormalizedRelayUrl, - forFilters: List?, - ) { - eventChannel.trySend(relay to event) - } - - override fun onEose( - relay: NormalizedRelayUrl, - forFilters: List?, - ) { - doneChannel.trySend(relay to "eose") - } - - override fun onClosed( - message: String, - relay: NormalizedRelayUrl, - forFilters: List?, - ) { - // 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?, - ) { - doneChannel.trySend(relay to "cannot:$message") - } - } - val collected = mutableListOf>() - 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> = + 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>, timeoutMs: Long = 30_000, maxConcurrentRelays: Int = 8, - ): List> { - if (filters.isEmpty()) return emptyList() - val collected = mutableListOf>() - // 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>(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> = + client.fetchAllPagesFromPoolWithHooks( + filters = filters, + timeoutMs = timeoutMs, + maxConcurrentRelays = maxConcurrentRelays, + ) { _, event -> verifyAndStore(event) } /** * Publish [request] to [relays], then wait for the FIRST event matching [responseFilter] diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cashu/ops/CashuWalletOps.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cashu/ops/CashuWalletOps.kt index 8ce9df8b77..a84604e7be 100644 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cashu/ops/CashuWalletOps.kt +++ b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cashu/ops/CashuWalletOps.kt @@ -1204,6 +1204,27 @@ class CashuWalletOps( existingSecrets: Set = 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 = 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 diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllWithHooksExt.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllWithHooksExt.kt new file mode 100644 index 0000000000..5c1156f5e9 --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/NostrClientFetchAllWithHooksExt.kt @@ -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:"` / `"cannot:"`), and what was collected. + */ +suspend fun INostrClient.fetchAllWithHooks( + filters: Map>, + timeoutMs: Long = 8_000L, + subscriptionId: String = newSubId(), + pendingOnAuthRequired: Boolean = false, + deadOut: MutableMap? = null, + onTimeout: ((stalled: Set, doneReasons: Map, collected: List>) -> Unit)? = null, + onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean, +): List> { + if (filters.isEmpty()) return emptyList() + val eventChannel = Channel>(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>(UNLIMITED) + val remaining = filters.keys.toMutableSet() + val doneReasons = HashMap() + val listener = + object : SubscriptionListener { + override fun onEvent( + event: Event, + isLive: Boolean, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + eventChannel.trySend(relay to event) + } + + override fun onEose( + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + doneChannel.trySend(relay to "eose") + } + + override fun onClosed( + message: String, + relay: NormalizedRelayUrl, + forFilters: List?, + ) { + // 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?, + ) { + doneChannel.trySend(relay to "cannot:$message") + } + } + val collected = mutableListOf>() + 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>, + timeoutMs: Long = 30_000L, + maxConcurrentRelays: Int = 8, + onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean, +): List> { + if (filters.isEmpty()) return emptyList() + val collected = mutableListOf>() + // 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>(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 +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/README.md b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/README.md index 46ce631095..a7adfef4a9 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/README.md +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/nip01Core/relay/client/accessories/README.md @@ -20,6 +20,8 @@ Import as `com.vitorpamplona.quartz.nip01Core.relay.client.accessories.` ( | `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`)