diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt index 09d44fd9bf..b9575826f1 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/Account.kt @@ -375,6 +375,18 @@ class Account( val mlsGroupStateStore: MlsGroupStateStore? = null, val marmotMessageStore: com.vitorpamplona.quartz.marmot.mls.group.MarmotMessageStore? = null, val marmotKeyPackageStore: com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageBundleStore? = null, + /** + * Durable publish obligations. Null means publish-before-apply does not + * survive a restart, so a commit interrupted mid-publish is replaced by a + * fresh one for the same epoch — a fork against the peers that took the + * first. + */ + val marmotPublishObligationStore: com.vitorpamplona.quartz.marmot.protocolCore.MarmotPublishObligationStore? = null, + /** + * Durable "already decided" markers for inbound events. Null means every + * backdated gift wrap is re-unwrapped on every sync. + */ + val marmotIngestDedupStore: com.vitorpamplona.quartz.marmot.MarmotIngestDedupStore? = null, val powQueue: () -> PoWPublishQueue? = { null }, relayAuthPermissionStore: RelayAuthPermissionStore = InMemoryRelayAuthPermissionStore(), signerPermissionStore: NostrSignerPermissionStore = InMemoryNostrSignerPermissionStore(), @@ -935,6 +947,8 @@ class Account( // acknowledged accept" rule; a plain `publish` would report // success for bytes nobody took. MarmotPublisher { event, relays -> client.publishAndConfirm(event, relays) }, + marmotPublishObligationStore, + marmotIngestDedupStore, scope = scope, ) } diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/accountsCache/AccountCacheState.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/accountsCache/AccountCacheState.kt index 0cb125444a..4a8f0f0b0a 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/accountsCache/AccountCacheState.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/accountsCache/AccountCacheState.kt @@ -32,9 +32,11 @@ import com.vitorpamplona.amethyst.commons.service.pow.PoWPublishQueue import com.vitorpamplona.amethyst.model.Account import com.vitorpamplona.amethyst.model.AccountSettings import com.vitorpamplona.amethyst.model.LocalCache +import com.vitorpamplona.amethyst.model.marmot.AndroidIngestDedupStore import com.vitorpamplona.amethyst.model.marmot.AndroidKeyPackageBundleStore import com.vitorpamplona.amethyst.model.marmot.AndroidMarmotMessageStore import com.vitorpamplona.amethyst.model.marmot.AndroidMlsGroupStateStore +import com.vitorpamplona.amethyst.model.marmot.AndroidPublishObligationStore import com.vitorpamplona.amethyst.service.location.LocationState import com.vitorpamplona.amethyst.service.relayClient.authCommand.model.DataStoreRelayAuthPermissionStore import com.vitorpamplona.quartz.nip01Core.core.HexKey @@ -266,6 +268,32 @@ class AccountCacheState( null } + val marmotPublishObligationStore = + try { + AndroidPublishObligationStore(accountDir) + } catch (e: Exception) { + Log.e( + "AccountCacheState", + "Failed to initialize AndroidPublishObligationStore " + + "(a Marmot commit interrupted mid-publish will NOT be retried after a restart)", + e, + ) + null + } + + val marmotIngestDedupStore = + try { + AndroidIngestDedupStore(accountDir) + } catch (e: Exception) { + Log.e( + "AccountCacheState", + "Failed to initialize AndroidIngestDedupStore " + + "(every backdated gift wrap will be re-decided on each sync)", + e, + ) + null + } + // Per-account NIP-42 ALLOW/DENY overrides live in this account's own dir, so a DENY for one // account never leaks into another (the store used to be a single app-wide file). val relayAuthPermissionStore = DataStoreRelayAuthPermissionStore(accountDir) @@ -291,6 +319,8 @@ class AccountCacheState( mlsGroupStateStore = mlsStore, marmotMessageStore = marmotMessageStore, marmotKeyPackageStore = marmotKeyPackageStore, + marmotPublishObligationStore = marmotPublishObligationStore, + marmotIngestDedupStore = marmotIngestDedupStore, powQueue = powQueue, relayAuthPermissionStore = relayAuthPermissionStore, signerPermissionStore = signerPermissionStore, diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidIngestDedupStore.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidIngestDedupStore.kt new file mode 100644 index 0000000000..b414d09ac3 --- /dev/null +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidIngestDedupStore.kt @@ -0,0 +1,90 @@ +/* + * 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.model.marmot + +import com.vitorpamplona.quartz.marmot.MarmotIngestDedupStore +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.utils.Log +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext +import java.io.File + +/** + * Android implementation of [MarmotIngestDedupStore] — one hex event id per + * line under `/marmot_ingested.ids`. + * + * Deliberately NOT encrypted: the file holds public relay event ids and no key + * material, and a marker lost to a decryption failure would silently cost a + * re-decision rather than fail loudly. + * + * Capped and trimmed oldest-first. The worst case for a forgotten marker is + * one wasted NIP-59 unwrap on the next sync, so bounding growth is worth more + * than remembering every id forever. + */ +class AndroidIngestDedupStore( + private val rootDir: File, + private val maxEntries: Int = 20_000, +) : MarmotIngestDedupStore { + private val mutex = Mutex() + + private fun file(): File = File(rootDir, "marmot_ingested.ids") + + override suspend fun mark(eventId: HexKey) = + withContext(Dispatchers.IO) { + mutex.withLock { + val target = file() + try { + target.parentFile?.mkdirs() + target.appendText(eventId + "\n") + if (target.length() > maxEntries.toLong() * 65L) { + val kept = target.readLines().filter { it.isNotBlank() }.takeLast(maxEntries / 2) + target.writeText(kept.joinToString("\n") + "\n") + } + } catch (e: Exception) { + Log.w(TAG, "could not record ingest marker: ${e.message}", e) + } + Unit + } + } + + override suspend fun loadAll(): Set = + withContext(Dispatchers.IO) { + mutex.withLock { + try { + file() + .takeIf { it.exists() } + ?.readLines() + ?.filter { it.isNotBlank() } + ?.toSet() + .orEmpty() + } catch (e: Exception) { + Log.w(TAG, "could not read ingest markers: ${e.message}", e) + emptySet() + } + } + } + + companion object { + private const val TAG = "AndroidIngestDedupStore" + } +} diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidPublishObligationStore.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidPublishObligationStore.kt new file mode 100644 index 0000000000..02baa30a3d --- /dev/null +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidPublishObligationStore.kt @@ -0,0 +1,121 @@ +/* + * 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.model.marmot + +import com.vitorpamplona.amethyst.model.preferences.KeyStoreEncryption +import com.vitorpamplona.quartz.marmot.protocolCore.MarmotPublishObligationStore +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.utils.Log +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext +import java.io.File + +/** + * Android implementation of [MarmotPublishObligationStore], encrypted at rest + * with [KeyStoreEncryption] like the group-state and KeyPackage stores. + * + * ``` + * /marmot_obligations/.obligation + * ``` + * + * Publish-before-apply only means anything if the record outlives the process. + * A commit is recorded, published, and only then applied; a crash inside that + * window has to leave a trace, or the next launch mints a REPLACEMENT commit + * for the same epoch and forks this device against every peer that accepted + * the first one. Android kills apps mid-work routinely, so "in memory" here is + * not a simplification — it is the common case. + * + * One file per obligation rather than one appended log: two groups can publish + * concurrently and resolve out of order, so removing one record must not + * rewrite another's. + */ +class AndroidPublishObligationStore( + private val rootDir: File, + private val encryption: KeyStoreEncryption = KeyStoreEncryption(), +) : MarmotPublishObligationStore { + private val mutex = Mutex() + + private fun dir(): File = File(rootDir, "marmot_obligations") + + private fun file(obligationId: String) = File(dir(), "$obligationId.obligation") + + override suspend fun save( + obligationId: HexKey, + bytes: ByteArray, + ) = withContext(Dispatchers.IO) { + mutex.withLock { + val target = file(obligationId) + try { + target.parentFile?.mkdirs() + val encrypted = encryption.encrypt(bytes) + val tmp = File(target.parentFile, "${target.name}.tmp") + tmp.writeBytes(encrypted) + if (!tmp.renameTo(target)) { + tmp.copyTo(target, overwrite = true) + if (!tmp.delete()) Log.w(TAG) { "could not delete temp file ${tmp.absolutePath}" } + } + } catch (e: Exception) { + // Failing to record is worse than failing to publish: an + // unrecorded commit that peers accept is a fork we cannot + // detect. Surface it rather than continuing to the publish. + Log.e(TAG, "save($obligationId) FAILED", e) + throw e + } + } + } + + override suspend fun delete(obligationId: HexKey) = + withContext(Dispatchers.IO) { + mutex.withLock { + val target = file(obligationId) + if (target.exists() && !target.delete()) { + Log.w(TAG) { "could not delete resolved obligation ${target.absolutePath}" } + } + Unit + } + } + + override suspend fun loadAll(): List = + withContext(Dispatchers.IO) { + mutex.withLock { + dir() + .listFiles { f -> f.isFile && f.name.endsWith(".obligation") } + ?.sortedBy { it.name } + ?.mapNotNull { file -> + try { + encryption.decrypt(file.readBytes()) + } catch (e: Exception) { + // One unreadable record must not cost us the + // others; the gate treats a missing obligation as + // "never confirmed", which is the safe direction. + Log.w(TAG, "unreadable obligation ${file.name}: ${e.message}", e) + null + } + }.orEmpty() + } + } + + companion object { + private const val TAG = "AndroidPublishObligationStore" + } +} diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/DecryptAndIndexProcessor.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/DecryptAndIndexProcessor.kt index ea018b04d9..6414d19f79 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/DecryptAndIndexProcessor.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/DecryptAndIndexProcessor.kt @@ -776,6 +776,17 @@ class GroupEventHandler( } } + is GroupEventResult.RefusedByLifecycle -> { + // Disbanded is absorbing and Unrecoverable needs a repair + // before anything more may be applied, so this input was + // refused before decryption. Nothing to render, nothing to + // retain, and nothing the user can do about it here. + Log.d("MarmotDbg") { + "GroupEventHandler.add: refused for group=${result.groupId.take(8)}… " + + "lifecycle=${result.lifecycle}" + } + } + is GroupEventResult.Error -> { Log.w("MarmotDbg") { "GroupEventHandler.add: ERROR ${result.message}" } } diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Config.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Config.kt index ba408f8158..2cd961752f 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Config.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Config.kt @@ -228,6 +228,17 @@ class DataDir( val groupsDir = File(marmotDir, "groups") val keyPackageBundleFile = File(marmotDir, "keypackages.bundle") + /** + * Unresolved publish obligations. Durable because publish-before-apply is + * only meaningful across a crash: without this, a commit recorded and then + * lost to a restart is replaced by a fresh one for the same epoch, forking + * us against the peers that accepted the first. + */ + val publishObligationsDir = File(marmotDir, "obligations") + + /** Inbound events this account has terminally decided about. */ + val ingestDedupFile = File(marmotDir, "ingested.ids") + /** * SQLite event-store DB file, a sibling of [eventsDir] under * `/shared/`. Used when the store backend is SQLite (the 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 e8c3c44946..be88e3861b 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/Context.kt @@ -21,9 +21,11 @@ package com.vitorpamplona.amethyst.cli import com.sun.management.UnixOperatingSystemMXBean +import com.vitorpamplona.amethyst.cli.stores.FileIngestDedupStore import com.vitorpamplona.amethyst.cli.stores.FileKeyPackageBundleStore import com.vitorpamplona.amethyst.cli.stores.FileMarmotMessageStore import com.vitorpamplona.amethyst.cli.stores.FileMlsGroupStateStore +import com.vitorpamplona.amethyst.cli.stores.FilePublishObligationStore import com.vitorpamplona.amethyst.commons.cashu.CashuWalletReader import com.vitorpamplona.amethyst.commons.cashu.ops.CashuWalletOps import com.vitorpamplona.amethyst.commons.cashu.ops.RestoreOutcome @@ -319,6 +321,8 @@ class Context( private val mlsStore by lazy { FileMlsGroupStateStore(dataDir.groupsDir) } private val keyPackageStore by lazy { FileKeyPackageBundleStore(dataDir.keyPackageBundleFile) } private val messageStore by lazy { FileMarmotMessageStore(dataDir.groupsDir) } + private val publishObligationStore by lazy { FilePublishObligationStore(dataDir.publishObligationsDir) } + private val ingestDedupStore by lazy { FileIngestDedupStore(dataDir.ingestDedupFile) } /** * Shared Nostr event store for this run, opened via [StoreFactory] @@ -356,6 +360,8 @@ class Context( // once a relay in the group's own scope returns OK true. Anything // weaker (queued, sent, no error yet) is explicitly not success. MarmotPublisher { event, relays -> client.publishAndConfirm(event, relays) }, + publishObligationStore, + ingestDedupStore, ) } diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt index 53dfa0ae18..2fe55c7dd4 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/stores/FileStores.kt @@ -22,9 +22,14 @@ package com.vitorpamplona.amethyst.cli.stores import com.vitorpamplona.amethyst.cli.SecureFileIO import com.vitorpamplona.amethyst.commons.util.deleteOrWarn +import com.vitorpamplona.quartz.marmot.MarmotIngestDedupStore import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageBundleStore import com.vitorpamplona.quartz.marmot.mls.group.MarmotMessageStore import com.vitorpamplona.quartz.marmot.mls.group.MlsGroupStateStore +import com.vitorpamplona.quartz.marmot.protocolCore.MarmotPublishObligationStore +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock import java.io.File /** @@ -148,3 +153,81 @@ class FileMarmotMessageStore( file(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group messages") } } + +/** + * Durable publish obligations, one file per obligation under [dir]. + * + * Publish-before-apply only means anything if the obligation outlives the + * process: the whole point is that a commit is prepared, recorded, published, + * and only then applied, so a crash between record and publish must leave a + * trace. With a non-durable store that window silently becomes "the commit + * never happened", and on relaunch the client mints a REPLACEMENT commit for + * the same epoch — forking itself against the peers that accepted the first + * one. + * + * A file per obligation rather than one appended log: obligations resolve out + * of order (two groups publish concurrently), and deleting one must not + * rewrite the others. + */ +class FilePublishObligationStore( + private val dir: File, +) : MarmotPublishObligationStore { + init { + SecureFileIO.secureMkdirs(dir) + } + + private fun file(obligationId: String) = File(dir, "$obligationId.obligation") + + override suspend fun save( + obligationId: HexKey, + bytes: ByteArray, + ) { + SecureFileIO.writeBytesAtomic(file(obligationId), bytes) + } + + override suspend fun delete(obligationId: HexKey) { + file(obligationId).deleteOrWarn("FilePublishObligationStore", "publish obligation") + } + + override suspend fun loadAll(): List = + dir + .listFiles { f -> f.isFile && f.name.endsWith(".obligation") } + ?.sortedBy { it.name } + ?.mapNotNull { runCatching { it.readBytes() }.getOrNull() } + .orEmpty() +} + +/** + * Durable "already decided" markers, one hex id per line. + * + * Append-only and capped: the point is to stop re-deciding backdated gift + * wraps forever, not to remember every event this account has ever seen. When + * the cap is hit the oldest half is dropped — the worst case for a forgotten + * marker is one wasted re-decision, so trading memory for exactness is the + * right way round. + */ +class FileIngestDedupStore( + private val file: File, + private val maxEntries: Int = 20_000, +) : MarmotIngestDedupStore { + private val mutex = Mutex() + + override suspend fun mark(eventId: HexKey) = + mutex.withLock { + SecureFileIO.appendText(file, eventId + "\n") + if (file.length() > maxEntries.toLong() * 65L) { + val kept = file.readLines().filter { it.isNotBlank() }.takeLast(maxEntries / 2) + SecureFileIO.writeBytesAtomic(file, (kept.joinToString("\n") + "\n").encodeToByteArray()) + } + } + + override suspend fun loadAll(): Set = + mutex.withLock { + file + .takeIf { it.exists() } + ?.readLines() + ?.filter { it.isNotBlank() } + ?.toSet() + .orEmpty() + } +} diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt index ac51cca3a3..54f51999d4 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotIngest.kt @@ -105,7 +105,28 @@ suspend fun MarmotManager.ingest(event: Event): MarmotIngestResult = else -> MarmotIngestResult.Ignored } -private suspend fun MarmotManager.ingestGiftWrap(wrap: GiftWrapEvent): MarmotIngestResult = +private suspend fun MarmotManager.ingestGiftWrap(wrap: GiftWrapEvent): MarmotIngestResult { + // A relay `since` cursor cannot skip a backdated event, and NIP-59 wraps + // are backdated by up to two days on purpose — so without a durable marker + // every wrap in that band is unwrapped and re-decided on every single sync. + if (isTerminallyIngested(wrap.id)) return MarmotIngestResult.Ignored + val result = ingestGiftWrapUncached(wrap) + when (result) { + // Joined, or already in the group: nothing more can come of this wrap. + is MarmotIngestResult.JoinedGroup, is MarmotIngestResult.AlreadyInGroup -> markTerminallyIngested(wrap.id) + + // A Welcome naming a KeyPackage whose private half we never held can + // never become processable: bundles are generated locally BEFORE the + // KeyPackage is published, so one we do not have is one we never will. + is MarmotIngestResult.Failure -> + if (result.message.contains("No matching KeyPackageBundle")) markTerminallyIngested(wrap.id) + + else -> Unit + } + return result +} + +private suspend fun MarmotManager.ingestGiftWrapUncached(wrap: GiftWrapEvent): MarmotIngestResult = try { // NIP-59 wraps carry two encryption layers (kind:1059 → kind:13 → rumor). // [unwrapAndUnsealOrNull] peels both so we land directly on the inner @@ -158,6 +179,9 @@ private suspend fun MarmotManager.ingestGroupEvent(ge: GroupEvent): MarmotIngest // Decrypted only on a losing branch: real protocol input (it may have // witnessed for that branch), but never application output. is GroupEventResult.AppMessageOnCandidateBranch, + // Disbanded or locally unrecoverable — refused before decryption, so + // there is nothing to deliver and nothing to retain. + is GroupEventResult.RefusedByLifecycle, -> { MarmotIngestResult.Ignored } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt index 08b968cf48..94aab47d43 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotManager.kt @@ -24,6 +24,7 @@ import com.vitorpamplona.amethyst.commons.model.marmotGroups.MarmotGroupChatroom import com.vitorpamplona.amethyst.commons.model.marmotGroups.MarmotGroupImage import com.vitorpamplona.quartz.marmot.GroupEventResult import com.vitorpamplona.quartz.marmot.MarmotInboundProcessor +import com.vitorpamplona.quartz.marmot.MarmotIngestDedupStore import com.vitorpamplona.quartz.marmot.MarmotOutboundProcessor import com.vitorpamplona.quartz.marmot.MarmotSubscriptionManager import com.vitorpamplona.quartz.marmot.MarmotWelcomeSender @@ -70,6 +71,8 @@ import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock import kotlin.io.encoding.Base64 import kotlin.io.encoding.ExperimentalEncodingApi @@ -102,6 +105,14 @@ class MarmotManager( */ val publisher: MarmotPublisher = MarmotPublisher { _, _ -> false }, publishObligationStore: MarmotPublishObligationStore? = null, + /** + * Durable "already decided" markers for inbound events. + * + * Null means every backdated gift wrap is re-unwrapped and re-decided on + * every sync, forever — correct, but it costs a full NIP-59 double + * decryption per wrap per sync and keeps the relay busy re-serving them. + */ + private val ingestDedupStore: MarmotIngestDedupStore? = null, /** * Scope used to carry an open convergence pass to its cutoff. * @@ -151,12 +162,93 @@ class MarmotManager( // Also restore previously-published KeyPackage bundles so that // Welcomes referencing them remain processable across restarts. keyPackageRotationManager.restoreFromStore() + ingestDedupStore?.loadAll()?.let { marks -> + terminallyIngestedMutex.withLock { terminallyIngested.addAll(marks) } + Unit + } + retryPendingPublishObligations() Log.d("MarmotManager") { "restoreAll(): done, ${activeIds.size} groups: $activeIds" } } catch (e: Exception) { Log.e("MarmotManager", "Failed to restore Marmot state", e) } } + /** + * Event ids this client has terminally decided about — see + * [MarmotIngestDedupStore]. Loaded once in [restoreAll]; the in-memory set + * is the hot path, the store only makes it survive a restart. + */ + private val terminallyIngested = mutableSetOf() + private val terminallyIngestedMutex = Mutex() + + suspend fun isTerminallyIngested(eventId: HexKey): Boolean = terminallyIngestedMutex.withLock { eventId in terminallyIngested } + + suspend fun markTerminallyIngested(eventId: HexKey) { + val added = terminallyIngestedMutex.withLock { terminallyIngested.add(eventId) } + if (added) { + try { + ingestDedupStore?.mark(eventId) + } catch (e: Exception) { + // A marker we failed to persist costs a re-decision next + // launch; it never costs correctness, so it is not worth + // failing the ingest over. + Log.w("MarmotManager", "could not persist ingest marker for ${eventId.take(8)}: ${e.message}", e) + } + } + } + + /** + * Re-emit every publish obligation a previous run left unresolved. + * + * The bytes are republished VERBATIM — the same signed kind-445, to the + * same recipient scope. That is the whole reason the obligation stores + * them rather than storing "there was a commit": a replacement commit for + * the same epoch is a fork against every peer that accepted the first one, + * and a peer that already has this event simply deduplicates it. + * + * A retry that still fails leaves the group in `PendingPublish`, which is + * the safe direction: it blocks new local commits until the group actually + * knows what happened to this one. + */ + suspend fun retryPendingPublishObligations() { + val pending = publishGate.allPending() + if (pending.isEmpty()) return + Log.d("MarmotManager") { "retryPendingPublishObligations(): ${pending.size} unresolved" } + for (obligation in pending) { + val event = + try { + Event.fromJson(obligation.outboundBytes.decodeToString()) + } catch (e: Exception) { + // A record we cannot decode can never be republished, and + // holding the group in PendingPublish forever helps nobody. + Log.w("MarmotManager", "unreadable publish obligation ${obligation.obligationId}: ${e.message}", e) + publishGate.resolve(obligation.obligationId, PublishOutcome.FAILED) + continue + } + val relays = obligation.recipientScope.mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) }.toSet() + val confirmed = + if (relays.isEmpty()) { + false + } else { + try { + publisher.publish(event, relays) + } catch (e: Exception) { + Log.w("MarmotManager", "publish retry failed for ${obligation.groupId}: ${e.message}", e) + false + } + } + val state = + publishGate.resolve( + obligation.obligationId, + if (confirmed) PublishOutcome.CONFIRMED else PublishOutcome.UNKNOWN, + ) + Log.d("MarmotManager") { + "retryPendingPublishObligations(): ${obligation.groupId.take(8)}… " + + "confirmed=$confirmed lifecycle=$state" + } + } + } + /** * Computes a per-group kind:445 subscription `since` from the newest * persisted decrypted message of each group. @@ -231,6 +323,7 @@ class MarmotManager( is GroupEventResult.Duplicate, is GroupEventResult.UndecryptableOuterLayer, is GroupEventResult.AppMessageOnCandidateBranch, + is GroupEventResult.RefusedByLifecycle, is GroupEventResult.Error, -> {} } @@ -267,7 +360,10 @@ class MarmotManager( suspend fun buildGroupMessage( nostrGroupId: HexKey, innerEvent: Event, - ): OutboundGroupEvent = outboundProcessor.buildGroupEvent(nostrGroupId, innerEvent) + ): OutboundGroupEvent { + requireOutboundAllowed(nostrGroupId, "send a message") + return outboundProcessor.buildGroupEvent(nostrGroupId, innerEvent) + } /** * Build a kind:9 chat-message GroupEvent from plain text. The inner @@ -562,6 +658,7 @@ class MarmotManager( relays: List, stage: suspend () -> MlsGroupManager.StagedCommit, ): CommitPublication { + requireOutboundAllowed(nostrGroupId, "commit a group-state change") check(publishGate.canPrepareLocalCommit(nostrGroupId)) { "Group $nostrGroupId cannot prepare a local commit " + "(lifecycle=${publishGate.lifecycle(nostrGroupId)}, gate=${publishGate.outboundGate(nostrGroupId)})" @@ -597,7 +694,10 @@ class MarmotManager( publishGate.resolve( obligation.obligationId, - if (confirmed) PublishOutcome.CONFIRMED else PublishOutcome.FAILED, + // Not confirmed is NOT the same as rejected: a timeout or a + // dropped connection leaves us unable to say whether a peer took + // the commit, so the obligation stays retryable. + if (confirmed) PublishOutcome.CONFIRMED else PublishOutcome.UNKNOWN, ) if (confirmed) { @@ -660,8 +760,55 @@ class MarmotManager( private val settlerRunning = MutableStateFlow(false) - /** Lifecycle state for a group, including any unresolved publish obligation. */ - suspend fun lifecycle(nostrGroupId: HexKey): GroupLifecycleState = publishGate.lifecycle(nostrGroupId) + /** + * The group's effective lifecycle state. + * + * Two components track lifecycle for different reasons and neither is the + * whole answer on its own: the publish gate owns the local-publish states + * (`PendingPublish`, `Merging`), convergence owns `Recovering` and the two + * terminal-ish states (`Disbanded`, `Unrecoverable`). Reading only the + * publish gate — as this used to — meant a disbanded or unrecoverable + * group still reported `Stable` and every outbound gate keyed on it + * happily let work through. + * + * Terminal wins: a group that convergence has terminalized is terminal + * regardless of what the publish gate is doing, because the publish that + * gate is tracking can no longer be applied to anything. + */ + suspend fun lifecycle(nostrGroupId: HexKey): GroupLifecycleState { + val converged = inboundProcessor.groupLifecycle(nostrGroupId) + if (converged == GroupLifecycleState.DISBANDED || converged == GroupLifecycleState.UNRECOVERABLE) { + return converged + } + val publish = publishGate.lifecycle(nostrGroupId) + return if (publish == GroupLifecycleState.STABLE) converged else publish + } + + /** + * Refuse outbound work a group's lifecycle does not permit. + * + * The check is here rather than at each call site because every outbound + * path has the same answer: a `Disbanded` group takes no further work of + * any kind, and an `Unrecoverable` one takes none until it is repaired — + * encrypting against state we do not trust produces a message the group + * will invalidate, which is worse than refusing. + */ + private suspend fun requireOutboundAllowed( + nostrGroupId: HexKey, + what: String, + ) { + when (val state = lifecycle(nostrGroupId)) { + GroupLifecycleState.DISBANDED -> + throw IllegalStateException("Group $nostrGroupId is disbanded; cannot $what") + + GroupLifecycleState.UNRECOVERABLE -> + throw IllegalStateException( + "Group $nostrGroupId is unrecoverable locally; repair, restore or rejoin before you $what", + ) + + else -> Log.d("MarmotManager") { "$what allowed for ${nostrGroupId.take(8)}… in $state" } + } + } /** * The group's own relay list, as the recipient scope for a publish diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishBeforeApplyTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishBeforeApplyTest.kt index 9100f51ea3..bc15a68cb5 100644 --- a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishBeforeApplyTest.kt +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishBeforeApplyTest.kt @@ -237,9 +237,17 @@ class MarmotPublishBeforeApplyTest { /** * A group whose publisher never acknowledges anything can still be read. * It simply cannot advance — which is the safe direction to fail. + * + * It also cannot start over. "No OK arrived" is not "no peer took it": a + * timeout or a dropped connection leaves us unable to say. Discarding the + * obligation and letting a SECOND commit be prepared for the same epoch is + * how that uncertainty becomes a permanent fork — the peer that did + * receive the first commit is at epoch 1, rejects our second, and the two + * copies never reconcile. So the obligation stays, the group stays held, + * and the retry republishes the SAME bytes. */ @Test - fun aFailedPublishLeavesTheGroupUsable() = + fun anUnconfirmedPublishHoldsTheGroupInsteadOfStartingOver() = runBlocking { val fx = fixture(accepts = false) fx.manager.addMember( @@ -250,15 +258,24 @@ class MarmotPublishBeforeApplyTest { relays = listOf(relay), ) - assertEquals(GroupLifecycleState.STABLE, fx.manager.lifecycle(fx.groupId)) - assertTrue( + assertEquals(GroupLifecycleState.PENDING_PUBLISH, fx.manager.lifecycle(fx.groupId)) + assertEquals( + 1, fx.manager.publishGate .pendingFor(fx.groupId) - .isEmpty(), - "a failed obligation is discarded, not left pending forever", + .size, + "an unconfirmed obligation stays retryable", ) - // A second attempt is allowed: nothing was consumed by the failure. - val retry = + // Reading is unaffected; only advancing the group is blocked. + assertEquals( + 0L, + fx.manager.groupManager + .getGroup(fx.groupId)!! + .epoch, + ) + + // A second, different commit for the same epoch is refused. + assertFailsWith { fx.manager.addMember( nostrGroupId = fx.groupId, memberPubKey = fx.bobPubKey, @@ -266,11 +283,8 @@ class MarmotPublishBeforeApplyTest { keyPackageEventId = "c".repeat(64), relays = listOf(relay), ) - assertEquals(2, fx.publisher.published.size) - assertTrue( - retry.first.signedEvent.id - .isNotEmpty(), - ) + } + assertEquals(1, fx.publisher.published.size, "no replacement commit was offered") } /** The group's own relay list is the default recipient scope. */ diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishDurabilityTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishDurabilityTest.kt new file mode 100644 index 0000000000..7f3e6998c3 --- /dev/null +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/MarmotPublishDurabilityTest.kt @@ -0,0 +1,233 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.amethyst.commons.marmot + +import com.vitorpamplona.quartz.marmot.mip01Groups.MarmotGroupData +import com.vitorpamplona.quartz.marmot.mls.group.MlsGroupStateStore +import com.vitorpamplona.quartz.marmot.protocolCore.GroupLifecycleState +import com.vitorpamplona.quartz.marmot.protocolCore.MarmotPublishObligationStore +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +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.relay.normalizer.RelayUrlNormalizer +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal +import kotlinx.coroutines.runBlocking +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * Publish-before-apply is only a real rule if it survives the process. + * + * A commit is recorded, published, and only then applied. A crash inside that + * window has to leave a durable trace, because the alternative is that the next + * launch has no memory of the commit, mints a REPLACEMENT for the same epoch, + * and forks this client against every peer that accepted the first one. So + * these tests restart the manager against the same stores and assert the + * obligation is still there, is retried with the SAME bytes, and only then + * lets the group move. + */ +class MarmotPublishDurabilityTest { + private val relay: NormalizedRelayUrl = RelayUrlNormalizer.normalizeOrNull("wss://relay.example.com")!! + + /** Answers with a verdict that can be flipped between "runs". */ + private class SwitchablePublisher( + var accepts: Boolean, + ) : MarmotPublisher { + val published = mutableListOf() + + override suspend fun publish( + event: Event, + relays: Set, + ): Boolean { + published.add(event) + return accepts + } + } + + private class MemoryObligationStore : MarmotPublishObligationStore { + val entries = LinkedHashMap() + + override suspend fun save( + obligationId: HexKey, + bytes: ByteArray, + ) { + entries[obligationId] = bytes + } + + override suspend fun delete(obligationId: HexKey) { + entries.remove(obligationId) + } + + override suspend fun loadAll(): List = entries.values.toList() + } + + /** Same MLS state across "restarts", like a real store on disk. */ + private class MemoryStateStore : MlsGroupStateStore { + private val states = mutableMapOf() + private val retained = mutableMapOf>() + + override suspend fun save( + nostrGroupId: String, + state: ByteArray, + ) { + states[nostrGroupId] = state + } + + override suspend fun load(nostrGroupId: String): ByteArray? = states[nostrGroupId] + + override suspend fun delete(nostrGroupId: String) { + states.remove(nostrGroupId) + retained.remove(nostrGroupId) + } + + override suspend fun listGroups(): List = states.keys.toList() + + override suspend fun saveRetainedEpochs( + nostrGroupId: String, + epochs: List, + ) { + retained[nostrGroupId] = epochs + } + + override suspend fun loadRetainedEpochs(nostrGroupId: String): List = retained[nostrGroupId].orEmpty() + } + + private fun manager( + signer: NostrSignerInternal, + store: MlsGroupStateStore, + obligations: MarmotPublishObligationStore, + publisher: MarmotPublisher, + ) = MarmotManager( + signer, + store, + publisher = publisher, + publishObligationStore = obligations, + ) + + @Test + fun anUnresolvedObligationOutlivesTheProcessAndIsRetriedVerbatim() = + runBlocking { + val signer = NostrSignerInternal(KeyPair()) + val store = MemoryStateStore() + val obligations = MemoryObligationStore() + val publisher = SwitchablePublisher(accepts = false) + val groupId = "b".repeat(64) + + val first = manager(signer, store, obligations, publisher) + first.createGroup( + groupId, + MarmotGroupData( + nostrGroupId = groupId, + adminPubkeys = listOf(signer.pubKey), + relays = listOf(relay.url), + ), + ) + val bob = KeyPair() + val bundle = + first.groupManager + .getGroup(groupId)!! + .createKeyPackage(bob.pubKey, ByteArray(0)) + + // The publish is refused, so the group must NOT move and the + // obligation must remain on disk. + runCatching { + first.addMember( + nostrGroupId = groupId, + memberPubKey = bob.pubKey.toHexKey(), + keyPackageBytes = bundle.keyPackage.toTlsBytes(), + keyPackageEventId = "d".repeat(64), + relays = listOf(relay), + ) + } + assertEquals(GroupLifecycleState.PENDING_PUBLISH, first.lifecycle(groupId)) + assertEquals(0L, first.groupManager.getGroup(groupId)!!.epoch) + assertEquals(1, obligations.entries.size, "an unacknowledged commit leaves its obligation durable") + val recordedBytes = + obligations.entries.values + .first() + .copyOf() + val firstAttempt = publisher.published.single() + + // Restart against the same stores. The relay accepts this time. + publisher.accepts = true + val second = manager(signer, store, obligations, publisher) + second.restoreAll() + + val retry = publisher.published.last() + assertEquals( + firstAttempt.toJson(), + retry.toJson(), + "the retry republishes the same event, not a replacement commit for the same epoch", + ) + assertEquals(GroupLifecycleState.STABLE, second.lifecycle(groupId)) + assertEquals(1L, second.groupManager.getGroup(groupId)!!.epoch, "the acknowledged commit applies") + assertTrue(obligations.entries.isEmpty(), "a resolved obligation is deleted") + assertTrue(recordedBytes.isNotEmpty()) + } + + @Test + fun aRetryThatFailsAgainKeepsTheGroupHeld() = + runBlocking { + val signer = NostrSignerInternal(KeyPair()) + val store = MemoryStateStore() + val obligations = MemoryObligationStore() + val publisher = SwitchablePublisher(accepts = false) + val groupId = "c".repeat(64) + + val first = manager(signer, store, obligations, publisher) + first.createGroup( + groupId, + MarmotGroupData( + nostrGroupId = groupId, + adminPubkeys = listOf(signer.pubKey), + relays = listOf(relay.url), + ), + ) + val carol = KeyPair() + val bundle = + first.groupManager + .getGroup(groupId)!! + .createKeyPackage(carol.pubKey, ByteArray(0)) + runCatching { + first.addMember( + nostrGroupId = groupId, + memberPubKey = carol.pubKey.toHexKey(), + keyPackageBytes = bundle.keyPackage.toTlsBytes(), + keyPackageEventId = "e".repeat(64), + relays = listOf(relay), + ) + } + + val second = manager(signer, store, obligations, publisher) + second.restoreAll() + + // Still refused: the group stays held rather than quietly moving + // on, which is what stops a new commit stacking on an epoch peers + // never accepted. + assertEquals(GroupLifecycleState.PENDING_PUBLISH, second.lifecycle(groupId)) + assertEquals(0L, second.groupManager.getGroup(groupId)!!.epoch) + assertEquals(1, obligations.entries.size) + assertTrue(bundle.keyPackage.toTlsBytes().isNotEmpty()) + } +} diff --git a/quartz/plans/2026-09-08-marmot-spec-resync.md b/quartz/plans/2026-09-08-marmot-spec-resync.md index 183c22ae0e..945b79efd7 100644 --- a/quartz/plans/2026-09-08-marmot-spec-resync.md +++ b/quartz/plans/2026-09-08-marmot-spec-resync.md @@ -584,46 +584,75 @@ Writing the producer side immediately found two bugs the reader-side tests could just started enforcing, so every KeyPackage we published would have been rejected by any conformant peer — including, once Stage 4 landed, by us. Now `now - 1h` to `+84 days`. -## What is NOT done +## Interop status (2026-09-09) -- **The interop harness now RUNS but does not pass.** It used to die in preflight; it now - builds MDK 0.9.20 (against the same OpenMLS fork rev `mdk-vector-gen` pins), boots - nostr-rs-relay, brings up both `wnd` daemons and `amy`, and executes all 17 scenarios. Every - one fails, all downstream of Test 01 (MDK cannot find A's KeyPackage). Four environment - blockers were fixed to get that far, all recorded in the harness: - - `protoc` is a build prerequisite MDK now needs. - - MDK 0.9.x requires `WN_ALLOW_LOOPBACK_RELAYS=1` before it will accept a `ws://` loopback - relay at all; without it `wnd` exits before creating its socket. - - MDK refuses to create its socket unless the socket's parent directory is `0700`. - - `wn --json whoami` moved to `{"ok":true,"result":{"accounts":[…]}}`; the harness's - `extract_pubkey` probed only the older array shapes and silently returned nothing. +The MDK 0.9.20 harness runs end to end. It went **1 passed / 12 failed** to +**9 passed** over this pass, and the failures that remain are named below rather +than lumped together. - Two real defects in our own code came out of the run, both fixed: - - `amy relay add` reported success from the DECISION to write, not the store's answer, so a - rejected or no-op write printed `added: yes`. - - Our own KeyPackage/DM relay lists were read back through the local-network filter. That - filter is right for someone else's list — it is attacker-supplied input, and it is also what - exempts a relay from Tor — but applying it to a list we published ourselves made a - deliberately configured local relay look like no configuration at all. The publisher then - fell back to a default set, and A's KeyPackage went to five PUBLIC relays instead of the - harness's loopback. `allRelays()` now exists for reading back our own lists, and the - KeyPackage publish goes only to the configured relay. +### Defects the harness found in our own code - **Still blocking Test 01:** the kind-10051 list persists under `relay key-package set` but not - under `relay add`/`relay key-package add`, so MDK finds no relay list to fetch A's KeyPackage - from. `verifyAndStore` returns true and kind 10050 works through the identical code path, so - this is a storage/CLI issue rather than a protocol one, and it needs its own focused pass. +Each of these was invisible to every same-implementation test we have, because +each is a place where two implementations have to agree on something one +implementation alone never disagrees with. - **Also unresolved, and it is a design conflict rather than a bug:** MDK accepts `ws://` ONLY - for a loopback host, while quartz strips exactly those hosts out of relay lists. No address - satisfies both, so a loopback-relay harness cannot work until one side moves. Changing a - Tor-adjacent privacy guard is a maintainer decision, not one to make in passing. -- **`MarmotManager.createGroup` (the MIP-era path) is still the one the UI calls.** - `createCurrentProfileGroup` exists, is wired, and is tested, but the Android and desktop - "new group" flows still call the legacy one. KeyPackage publishing HAS switched: it now - defaults to the current profile, which is the half that decides whether anyone can invite us. -- **Lifecycle enforcement covers the publish path, not everything.** `PendingPublish`, - `Merging` and the outbound gates are enforced; `Unrecoverable` and `Disbanded` are still a - correct model with nothing driving them. -- Stage 7: durability/restart conformance, app payload kinds `1009`/`1210`, encrypted-media v2, - the push owner proof (kind `451`). +1. **A Commit rebuilt our own leaf from defaults.** The UpdatePath leaf replaces + OUR leaf — same member, new key material — so its capabilities and extensions + must carry over. `buildLeafNode` was called with neither, so the very first + invite we ever sent dropped the `account-identity-proof` (a LEAF extension no + proposal can restore) and stopped advertising the extension and proposal the + group's own `required_capabilities` demanded. MDK reported + `PublicGroupError(LeafNodeValidation(UnsupportedExtensions))` and dropped the + Welcome minted by that same commit; the invitee simply never saw an invite, + with nothing logged on either side. + +2. **We held private keys for our leaf only, not our direct path.** RFC 9420 + §7.6 lets a committer address us at any node in the copath resolution whose + key we hold, and a merged subtree resolves to its PARENT — so from three + members on, the ciphertext meant for us stops naming our leaf. Worked for + two members, failed for three, which is exactly why it survived every + two-party test. `MlsGroupState` v3 persists them. + +3. **Three readers only understood the legacy `0xF2EE` extension**, so they did + nothing at all on a current-profile group: the Welcome's `nostr_group_id` + (without which a joiner cannot even subscribe), the admin gate on + GroupContextExtensions changes (which skipped rather than failed closed), and + disappearing-message expiration. Group metadata reads and writes now go + through `MarmotManager.groupView` / `setGroupProfile` / `setGroupAdmins` / + `setGroupImage`, which dispatch on the profile the group actually uses. + +4. **A non-confirmed publish discarded its obligation**, letting a REPLACEMENT + commit be prepared for the same epoch. "No OK arrived" is not "no peer took + it": a timeout or a dropped connection leaves it unknown, and minting a + second commit for an epoch a peer may already hold is the fork the gate + exists to prevent. `PublishOutcome.UNKNOWN` keeps the record and holds the + group. + +5. **The harness was parsing a wire format `wn` no longer speaks.** MDK 0.9.x + returns `{"ok":true,"result":{"invites":[…]}}`; iterating `.result` walked + that object's three VALUES, so every poll matched nothing and reported "never + received invite" for welcomes that had arrived and been accepted. + +### What is NOT done + +- **Test 03 gets further but does not pass.** We join MDK's group and see its + name; the first kind-445 after the join is not delivered. Tests 05/12/14/15 + fail behind it with "A never received invite". +- **Tests 13 and 16 (KeyPackage rotation) fail.** `wn keys check` finds no prior + KeyPackage for A at that point in the run, and amy keeps seeing B's + pre-rotation KeyPackage. +- **Test 09 fails on the relay, not on us**: `disconnected before OK` from the + local nostr-rs-relay under the load of a full run. The durable ingest markers + added in this pass cut a large part of that load (a backdated gift wrap used + to be re-unwrapped on every sync, forever) but the test has not been re-run + since. +- **Agent-text-stream is receive-only.** We decode the `0x8006` policy, derive + per-stream record keys, open records and fold the transcript, and we advertise + the `0xF2D1` receive capability. We do NOT advertise `send` (`0xF2D2`) or + `fanout` (`0xF2D4`): publishing needs durable per-stream sequence state to + avoid reusing an AEAD nonce across a restart, and there is none. A group whose + policy requires `send` is refused at join rather than joined into a state + every peer would reject us from. +- **The QUIC transport for agent text streams is not wired.** The record layer + and the kind-1200 anchor are implemented; nothing yet opens a WebTransport + session to a broker and feeds it records. diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotInboundProcessor.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotInboundProcessor.kt index 9a876ca8f2..0ec8f0f677 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotInboundProcessor.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotInboundProcessor.kt @@ -133,6 +133,21 @@ sealed class GroupEventResult { val countedAsWitness: Boolean, ) : GroupEventResult() + /** + * The group's lifecycle state refuses this input outright. + * + * `Disbanded` is absorbing — no later branch supersedes a terminalized + * disband, so there is nothing a subsequent kind-445 could do but be + * retained forever. `Unrecoverable` means this client cannot safely apply + * more traffic until it repairs, restores or rejoins; retaining input + * against material we no longer trust is how a client talks itself into + * settling for a state it should have refused. + */ + data class RefusedByLifecycle( + val groupId: HexKey, + val lifecycle: GroupLifecycleState, + ) : GroupEventResult() + /** * The event could not be processed. */ @@ -264,6 +279,15 @@ class MarmotInboundProcessor( return GroupEventResult.Error(groupId, "Not a member of group $groupId") } + // Lifecycle BEFORE anything is retained. A disbanded group is + // terminal and an unrecoverable one cannot safely apply traffic, so + // input for either is refused here rather than decrypted, retained, + // and then quietly never acted on. + val lifecycle = convergence.lifecycle(groupId) + if (lifecycle == GroupLifecycleState.DISBANDED || lifecycle == GroupLifecycleState.UNRECOVERABLE) { + return GroupEventResult.RefusedByLifecycle(groupId, lifecycle) + } + // Settle FIRST, so this event is processed against resolved state // rather than against a branch a pass is about to abandon. Inbound // traffic is only an opportunistic tick, though: a group that goes diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotIngestDedupStore.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotIngestDedupStore.kt new file mode 100644 index 0000000000..15b970c9bd --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/MarmotIngestDedupStore.kt @@ -0,0 +1,60 @@ +/* + * 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.marmot + +import com.vitorpamplona.quartz.nip01Core.core.HexKey + +/** + * Durable record of inbound events this client has TERMINALLY decided about. + * + * A relay subscription cannot express "everything after event X" — only + * `since`, a timestamp — and NIP-59 gift wraps are deliberately backdated by up + * to two days, so a `since` cursor that has caught up still re-delivers every + * wrap in that band on every sync. Without a durable marker each of those is + * unwrapped, decrypted and re-decided from scratch, forever: a client that has + * been in a few groups spends most of its sync budget re-deciding events it + * already resolved, and hammers the relay doing it. + * + * "Terminally decided" is a narrow claim. It covers outcomes that cannot change + * with more information — a Welcome we joined from, and one naming a + * KeyPackage whose private half we never held (bundles are generated locally + * before the KeyPackage is published, so a bundle we do not have is one we + * never will). It does NOT cover an event that is merely undecryptable right + * now: a kind-445 encrypted under an epoch we have not reached yet becomes + * readable the moment the commit arrives, and marking it would lose the + * message permanently. + */ +interface MarmotIngestDedupStore { + suspend fun mark(eventId: HexKey) + + suspend fun loadAll(): Set +} + +/** Non-durable default. Re-decides everything after a restart. */ +class InMemoryIngestDedupStore : MarmotIngestDedupStore { + private val ids = LinkedHashSet() + + override suspend fun mark(eventId: HexKey) { + ids.add(eventId) + } + + override suspend fun loadAll(): Set = ids.toSet() +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroup.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroup.kt index 950e542230..b44a71cbe1 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroup.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroup.kt @@ -135,6 +135,20 @@ class MlsGroup private constructor( /** Staged keys from proposeSigningKeyRotation — only promoted on successful commit */ private var pendingSigningKey: ByteArray? = null, private var pendingEncryptionKey: ByteArray? = null, + /** + * HPKE private keys for the PARENT nodes on our own direct path, keyed by + * node index. + * + * RFC 9420 §7.6 does not say "the committer encrypts to your leaf" — it + * says the committer encrypts one path secret per node in the copath + * resolution, and you decrypt at whichever of those nodes you hold a key + * for. A merged subtree resolves to its PARENT, so as soon as a group has + * three members the commits addressed to us stop naming our leaf at all. + * Keeping only the leaf key is why every MDK commit after a three-member + * Add failed with "UpdatePath at common ancestor carries no ciphertext + * for us". + */ + private val pathPrivateKeys: MutableMap = mutableMapOf(), ) { val groupId: ByteArray get() = groupContext.groupId val epoch: Long get() = groupContext.epoch @@ -295,9 +309,38 @@ class MlsGroup private constructor( // rewind our own generation counter to 0 and reuse an AEAD // key+nonce within this epoch (RFC 9420 §9). senderRatchetStates = secretTree.exportSenderStates(), + pathPrivateKeys = pathPrivateKeys.toMap(), ) } + /** + * Record the HPKE private keys our direct path gained from [pathSecret] at + * [fromNodeIndex] and every node above it, up to the root. + * + * A path secret ratchets one KDF step per level regardless of filtering, + * and each level's node keypair is `DeriveKeyPair(DeriveSecret(secret, + * "node"))` — the same derivation the committer used, which is what makes + * the keys we store here the ones a later committer will encrypt to. + */ + private fun rememberPathKeys( + fromNodeIndex: Int, + pathSecret: ByteArray, + ) { + val fullPath = BinaryTree.directPath(myLeafIndex, tree.leafCount) + val start = fullPath.indexOf(fromNodeIndex) + if (start < 0) return + var secret = pathSecret + for (i in start until fullPath.size) { + val nodeSecret = MlsCryptoProvider.deriveSecret(secret, "node") + pathPrivateKeys[fullPath[i]] = Hpke.deriveKeyPair(nodeSecret).privateKey + secret = MlsCryptoProvider.deriveSecret(secret, "path") + } + // Anything no longer on our direct path (the tree reshaped under us) + // can never be addressed to us again; holding it would only make a + // stale key look usable at the next resolution scan. + pathPrivateKeys.keys.retainAll(fullPath.toSet()) + } + /** * Extract retained epoch secrets for late-message decryption. * @@ -634,6 +677,17 @@ class MlsGroup private constructor( val leafSecret = MlsCryptoProvider.randomBytes(MlsCryptoProvider.HASH_OUTPUT_LENGTH) val pathSecrets = tree.derivePathSecrets(myLeafIndex, leafSecret) + // We just minted the keys for our whole direct path. Keep the private + // halves: the next committer will address us at one of these nodes, + // not at our leaf, as soon as our subtree is merged. + run { + val fullPath = BinaryTree.directPath(myLeafIndex, tree.leafCount) + pathPrivateKeys.keys.retainAll(fullPath.toSet()) + for ((i, nodeIdx) in fullPath.withIndex()) { + pathSecrets.getOrNull(i)?.let { pathPrivateKeys[nodeIdx] = it.privateKey } + } + } + // RFC 9420 §12.4.1: newly-added leaves (from Add proposals in THIS commit) // MUST be excluded from the copath resolution — they join via the Welcome // at epoch N+1 and don't need the path secret. Keeping them in the list @@ -1810,24 +1864,67 @@ class MlsGroup private constructor( BinaryTree.nodeToLeaf(resNode) in newLeavesInCommit } - // Find which encrypted secret corresponds to our position + // RFC 9420 §7.6: the committer encrypts one path secret per node + // in the copath resolution, and we decrypt at whichever of those + // nodes we hold a private key for. That is usually NOT our leaf — + // a merged subtree resolves to its parent, so from three members + // on we are addressed at an ancestor. Scan the resolution for a + // key we actually have rather than assuming our own leaf node. val myNodeIdx = BinaryTree.leafToNode(myLeafIndex) - val myResIdx = resolution.indexOf(myNodeIdx) - check(myResIdx in 0 until pathNode.encryptedPathSecret.size) { + val candidates = + resolution.withIndex().mapNotNull { (i, resNode) -> + if (i >= pathNode.encryptedPathSecret.size) { + null + } else { + val key = if (resNode == myNodeIdx) encryptionPrivateKey else pathPrivateKeys[resNode] + key?.let { Triple(i, resNode, it) } + } + } + check(candidates.isNotEmpty()) { "UpdatePath at common ancestor carries no ciphertext for us " + "(my_leaf=$myLeafIndex, my_node=$myNodeIdx, resolution=$resolution, " + + "held_path_nodes=${pathPrivateKeys.keys.sorted()}, " + "encrypted_path_secrets=${pathNode.encryptedPathSecret.size})" } - val ct = pathNode.encryptedPathSecret[myResIdx] - val pathSecret = - MlsCryptoProvider.decryptWithLabel( - encryptionPrivateKey, - "UpdatePathNode", - pathDecContextBytes, - ct.kemOutput, - ct.ciphertext, - ) + // Try each node we hold a key for rather than committing to the + // first. An Add or Remove renumbers nodes, so a retained key can + // outlive the node it belonged to; a stale one fails the AEAD + // rather than producing a wrong secret, so trying the next + // candidate is exact, not a guess. + var pathSecret: ByteArray? = null + var decryptedAt = -1 + for ((i, _, key) in candidates) { + val ct = pathNode.encryptedPathSecret[i] + pathSecret = + try { + MlsCryptoProvider.decryptWithLabel( + key, + "UpdatePathNode", + pathDecContextBytes, + ct.kemOutput, + ct.ciphertext, + ) + } catch (_: Exception) { + null + } + if (pathSecret != null) { + decryptedAt = i + break + } + } + val recoveredPathSecret = + checkNotNull(pathSecret) { + "UpdatePath at common ancestor did not decrypt with any key we hold " + + "(my_leaf=$myLeafIndex, resolution=$resolution, slot_tried=$decryptedAt, " + + "tried=${candidates.map { it.second }})" + } + + // The path secret we just recovered belongs to the common ancestor + // and ratchets up to the root, so it hands us the private key for + // every node above it on our own direct path. Those are exactly + // the nodes a later committer may address us at. + rememberPathKeys(commonAncestorNode, recoveredPathSecret) // Derive remaining path secrets from common ancestor up to root, // then one more step to reach the `commit_secret` (RFC 9420 §9.2: @@ -1842,7 +1939,7 @@ class MlsGroup private constructor( // UpdatePath node, so filtering changes which nodes carry // ciphertext but not the number of KDF steps. val stepsToRoot = unfilteredDirectPath.size - commonAncestorUnfilteredIdx - 1 - var currentSecret = pathSecret + var currentSecret = recoveredPathSecret repeat(stepsToRoot) { currentSecret = MlsCryptoProvider.deriveSecret(currentSecret, "path") } @@ -3991,6 +4088,7 @@ class MlsGroup private constructor( signingPrivateKey = state.signingPrivateKey, encryptionPrivateKey = state.encryptionPrivateKey, interimTranscriptHash = state.interimTranscriptHash, + pathPrivateKeys = state.pathPrivateKeys.toMutableMap(), ) } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroupState.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroupState.kt index cef888ddf1..518631e7fa 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroupState.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/mls/group/MlsGroupState.kt @@ -60,6 +60,17 @@ data class MlsGroupState( val interimTranscriptHash: ByteArray, val encryptionSecret: ByteArray, val senderRatchetStates: Map = emptyMap(), + /** + * HPKE private keys for the PARENT nodes on our own direct path, by node + * index (STATE_VERSION 3+). + * + * RFC 9420 §7.6 lets a committer address us at any node in the copath + * resolution whose key we hold — usually an ancestor rather than our leaf, + * because a merged subtree resolves to its parent. A member that keeps + * only its leaf key cannot decrypt those commits at all, and losing these + * across a restart makes the same group undecryptable on relaunch. + */ + val pathPrivateKeys: Map = emptyMap(), ) { fun encodeTls(): ByteArray { val writer = TlsWriter() @@ -115,6 +126,13 @@ data class MlsGroupState( writer.putUint32(ratchet.applicationGeneration.toLong()) } + // Direct-path node private keys (STATE_VERSION 3+). + writer.putUint32(pathPrivateKeys.size.toLong()) + for ((nodeIndex, key) in pathPrivateKeys) { + writer.putUint32(nodeIndex.toLong()) + writer.putOpaqueVarInt(key) + } + return writer.toByteArray() } @@ -135,8 +153,11 @@ data class MlsGroupState( * v1: original layout (no SecretTree ratchet positions). * v2: appends [senderRatchetStates] so restores don't reset the * ratchet to generation 0. v1 blobs still decode (empty map). + * v3: appends [pathPrivateKeys] so a restore can still decrypt an + * UpdatePath addressed at one of our ancestors. Older blobs decode + * with an empty map and refill on the next commit we process. */ - private const val STATE_VERSION = 2 + private const val STATE_VERSION = 3 fun decodeTls(data: ByteArray): MlsGroupState { val reader = TlsReader(data) @@ -197,6 +218,22 @@ data class MlsGroupState( emptyMap() } + // v3+: direct-path node private keys. Absent for older blobs, + // which restore able to decrypt only commits addressed at their + // own leaf until the next commit refills the path. + val pathPrivateKeys = + if (version >= 3 && reader.hasRemaining) { + val count = reader.readUint32().toInt() + buildMap { + repeat(count) { + val nodeIndex = reader.readUint32().toInt() + put(nodeIndex, reader.readOpaqueVarInt()) + } + } + } else { + emptyMap() + } + return MlsGroupState( groupContext = groupContext, treeBytes = treeBytes, @@ -208,6 +245,7 @@ data class MlsGroupState( interimTranscriptHash = interimTranscriptHash, encryptionSecret = encryptionSecret, senderRatchetStates = senderRatchetStates, + pathPrivateKeys = pathPrivateKeys, ) } } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotConvergenceEngine.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotConvergenceEngine.kt index 17d6dfa8f7..6965b53086 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotConvergenceEngine.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotConvergenceEngine.kt @@ -26,6 +26,7 @@ import com.vitorpamplona.quartz.marmot.mls.group.MlsGroupManager import com.vitorpamplona.quartz.marmot.mls.group.MlsGroupState import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.toHexKey +import com.vitorpamplona.quartz.utils.Log import com.vitorpamplona.quartz.utils.sha256.sha256 import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock @@ -227,8 +228,57 @@ class MarmotConvergenceEngine( } ctx.canonicalCommits.addLast(candidateOf(commitBytes, sourceEpoch)) trim(ctx) + terminalizeIfDisbanded(groupId, ctx) } + /** + * Move a group to `Disbanded` once its lifecycle component says so. + * + * `Disbanded` is absorbing: there is no outgoing transition, no later + * branch supersedes a terminalized disband, and a replacement conversation + * is a new MLS group. Deriving it from the applied state rather than from + * a transport claim is the whole point — the only thing that can disband a + * group is an authenticated Commit that every member replays identically. + */ + private fun terminalizeIfDisbanded( + groupId: HexKey, + ctx: GroupContext, + ) { + if (ctx.lifecycle == GroupLifecycleState.DISBANDED) return + val disbanded = groupManager.getGroup(groupId)?.currentGroupState()?.isDisbanded == true + if (disbanded) ctx.lifecycle = GroupLifecycleState.DISBANDED + } + + /** + * Mark a group locally unrecoverable. + * + * Local to ONE client: it does not mean the group is dead, it means this + * client cannot safely apply more traffic until it repairs, restores, + * rejoins or discards its copy. Settling for the current local state just + * because it is the only one available is exactly what this state exists + * to prevent, so it also drops any pass in flight rather than letting it + * resolve against material we no longer trust. + */ + suspend fun markUnrecoverable(groupId: HexKey) = + mutex.withLock { + val ctx = contexts.getOrPut(groupId) { GroupContext() } + if (ctx.lifecycle == GroupLifecycleState.DISBANDED) return@withLock + ctx.lifecycle = GroupLifecycleState.UNRECOVERABLE + ctx.pass = null + } + + /** + * Clear `Unrecoverable` after a verified repair — a replacement Welcome, a + * restore, or a rejoin. `Disbanded` is NOT clearable. + */ + suspend fun markRepaired(groupId: HexKey) = + mutex.withLock { + val ctx = contexts[groupId] ?: return@withLock + if (ctx.lifecycle == GroupLifecycleState.UNRECOVERABLE) { + ctx.lifecycle = GroupLifecycleState.STABLE + } + } + /** * Offer a commit that did NOT extend the canonical tip. * @@ -431,7 +481,25 @@ class MarmotConvergenceEngine( val rewound = selectedTipId != null && selectedTipId != inputs.tipId if (rewound) { - groupManager.installState(groupId, graph.statesById.getValue(selectedTipId)) + // The selected tip's state must be rebuildable from retained + // material. When it is not — the anchor the rewind needs fell out + // of the window, or a retained state failed to replay — this + // client cannot reach the branch the group selected, and the one + // thing it must NOT do is keep its own losing branch and call that + // settled. That is exactly `Unrecoverable`: local, repairable, and + // never resolved by pretending the pass succeeded. + val target = graph.statesById[selectedTipId] + if (target == null) { + markUnrecoverable(groupId) + return null + } + try { + groupManager.installState(groupId, target) + } catch (e: Exception) { + Log.w("MarmotConvergence", "rewind of $groupId to the selected branch failed: ${e.message}", e) + markUnrecoverable(groupId) + return null + } } return mutex.withLock { @@ -458,6 +526,7 @@ class MarmotConvergenceEngine( trimCandidates(ctx) ctx.divergent.clear() ctx.pass = null + terminalizeIfDisbanded(groupId, ctx) val epoch = groupManager.getGroup(groupId)?.epoch ?: inputs.tipEpoch ConvergenceResolution( groupId = groupId, diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotPublishGate.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotPublishGate.kt index 8f8fce335c..dfdcee4033 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotPublishGate.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/marmot/protocolCore/MarmotPublishGate.kt @@ -122,7 +122,25 @@ enum class PublishOutcome { /** At least one endpoint in the recipient scope acknowledged an accept. */ CONFIRMED, - /** Every endpoint rejected, or the attempt is known not to have been accepted. */ + /** + * The attempt did not confirm, and we cannot prove no endpoint took it. + * + * A timeout, a dropped connection, or an OK that never arrived all land + * here, and none of them means "no peer has this commit". The obligation + * stays durable and the group stays held, because minting a REPLACEMENT + * commit for the same epoch would fork us against whoever did receive the + * first one. This is the default for a publisher that answers with a + * boolean: `false` is "unconfirmed", not "rejected". + */ + UNKNOWN, + + /** + * The attempt can never succeed and there is nothing to retry — an + * unreadable record, or bytes no endpoint could ever accept. + * + * Discards the obligation. Use it only when retrying is impossible, never + * as a synonym for "did not get an OK". + */ FAILED, } @@ -280,6 +298,17 @@ class MarmotPublishGate( groupManager.installState(obligation.groupId, obligation.pendingState) } + if (outcome == PublishOutcome.UNKNOWN) { + // Keep the record and keep the group held. The staged commit was + // never applied, so nothing local is wrong — what is unknown is + // whether a PEER took it, and preparing a fresh commit while that + // is unknown is exactly the fork this gate exists to prevent. + return mutex.withLock { + lifecycles[obligation.groupId] = GroupLifecycleState.PENDING_PUBLISH + GroupLifecycleState.PENDING_PUBLISH + } + } + store.delete(obligationId) return mutex.withLock { pending.remove(obligationId) diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/MlsGroupStateTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/MlsGroupStateTest.kt index 34be9d3317..080f997502 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/MlsGroupStateTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/MlsGroupStateTest.kt @@ -173,12 +173,14 @@ class MlsGroupStateTest { val state = group.saveState() val bytes = state.encodeTls() - // First two bytes should be the version (uint16 = 2) + // First two bytes are the state version (uint16). v3 appends the + // direct-path node private keys; older blobs still decode, so the + // version only ever moves forward when the layout gains a field. val reader = com.vitorpamplona.quartz.marmot.mls.codec .TlsReader(bytes) val version = reader.readUint16() - assertEquals(2, version) + assertEquals(3, version) } @Test diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/group/UpdatePathAncestorTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/group/UpdatePathAncestorTest.kt new file mode 100644 index 0000000000..e078989678 --- /dev/null +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/marmot/mls/group/UpdatePathAncestorTest.kt @@ -0,0 +1,134 @@ +/* + * 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.marmot.mls.group + +import com.vitorpamplona.quartz.marmot.mls.messages.KeyPackageBundle +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertTrue + +/** + * RFC 9420 §7.6 does not say "the committer encrypts the path secret to your + * leaf". It says the committer encrypts one secret per node in the copath + * RESOLUTION, and each member decrypts at whichever of those nodes it holds a + * private key for. A merged subtree resolves to its parent, so from three + * members on the ciphertext meant for us is addressed at an ANCESTOR. + * + * Keeping only our leaf key therefore worked for two members and failed for + * three, which is exactly how it survived every same-implementation test: + * `UpdatePath at common ancestor carries no ciphertext for us (my_leaf=1, + * my_node=2, resolution=[1], encrypted_path_secrets=1)`. Node 1 was the parent + * we had held a key for since the commit that merged us, and we never stored + * it. + */ +class UpdatePathAncestorTest { + private fun bundleFor(seed: Byte): KeyPackageBundle = + MlsGroup + .create(identity = ByteArray(32) { seed }) + .createKeyPackage(identity = ByteArray(32) { seed }, signingKey = ByteArray(32) { seed }) + + @Test + fun aThirdPartyCommitReachesUsAtAMergedAncestor() { + // Alice creates, adds Bob, then adds Carol. After Bob's own commit the + // {alice, bob} subtree is merged, so a later commit from Carol + // resolves that subtree to its parent rather than to Bob's leaf. + val alice = MlsGroup.create(identity = ByteArray(32) { 0x0a }) + val bobBundle = bundleFor(0x0b) + val carolBundle = bundleFor(0x0c) + + alice.proposeAdd(bobBundle.keyPackage.toTlsBytes()) + val addBob = alice.commit() + val bob = MlsGroup.processWelcome(assertNotNull(addBob.welcomeBytes), bobBundle) + + alice.proposeAdd(carolBundle.keyPackage.toTlsBytes()) + val addCarol = alice.commit() + bob.processFramedCommit(addCarol.framedCommitBytes) + val carol = MlsGroup.processWelcome(assertNotNull(addCarol.welcomeBytes), carolBundle) + + // Bob commits: this merges Bob's direct path and hands him the parent + // keys the next committer will address him at. + val bobCommit = bob.commit() + alice.processFramedCommit(bobCommit.framedCommitBytes) + carol.processFramedCommit(bobCommit.framedCommitBytes) + + // Carol commits. Bob is now inside a merged subtree, so Carol's + // UpdatePath addresses him at a parent node, not at his leaf. + val carolCommit = carol.commit() + alice.processFramedCommit(carolCommit.framedCommitBytes) + bob.processFramedCommit(carolCommit.framedCommitBytes) + + assertEquals(carol.epoch, bob.epoch) + assertEquals(carol.epoch, alice.epoch) + assertContentEquals( + carol.exporterSecret("marmot", "group-event".encodeToByteArray(), 32), + bob.exporterSecret("marmot", "group-event".encodeToByteArray(), 32), + ) + assertContentEquals( + carol.exporterSecret("marmot", "group-event".encodeToByteArray(), 32), + alice.exporterSecret("marmot", "group-event".encodeToByteArray(), 32), + ) + } + + /** + * The keys are useless if a restart drops them: the group would keep + * working until the next commit addressed us at an ancestor and then stop + * dead, with nothing in the log tying the failure to the relaunch. + */ + @Test + fun theAncestorKeysSurviveASaveAndRestore() { + val alice = MlsGroup.create(identity = ByteArray(32) { 0x0a }) + val bobBundle = bundleFor(0x0b) + val carolBundle = bundleFor(0x0c) + + alice.proposeAdd(bobBundle.keyPackage.toTlsBytes()) + val addBob = alice.commit() + var bob = MlsGroup.processWelcome(assertNotNull(addBob.welcomeBytes), bobBundle) + + alice.proposeAdd(carolBundle.keyPackage.toTlsBytes()) + val addCarol = alice.commit() + bob.processFramedCommit(addCarol.framedCommitBytes) + val carol = MlsGroup.processWelcome(assertNotNull(addCarol.welcomeBytes), carolBundle) + + val bobCommit = bob.commit() + alice.processFramedCommit(bobCommit.framedCommitBytes) + carol.processFramedCommit(bobCommit.framedCommitBytes) + + // Round-trip Bob through the persisted blob, exactly as a relaunch does. + val saved = bob.saveState() + assertTrue(saved.pathPrivateKeys.isNotEmpty(), "a committer holds keys for its own direct path") + bob = MlsGroup.restore(MlsGroupStateCodec.roundTrip(saved)) + + val carolCommit = carol.commit() + bob.processFramedCommit(carolCommit.framedCommitBytes) + + assertEquals(carol.epoch, bob.epoch) + assertContentEquals( + carol.exporterSecret("marmot", "group-event".encodeToByteArray(), 32), + bob.exporterSecret("marmot", "group-event".encodeToByteArray(), 32), + ) + } +} + +private object MlsGroupStateCodec { + fun roundTrip(state: MlsGroupState): MlsGroupState = MlsGroupState.decodeTls(state.encodeTls()) +}