Merge pull request #4232 from vitorpamplona/fix/marmot-updatepath-fork

fix(marmot): treat the relay echo of our own unconfirmed commit as its confirmation
This commit is contained in:
Vitor Pamplona
2026-09-27 16:07:27 -04:00
committed by GitHub
5 changed files with 186 additions and 26 deletions
@@ -210,6 +210,7 @@ ALL_TESTS=(
test_27_deletion_wn_to_amy
test_28_retention_wn_to_amy
test_29_disband_amy_to_wn
test_30_wn_commit_after_app_data_update
test_31_reaction_materializes_on_wn
)
+51
View File
@@ -251,6 +251,57 @@ test_17_group_image_commit() {
fi
}
# The device sequence that forked Amethyst out of a group White Noise joined:
# amy creates, invites wn, commits an AppDataUpdate (the app's "Use encrypted
# attachments" is one; set-retention is the same proposal type), promotes wn,
# and then wn makes its FIRST commit — a rename carrying an UpdatePath. On the
# device Amethyst refused that commit ("UpdatePath at common ancestor carries
# no ciphertext for us") and every later wn message failed to decrypt.
test_30_wn_commit_after_app_data_update() {
banner "Test 30 — wn's first commit after an amy AppDataUpdate commit"
local id="30 wn commit after app-data"
local out gid mls_gid b_gid
out=$(amy_json marmot group create --name "Interop-30") || {
record_result "$id" fail "amy group create failed"; return
}
gid=$(printf '%s' "$out" | jq -r '.group_id')
mls_gid=$(printf '%s' "$out" | jq -r '.mls_group_id')
amy_json marmot group add "$gid" "$B_NPUB" >/dev/null || {
record_result "$id" fail "amy could not invite wn"; return
}
b_gid=$(wait_for_invite B 60) || { record_result "$id" fail "wn never received the Welcome"; return; }
wn_b groups accept "$b_gid" >/dev/null 2>&1 || true
wn_group_field_becomes "$mls_gid" '.group.group_id // empty' "$mls_gid" 120 || {
record_result "$id" fail "wn never surfaced the group"; return
}
wn_b messages send "$mls_gid" "30 before" >/dev/null 2>&1 || true
amy_json marmot await message "$gid" --match "30 before" --timeout 90 >/dev/null || {
record_result "$id" fail "amy never received wn's first message"; return
}
amy_json marmot group set-retention "$gid" 3600 >/dev/null || {
record_result "$id" fail "amy set-retention failed"; return
}
sleep 3
amy_json marmot group promote "$gid" "$B_NPUB" >/dev/null || {
record_result "$id" fail "amy promote failed"; return
}
sleep 5
wn_b groups rename "$mls_gid" "Interop-30-by-wn" >/dev/null 2>&1 || true
if ! amy_json marmot await rename "$gid" --name "Interop-30-by-wn" --timeout 120 >/dev/null; then
record_result "$id" fail "amy did not apply wn's rename (commit after app-data update)"; return
fi
wn_b messages send "$mls_gid" "30 after" >/dev/null 2>&1 || true
if amy_json marmot await message "$gid" --match "30 after" --timeout 90 >/dev/null; then
record_result "$id" pass
else
record_result "$id" fail "amy applied the rename but cannot decrypt wn's next message (forked)"
fi
}
# A reaction White Noise can SEE. Test 09 only proves wn's raw event log holds a
# kind:7; the app renders the materialized timeline, which attaches a reaction
# to its target by its own rules. An amy reaction the raw log kept but the
@@ -98,6 +98,7 @@ import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlin.coroutines.cancellation.CancellationException
import kotlin.io.encoding.Base64
import kotlin.io.encoding.ExperimentalEncodingApi
@@ -277,32 +278,11 @@ class MarmotManager(
}
}
val state =
publishGate.resolve(
obligation.obligationId,
if (confirmed) PublishOutcome.CONFIRMED else PublishOutcome.UNKNOWN,
)
// A confirmed retry makes the commit canonical exactly as
// [commitAndPublish] would, so it owes the same follow-up. Resolving
// the obligation and stopping there was enough to install the state
// and no more: the relay's echo of THIS event was never marked
// processed, so the inbound pipeline met an unknown kind:445 at an
// epoch we had already merged and opened a convergence pass against
// ourselves — a restart could put a healthy group into Recovering
// purely by succeeding.
if (confirmed) {
val framedCommit = framedCommitOf(obligation, event)
if (framedCommit != null) {
inboundProcessor.markMessageProcessed(sha256(framedCommit).toHexKey())
inboundProcessor.recordLocalCommit(
groupId = obligation.groupId,
framedCommitBytes = framedCommit,
sourceEpoch = obligation.priorState.groupContext.epoch,
preState = obligation.priorState,
)
if (confirmed) {
confirmObligation(obligation, event)
} else {
publishGate.resolve(obligation.obligationId, PublishOutcome.UNKNOWN)
}
recordRetentionForCurrentEpoch(obligation.groupId)
syncGroupSystemRows(obligation.groupId, actor = signer.pubKey)
}
Log.d("MarmotManager") {
"retryPendingPublishObligations(): ${obligation.groupId.take(8)}… " +
"confirmed=$confirmed lifecycle=$state"
@@ -310,6 +290,58 @@ class MarmotManager(
}
}
/**
* Make a confirmed obligation's commit canonical, with every follow-up
* [commitAndPublish] owes a commit a relay acknowledged.
*
* Resolving the obligation alone installs the state and no more: the
* relay's echo of THIS event would not be marked processed, so the inbound
* pipeline would meet an unknown kind:445 at an epoch we had already merged
* and open a convergence pass against ourselves.
*/
private suspend fun confirmObligation(
obligation: MarmotPublishObligation,
event: Event,
): GroupLifecycleState {
val state = publishGate.resolve(obligation.obligationId, PublishOutcome.CONFIRMED)
val framedCommit = framedCommitOf(obligation, event)
if (framedCommit != null) {
inboundProcessor.markMessageProcessed(sha256(framedCommit).toHexKey())
inboundProcessor.recordLocalCommit(
groupId = obligation.groupId,
framedCommitBytes = framedCommit,
sourceEpoch = obligation.priorState.groupContext.epoch,
preState = obligation.priorState,
)
}
recordRetentionForCurrentEpoch(obligation.groupId)
syncGroupSystemRows(obligation.groupId, actor = signer.pubKey)
return state
}
/**
* The relay's copy of one of our own unconfirmed commits, or null.
*
* A relay only echoes what it stored, so seeing our pending kind:445 come
* back is the acceptance whose OK never arrived (a timeout, a dropped
* socket, a slow Tor circuit). It must confirm the obligation. Handed to
* the inbound pipeline instead, it reads as a peer's commit whose sender is
* our own leaf; a committer never encrypts the path secret to itself, so it
* fails ("no ciphertext for us") while every peer applies it, and the group
* forks with us one epoch behind.
*/
private suspend fun pendingObligationEchoed(groupEvent: GroupEvent): MarmotPublishObligation? {
val groupId = groupEvent.groupId() ?: return null
return publishGate.pendingFor(groupId).firstOrNull { obligation ->
try {
Event.fromJson(obligation.outboundBytes.decodeToString()).id == groupEvent.id
} catch (e: Exception) {
if (e is CancellationException) throw e
false
}
}
}
/**
* Recover the framed MLS commit from a stored obligation.
*
@@ -391,6 +423,12 @@ class MarmotManager(
* Returns the inner event JSON if it was an application message.
*/
suspend fun processGroupEvent(groupEvent: GroupEvent): GroupEventResult {
pendingObligationEchoed(groupEvent)?.let { obligation ->
confirmObligation(obligation, groupEvent)
subscriptionManager.updateGroupSince(obligation.groupId, groupEvent.createdAt)
return GroupEventResult.CommitProcessed(obligation.groupId, obligation.pendingState.groupContext.epoch)
}
val result = inboundProcessor.processGroupEvent(groupEvent)
// A fork just opened a bounded pass. Inbound traffic settles it
@@ -20,7 +20,9 @@
*/
package com.vitorpamplona.amethyst.commons.marmot
import com.vitorpamplona.quartz.marmot.GroupEventResult
import com.vitorpamplona.quartz.marmot.mip01Groups.MarmotGroupData
import com.vitorpamplona.quartz.marmot.mip03GroupMessages.GroupEvent
import com.vitorpamplona.quartz.marmot.protocolCore.GroupLifecycleState
import com.vitorpamplona.quartz.marmot.protocolCore.LocalOutboundGate
import com.vitorpamplona.quartz.nip01Core.core.Event
@@ -439,6 +441,72 @@ class MarmotPublishBeforeApplyTest {
assertEquals(listOf(relay), fx.manager.groupRelays(fx.groupId))
}
/**
* The relay echo of our own unconfirmed commit is the confirmation that
* never arrived as an OK.
*
* Reproduced against White Noise on a slow (Tor) link: the admin-grant
* publish timed out, so the commit stayed an unresolved obligation and the
* group stayed at its old epoch. The relay HAD stored it, the peer applied
* it, and when the relay echoed it back our inbound pipeline processed it as
* someone else's commit. Its sender is our own leaf, and a committer does not
* encrypt the path secret to itself, so it failed with "UpdatePath at common
* ancestor carries no ciphertext for us" and the group forked: the peer at
* epoch n+1, us still at n, every later message undecryptable.
*/
@Test
fun theRelayEchoOfAnUnconfirmedCommitConfirmsIt() =
runBlocking<Unit> {
val fx = foundedFixture(accepts = false)
val epochBefore =
fx.manager.groupManager
.getGroup(fx.groupId)!!
.epoch
fx.manager.updateGroupMetadata(
fx.groupId,
MarmotGroupData(nostrGroupId = fx.groupId, name = "renamed", adminPubkeys = listOf(fx.manager.signer.pubKey)),
listOf(relay),
)
assertEquals(GroupLifecycleState.PENDING_PUBLISH, fx.manager.lifecycle(fx.groupId))
// What the relay sends back: the same signed kind:445, byte for byte.
val echo =
Event.fromJson(
fx.publisher.published
.single()
.toJson(),
) as GroupEvent
val result = fx.manager.processGroupEvent(echo)
assertTrue(result !is GroupEventResult.Error, "our own commit's echo must not be processed as a foreign commit: $result")
assertEquals(
epochBefore + 1,
fx.manager.groupManager
.getGroup(fx.groupId)!!
.epoch,
"the echo proves a relay took the commit, so it becomes canonical",
)
assertEquals(GroupLifecycleState.STABLE, fx.manager.lifecycle(fx.groupId))
// A second echo (another relay, or a retry) is just a duplicate.
val again =
fx.manager.processGroupEvent(
Event.fromJson(
fx.publisher.published
.single()
.toJson(),
) as GroupEvent,
)
assertTrue(again !is GroupEventResult.Error, "a repeated echo must stay harmless: $again")
assertEquals(
epochBefore + 1,
fx.manager.groupManager
.getGroup(fx.groupId)!!
.epoch,
)
}
private class InMemoryStateStore : com.vitorpamplona.quartz.marmot.groups.MlsGroupStateStore {
private val states = mutableMapOf<String, ByteArray>()
private val retained = mutableMapOf<String, List<ByteArray>>()
@@ -1992,7 +1992,9 @@ class MlsGroup private constructor(
}
check(candidates.isNotEmpty()) {
"UpdatePath at common ancestor carries no ciphertext for us " +
"(my_leaf=$myLeafIndex, my_node=$myNodeIdx, resolution=$resolution, " +
"(sender_leaf=$senderLeafIndex, leaf_count=${tree.leafCount}, filtered_dp=$filteredDp, " +
"filtered_cp=$filteredCp, common_ancestor=$commonAncestorNode, new_leaves=$newLeavesInCommit, " +
"my_leaf=$myLeafIndex, my_node=$myNodeIdx, resolution=$resolution, " +
"held_path_nodes=${pathPrivateKeys.keys.sorted()}, " +
"encrypted_path_secrets=${pathNode.encryptedPathSecret.size})"
}