fix(marmot,relay): end the recovery, gate the composer, extend a retry's deadline

Four follow-ups from the disband work.

**A settled pass no longer reads as `Recovering` forever.** `settle()` cleared
the pass and terminalized a disband but never restored `Stable`, so after any
fork the group reported a recovery that had already finished. The publish gate
is what actually decides whether a commit may be prepared, so nothing locked up
— it was a lie in the reporting, which is worse in its way, since the next
reader to gate on it would have found a group that looked permanently stuck.
`endRecovery` clears only `Recovering`; `Unrecoverable` is not a pass outcome
and `Disbanded` is absorbing.

**The composer is disabled when an outbound gate is up.** Sending into a
disbanding, leaving or removed group throws behind the gate, and the UI found
out by tapping. The gate is mirrored onto `MarmotGroupChatroom` and the chat
view replaces the input with the reason. The gate map is guarded by a mutex and
the refresh path is not suspending, so `MarmotPublishGate` now publishes an
immutable snapshot for non-suspending readers rather than pushing `suspend` up
through every caller for one flag.

**A retry now outlives the deadline it was issued at.** The transport retry
already existed (7e187e39) and was still losing publishes: it shared the
caller's original budget, so the retry was issued, the clock ran out, and the
publish was reported failed having done the work and thrown the answer away —
which is how a healthy loopback relay kept costing the interop harness a
message a run to `disconnected before OK`. Issuing a retry now extends the
deadline by `TRANSPORT_RETRY_GRACE_MS`. Bounded by construction: only a relay
that gave a transport failure earns it, and only as often as the retry budget
allows. The new test fails without the change with exactly the harness's error.

**The three blocked interop directions are closed, with the reason recorded.**
The obvious next idea is to bypass `wn`'s missing verbs through its daemon, and
it does not work: `wnd`'s protocol carries `Ping`, `Status`, `Shutdown`, four
`*Subscribe` variants and `Execute { cli: Box<Cli> }` — the same clap tree `wn`
parses. A verb missing from `Cli` is unreachable through the socket too, so
closing them needs a verb upstream or a driver linked against
`marmot-uniffi`/`marmot-c`. Written down in cli/tests/README.md so nobody
re-investigates.

