fix(marmot): retry an invite whose first attempt failed instead of waiting for a restart

A Welcome is only processed when its seal is first opened. When a relay
re-delivered the wrap (every re-subscription replays the last two days of gift
wraps), processExistingSeal handed the stored kind-444 rumor to the generic
event processor, which does nothing with it. So a Welcome that failed once
(the signer busy, storage not ready, an exception inside the join) stayed
unjoined until the next app start opened the seal again.

- A replayed seal holding a Welcome goes back through the Welcome flow.
- A Welcome that joined, found us already joined, or names a KeyPackage we never
  held is marked terminally ingested (the policy amy already uses), so replays
  return at once and never re-notify.
- Any other failure is retried on the account scope after 30s, 2m and 10m, with
  at most one pending retry chain per Welcome.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Vitor Pamplona
2026-09-27 23:43:32 -04:00
co-authored by Claude Opus 5.5
parent 4ac9702e3b
commit a2ad02f58c
2 changed files with 59 additions and 2 deletions
@@ -22,6 +22,8 @@ package com.vitorpamplona.amethyst.model
import com.vitorpamplona.amethyst.commons.model.Note
import com.vitorpamplona.amethyst.commons.model.mediaServers.ServerType
import com.vitorpamplona.amethyst.commons.util.KmpLock
import com.vitorpamplona.amethyst.commons.util.withLock
import com.vitorpamplona.quartz.marmot.appComponents.BlobStoreEndpointV2
import com.vitorpamplona.quartz.marmot.appComponents.EncryptedMediaPolicyV2
import com.vitorpamplona.quartz.marmot.appComponents.GroupAvatarUrlV1
@@ -89,6 +91,18 @@ class AccountMarmotActions(
/** Last (timestamp, answer) from [latestKeyPackageOwner], for passive callers. */
private var lastOwnerCheck: Pair<Long, LatestKeyPackageOwner>? = null
// Welcome rumor ids with a retry already scheduled, so a relay re-delivering the same
// wrap while one waits does not start a second chain of retries.
private val welcomeRetries = mutableSetOf<HexKey>()
private val welcomeRetriesLock = KmpLock()
/** True when [welcomeId] had no retry pending and now has one. */
fun claimWelcomeRetry(welcomeId: HexKey): Boolean = welcomeRetriesLock.withLock { welcomeRetries.add(welcomeId) }
fun releaseWelcomeRetry(welcomeId: HexKey) {
welcomeRetriesLock.withLock { welcomeRetries.remove(welcomeId) }
}
/**
* Resolve the relay set for a Marmot group. Prefer the relays carried in
* the MLS GroupContext metadata so every member converges on the same
@@ -55,6 +55,9 @@ import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallRenegotiateEvent
import com.vitorpamplona.quartz.nipC7Chats.ChatEvent
import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
@@ -394,11 +397,15 @@ class GiftWrapEventHandler(
private suspend fun processMarmotWelcomeFlow(
innerEvent: Event,
account: Account,
attempt: Int = 0,
) {
val manager = account.marmotManager ?: return
if (innerEvent !is WelcomeEvent) {
return
}
// A Welcome that already joined (or never can) is done. Replays of its wrap are routine:
// every re-subscription re-delivers the last two days of gift wraps.
if (manager.isTerminallyIngested(innerEvent.id)) return
// "h" tag is optional per MIP-02 — some senders (e.g. whitenoise-rs) omit it.
// nostrGroupId is derived from the MLS GroupContext's NostrGroupData extension instead.
@@ -407,6 +414,7 @@ private suspend fun processMarmotWelcomeFlow(
when (result) {
is WelcomeResult.Joined -> {
manager.markTerminallyIngested(innerEvent.id)
Log.d("MarmotDbg") {
"processMarmotWelcomeFlow: Joined ${result.nostrGroupId.take(8)}… needsKeyPackageRotation=${result.needsKeyPackageRotation}"
}
@@ -437,6 +445,7 @@ private suspend fun processMarmotWelcomeFlow(
}
is WelcomeResult.AlreadyJoined -> {
manager.markTerminallyIngested(innerEvent.id)
// Benign replay of a gift-wrapped Welcome (kind:1059) we already
// processed in a prior session — the relay is just re-delivering
// it after app restart. Log at DEBUG, not WARN.
@@ -446,11 +455,39 @@ private suspend fun processMarmotWelcomeFlow(
}
is WelcomeResult.Error -> {
Log.w("MarmotDbg") { "processMarmotWelcomeFlow: ERROR ${result.message}" }
Log.w("MarmotDbg") { "processMarmotWelcomeFlow: ERROR (attempt ${attempt + 1}) ${result.message}" }
if (result.message.contains("No matching KeyPackageBundle")) {
// Bundles are generated before their KeyPackage is published, so a Welcome
// for one we never held can never become processable.
manager.markTerminallyIngested(innerEvent.id)
} else {
scheduleWelcomeRetry(innerEvent, account, attempt)
}
}
}
}
/**
* Backoff for a Welcome that failed for a reason that can pass (the signer was busy, the
* KeyPackage store was still loading, a relay round trip failed). Before this a failed
* Welcome was only retried on the next app start: a replayed wrap skipped the flow.
*/
private val WELCOME_RETRY_DELAYS_MS = longArrayOf(30_000L, 120_000L, 600_000L)
private fun scheduleWelcomeRetry(
welcome: WelcomeEvent,
account: Account,
attempt: Int,
) {
if (attempt >= WELCOME_RETRY_DELAYS_MS.size) return
if (!account.marmot.claimWelcomeRetry(welcome.id)) return
account.scope.launch(Dispatchers.IO) {
delay(WELCOME_RETRY_DELAYS_MS[attempt])
account.marmot.releaseWelcomeRetry(welcome.id)
processMarmotWelcomeFlow(welcome, account, attempt + 1)
}
}
class SealEventHandler(
private val account: Account,
private val cache: LocalCache,
@@ -529,7 +566,13 @@ class SealEventHandler(
cache.copyRelaysFromTo(publicNote, rumorId)
val innerRumorNote = cache.getOrCreateNote(rumorId)
innerRumorNote.event?.let { innerRumor ->
eventProcessor.consumeEvent(innerRumor, innerRumorNote, publicNote)
// A re-delivered Welcome goes back through the MLS flow, which returns at once
// when it already joined; one that failed the first time gets another chance.
if (MarmotInboundProcessor.isWelcomeEvent(innerRumor)) {
processMarmotWelcomeFlow(innerRumor, account)
} else {
eventProcessor.consumeEvent(innerRumor, innerRumorNote, publicNote)
}
}
}
}