10,998 tests green.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
Claude
2026-09-11 02:19:43 +00:00
parent 3c65e601df
commit 105d5bd217
11 changed files with 257 additions and 14 deletions
@@ -21,6 +21,7 @@
package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.marmotGroup
import android.widget.Toast
import androidx.compose.foundation.layout.Arrangement
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.Row
import androidx.compose.foundation.layout.Spacer
@@ -43,7 +44,9 @@ import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.platform.LocalContext
import androidx.compose.ui.text.style.TextAlign
import androidx.compose.ui.unit.dp
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import androidx.lifecycle.viewmodel.compose.viewModel
import com.vitorpamplona.amethyst.R
import com.vitorpamplona.amethyst.commons.resources.Res
@@ -72,6 +75,7 @@ import com.vitorpamplona.amethyst.ui.theme.EditFieldModifier
import com.vitorpamplona.amethyst.ui.theme.EditFieldTrailingIconModifier
import com.vitorpamplona.amethyst.ui.theme.SuggestionListDefaultHeightChat
import com.vitorpamplona.amethyst.ui.theme.placeholderText
import com.vitorpamplona.quartz.marmot.protocolCore.LocalOutboundGate
import com.vitorpamplona.quartz.nip01Core.core.HexKey
import kotlinx.collections.immutable.ImmutableList
import kotlinx.collections.immutable.persistentListOf
@@ -98,6 +102,11 @@ fun MarmotGroupChatView(
WatchLifecycleAndUpdateModel(feedViewModel)
val chatroom =
remember(nostrGroupId) {
accountViewModel.account.marmotGroupList.getOrCreateGroup(nostrGroupId)
}
val newMessageModel: MarmotNewMessageViewModel = viewModel(key = nostrGroupId + "MarmotNewMessageViewModel")
newMessageModel.init(accountViewModel)
newMessageModel.load(nostrGroupId)
@@ -157,15 +166,27 @@ fun MarmotGroupChatView(
Spacer(modifier = DoubleVertSpacer)
MarmotGroupMessageComposer(
nostrGroupId = nostrGroupId,
newMessageModel = newMessageModel,
accountViewModel = accountViewModel,
nav = nav,
onMessageSent = {
feedViewModel.feedState.sendToTop()
},
)
// A durable outbound gate means the group takes no new work: an
// unresolved disband request, a SelfRemove already sent, a realized
// removal. Sending would throw behind it, so the composer is replaced
// by the reason rather than left there to fail on tap — the history
// stays readable either way, which is the point of a gate that is not
// a terminal state.
val outboundGate by chatroom.outboundGate.collectAsStateWithLifecycle()
val gate = outboundGate
if (gate != null) {
MarmotGroupClosedComposer(gate)
} else {
MarmotGroupMessageComposer(
nostrGroupId = nostrGroupId,
newMessageModel = newMessageModel,
accountViewModel = accountViewModel,
nav = nav,
onMessageSent = {
feedViewModel.feedState.sendToTop()
},
)
}
}
}
@@ -349,3 +370,32 @@ private fun MarmotGroupFileUploadDialog(
isNip17 = false,
)
}
/**
* Stands in for the composer when an outbound gate is up.
*
* Deliberately a statement rather than a disabled text field: a greyed-out
* input still invites typing, and the three reasons are not the same — one is
* waiting on the group, one on a commit, and one is over. The group's history
* stays on screen above it.
*/
@Composable
private fun MarmotGroupClosedComposer(gate: LocalOutboundGate) {
val message =
when (gate) {
LocalOutboundGate.DISBANDING -> stringRes(R.string.marmot_group_composer_disbanding)
LocalOutboundGate.LEAVING -> stringRes(R.string.marmot_group_composer_leaving)
LocalOutboundGate.REMOVED -> stringRes(R.string.marmot_group_composer_removed)
}
Row(
modifier = EditFieldModifier.fillMaxWidth(),
horizontalArrangement = Arrangement.Center,
) {
Text(
text = message,
color = MaterialTheme.colorScheme.placeholderText,
style = MaterialTheme.typography.bodySmall,
textAlign = TextAlign.Center,
)
}
}
+3
View File
@@ -2701,6 +2701,9 @@
<string name="marmot_failed_to_enable_encrypted_media">Could not switch this group to encrypted attachments: %1$s</string>
<string name="marmot_group_disbanded_toast">Group disbanded</string>
<string name="marmot_group_disbanding_toast">Ending the group. It finishes once the group agrees.</string>
<string name="marmot_group_composer_disbanding">This group is being ended. You can still read it.</string>
<string name="marmot_group_composer_leaving">You are leaving this group.</string>
<string name="marmot_group_composer_removed">You are no longer a member of this group.</string>
<string name="marmot_failed_to_disband">Could not disband the group: %1$s</string>
<string name="marmot_adding_user">Adding %1$s…</string>
<string name="marmot_failed_to_add_user">Failed to add %1$s: %2$s</string>
+10
View File
@@ -111,6 +111,16 @@ The Marmot harnesses come in two flavours, same scenarios:
and uniffi surface (the apps call it) but has no `wn groups` verb, so
test 29 runs one way only.
**The daemon is not a way around this**, which is worth stating because it
is the obvious next idea. `wnd`'s socket protocol
(`crates/cli/src/daemon/protocol.rs`) carries `Ping`, `Status`, `Shutdown`,
four `*Subscribe` variants, and `Execute { cli: Box<Cli> }` — and that last
one takes the same clap command tree `wn` parses. The daemon is a persistent
host for the CLI's verbs, not a richer RPC, so a verb missing from `Cli` is
unreachable through the socket too. Closing these three needs either a verb
upstream in MDK or a driver linked against `marmot-uniffi`/`marmot-c`; both
are out of scope for a harness that deliberately builds MDK unpatched.
A third, slimmer harness covers the NIP-17 DM surface:
- **`dm/dm-interop-headless.sh`** — two `amy` processes (Identity A and
@@ -2539,6 +2539,10 @@ class MarmotManager(
chatroom.isCurrentProfile.value = view.isCurrentProfile
chatroom.hasEncryptedMediaPolicy.value = encryptedMediaPolicy(nostrGroupId) != null
}
// Read every sync, because a gate is raised and cleared by protocol
// events the UI never sees directly — a disband request resolving, a
// removal being realized.
chatroom.outboundGate.value = publishGate.outboundGateNow(nostrGroupId)
val previousCount = chatroom.members.value.size
val members = memberPubkeys(nostrGroupId)
chatroom.members.value = members
@@ -30,6 +30,7 @@ import com.vitorpamplona.amethyst.commons.util.KmpLock
import com.vitorpamplona.amethyst.commons.util.WeakReference
import com.vitorpamplona.amethyst.commons.util.withLock
import com.vitorpamplona.quartz.marmot.appComponents.GroupAvatarUrlV1
import com.vitorpamplona.quartz.marmot.protocolCore.LocalOutboundGate
import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import kotlinx.coroutines.channels.BufferOverflow
@@ -98,6 +99,18 @@ class MarmotGroupChatroom(
*/
var hasEncryptedMediaPolicy = MutableStateFlow(false)
/**
* Why this group takes no new outbound work, or null when it does.
*
* A durable outbound gate is not a lifecycle state: the member is still in
* the tree and the group is not terminal, but nothing new may be sent —
* an unresolved disband request, a SelfRemove already sent, a realized
* removal. A front end needs it separately from [isCurrentProfile] and the
* lifecycle so it can DISABLE the composer with a reason rather than let a
* send throw and surface as an error after the fact.
*/
var outboundGate = MutableStateFlow<LocalOutboundGate?>(null)
var adminPubkeys = MutableStateFlow<List<HexKey>>(emptyList())
var relays = MutableStateFlow<List<String>>(emptyList())
var memberCount = MutableStateFlow(0)
@@ -20,10 +20,12 @@
*/
package com.vitorpamplona.amethyst.commons.marmot
import com.vitorpamplona.amethyst.commons.model.marmotGroups.MarmotGroupChatroom
import com.vitorpamplona.quartz.marmot.appComponents.GroupProfileV1
import com.vitorpamplona.quartz.marmot.mip01Groups.MarmotGroupData
import com.vitorpamplona.quartz.marmot.protocolCore.GroupLifecycleState
import com.vitorpamplona.quartz.marmot.protocolCore.InMemoryPublishObligationStore
import com.vitorpamplona.quartz.marmot.protocolCore.LocalOutboundGate
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import kotlinx.coroutines.runBlocking
@@ -31,6 +33,7 @@ import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertFalse
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
@@ -237,6 +240,26 @@ class MarmotDisbandTest {
Unit
}
@Test
fun `the gate reaches the front end's group state`() =
runBlocking {
// The UI cannot ask a suspending publish gate from the synchronous
// path that refreshes a conversation, so the gate is mirrored onto
// the chatroom. If that mirror is missing, the composer stays
// enabled on a group that refuses every send and the user finds out
// by tapping.
val f = Fixture()
f.createCurrentProfile()
val chatroom = MarmotGroupChatroom(nostrGroupId)
f.manager.syncMetadataTo(nostrGroupId, chatroom)
assertNull(chatroom.outboundGate.value, "a live group has no gate")
f.manager.disbandGroup(nostrGroupId)
f.manager.syncMetadataTo(nostrGroupId, chatroom)
assertEquals(LocalOutboundGate.DISBANDING, chatroom.outboundGate.value)
}
@Test
fun `a pending disband request outlives a restart`() =
runBlocking {
@@ -265,6 +265,27 @@ class MarmotConvergenceEngine(
ctx.lifecycle = pass.lifecycleWhileRunning(ctx.lifecycle)
}
/**
* End the recovery a settled pass was running.
*
* `Recovering` describes a pass IN FLIGHT — a fork being resolved, or a
* disband candidate waiting for selection. Once the pass settles the group
* has a selected branch again and is `Stable`, so leaving the state behind
* makes every later reader believe a recovery is still running: the group
* reads as recovering forever, and anything that gates on the reported
* lifecycle (rather than on the publish gate, which is the authority for
* whether a commit may be prepared) silently stops working.
*
* Only `Recovering` is cleared. `Unrecoverable` is not a pass outcome —
* it means this client cannot safely apply more traffic and is cleared by a
* verified repair — and `Disbanded` is absorbing.
*/
private fun endRecovery(ctx: GroupContext) {
if (ctx.lifecycle == GroupLifecycleState.RECOVERING) {
ctx.lifecycle = GroupLifecycleState.STABLE
}
}
/**
* Move a group to `Disbanded` once its lifecycle component says so.
*
@@ -562,6 +583,7 @@ class MarmotConvergenceEngine(
trimCandidates(ctx)
ctx.divergent.clear()
ctx.pass = null
endRecovery(ctx)
terminalizeIfDisbanded(groupId, ctx)
val epoch = groupManager.getGroup(groupId)?.epoch ?: inputs.tipEpoch
ConvergenceResolution(
@@ -599,6 +621,7 @@ class MarmotConvergenceEngine(
pass.freeze()
ctx.pass = null
endRecovery(ctx)
terminalizeIfDisbanded(groupId, ctx)
ConvergenceResolution(
groupId = groupId,
@@ -27,6 +27,8 @@ 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.sha256.sha256
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
@@ -218,6 +220,24 @@ class MarmotPublishGate(
private val gates = mutableMapOf<HexKey, LocalOutboundGate>()
private val lifecycles = mutableMapOf<HexKey, GroupLifecycleState>()
/**
* An immutable copy of [gates], republished on every change.
*
* The map itself is mutable state guarded by [mutex], so it cannot be read
* from a non-suspending caller without a data race. A front end needs the
* gate on the synchronous path that refreshes a conversation's state — and
* making that path suspend would push `suspend` up through every caller for
* one flag — so the authority stays behind the lock and this is what
* everyone else reads.
*/
private val gateSnapshot = MutableStateFlow<Map<HexKey, LocalOutboundGate>>(emptyMap())
/** Lock-free view of the outbound gates, for non-suspending readers. */
val outboundGates: StateFlow<Map<HexKey, LocalOutboundGate>> get() = gateSnapshot
/** The gate blocking [groupId] right now, read without suspending. */
fun outboundGateNow(groupId: HexKey): LocalOutboundGate? = gateSnapshot.value[groupId]
/** Reload unresolved obligations and outbound gates. Call once at startup, after group restore. */
suspend fun restore() =
mutex.withLock {
@@ -228,6 +248,7 @@ class MarmotPublishGate(
// owner never asked to end it.
LocalOutboundGate.entries.firstOrNull { it.name == name }?.let { gates[groupId] = it }
}
gateSnapshot.value = gates.toMap()
store.loadAll().forEach { bytes ->
try {
val obligation = MarmotPublishObligation.decodeTls(bytes)
@@ -287,7 +308,7 @@ class MarmotPublishGate(
store.saveGate(groupId, gate.name)
mutex.withLock {
gates[groupId] = gate
Unit
gateSnapshot.value = gates.toMap()
}
}
@@ -296,7 +317,7 @@ class MarmotPublishGate(
store.deleteGate(groupId)
mutex.withLock {
gates.remove(groupId)
Unit
gateSnapshot.value = gates.toMap()
}
}
@@ -78,6 +78,23 @@ class PublishResult(
*/
const val DEFAULT_TRANSPORT_RETRIES = 1
/**
* Extra wall-clock granted once per retry that is actually issued.
*
* A retry is not free time: it starts a fresh dial-and-flush cycle at the
* moment the relay hung up, and by then most of the caller's budget is usually
* spent. Sharing the original deadline meant the retry was issued and then
* timed out before the relay could answer — the publish was reported failed
* having done the work and thrown the answer away, which is how a healthy
* loopback relay could still cost a message a run.
*
* Bounded by construction: only a relay that gave a transport failure earns it,
* and only as many times as [DEFAULT_TRANSPORT_RETRIES] allows, so the worst
* case is the caller's timeout plus this much per retried relay rather than an
* open-ended wait on something unreachable.
*/
const val TRANSPORT_RETRY_GRACE_MS = 5_000L
/**
* Internal channel marker for "this relay is back up", so the wait loop — the
* one coroutine that owns the retry bookkeeping — can re-issue a send that a
@@ -207,11 +224,19 @@ suspend fun INostrClient.publishAndCollectResults(
// that is healthy a moment later. The retries live inside the
// caller's existing timeout, so nothing waits longer than before.
val retriesLeft = relayList.associateWith { transportRetries }.toMutableMap()
// The withTimeout block will cancel the coroutine if the loop takes too long
withTimeoutOrNull(timeoutInSeconds * 1000) {
// The deadline is a moving target rather than one
// `withTimeout` around the whole loop: issuing a retry extends it by
// [TRANSPORT_RETRY_GRACE_MS], because the retry needs time the
// original budget has already spent. Everything else about the wait
// is unchanged — a relay that simply never answers still falls out at
// the caller's timeout.
var deadlineMs = timeoutInSeconds * 1000
run {
val awaitingReconnect = mutableSetOf<NormalizedRelayUrl>()
while (receivedResults.size < relayList.size) {
val result = resultChannel.receive()
val remainingMs = deadlineMs - mark.elapsedNow().inWholeMilliseconds
if (remainingMs <= 0) break
val result = withTimeoutOrNull(remainingMs) { resultChannel.receive() } ?: break
if (result.message == RECONNECTED) {
// The pool flushes what it still owes a relay as part of
@@ -242,6 +267,7 @@ suspend fun INostrClient.publishAndCollectResults(
// report the transport failure as before.
receivedResults.remove(result.relay)
awaitingReconnect.add(result.relay)
deadlineMs += TRANSPORT_RETRY_GRACE_MS
Log.d("publishAndConfirm") {
"Retrying ${event.id} on ${result.relay} after ${recorded.message}"
}
@@ -276,6 +276,37 @@ class MarmotConvergenceWiringTest {
assertTrue(obs.inbound.openConvergencePasses().isEmpty())
}
/**
* `Recovering` describes a pass in flight, so a settled pass must leave it.
*
* Leaving the state behind reads as a recovery that never ends: the group
* reports `Recovering` for the rest of the session, and anything that gates
* on the reported lifecycle rather than on the publish gate quietly stops
* working. The publish gate is what actually decides whether a commit may
* be prepared, which is why this was invisible — it is a lie in the
* reporting, not a lock, and only a reader notices.
*/
@Test
fun aSettledPassLeavesRecovering() =
runBlocking<Unit> {
val fork = buildFork()
var now = 0L
val obs = observer(fork.observerStateBytes) { now }
obs.inbound.processGroupEvent(fork.commitA)
obs.inbound.processGroupEvent(fork.commitB)
assertEquals(GroupLifecycleState.RECOVERING, obs.inbound.groupLifecycle(groupId))
now = ConvergencePolicy.V1.settlementQuiescenceMs
val settled = obs.inbound.settleDueConvergence()
assertEquals(1, settled.size)
// The resolution reports it, and so does the engine afterwards —
// a caller that reads either one has to see a group it can use.
assertEquals(GroupLifecycleState.STABLE, settled.single().lifecycle)
assertEquals(GroupLifecycleState.STABLE, obs.inbound.groupLifecycle(groupId))
}
/**
* A plain relay redelivery is still just a duplicate. Convergence only
* claims a commit when a RETAINED state authenticates it, and an echo's
@@ -74,6 +74,12 @@ class PublishRetriesTransportFailureTest {
private class HangsUpOnFirstEvent(
private val delegate: WebsocketBuilder,
private val dropFirstConnections: Int,
/**
* How long the relay takes to come back after it hung up. A real one
* does not reappear instantly, and the reconnect is the part the retry
* has to have time for.
*/
private val reconnectDelayMs: Long = 0,
) : WebsocketBuilder {
val eventFramesSeen = AtomicInteger(0)
private val socketsBuilt = AtomicInteger(0)
@@ -83,6 +89,7 @@ class PublishRetriesTransportFailureTest {
out: WebSocketListener,
): WebSocket {
val index = socketsBuilt.getAndIncrement()
if (index >= dropFirstConnections && reconnectDelayMs > 0) Thread.sleep(reconnectDelayMs)
val inner = delegate.build(url, out)
val hangUp = index < dropFirstConnections
return object : WebSocket by inner {
@@ -127,6 +134,38 @@ class PublishRetriesTransportFailureTest {
client.disconnect()
}
@Test
fun aRetryOutlivesTheCallersOriginalDeadline() =
runBlocking {
// The bug this pins: the retry was issued and then timed out before
// the relay could answer, so the publish did the work and threw the
// answer away. The hang-up costs the whole one-second budget here —
// only the grace granted when the retry is issued can land it.
val builder =
HangsUpOnFirstEvent(
hub,
dropFirstConnections = 1,
reconnectDelayMs = 1_500,
)
val client = NostrClient(builder, scope)
val event = NostrSignerInternal(KeyPair()).sign(TextNoteEvent.build("slow to come back"))
val results =
client.publishAndCollectResults(
event = event,
relayList = setOf(InProcessRelays.DEFAULT_URL),
timeoutInSeconds = 1,
)
val result = results.values.single()
assertTrue(
result.accepted,
"a retry issued at the edge of the budget must be given time to answer " +
"(got \"${result.message}\")",
)
client.disconnect()
}
@Test
fun aRelayThatKeepsHangingUpStillReportsTheTransportFailure() =
runBlocking {