feat(marmot): pin retention to the delivering epoch, and set it from amy

Three related pieces.

**The pinning was approximate.** Expiry was pinned from the group's CURRENT
retention at persist time, which is right for a message delivered under the
current epoch and wrong for one delivered under an older one — a kind:445 held
back as a retained candidate, or replayed after a restart, is decrypted under
an epoch the group has since moved past. Pinning that to today's setting is
precisely what the component forbids. The store now keeps a small epoch →
retention history, written wherever the epoch may have advanced, and the pin
reads the delivering epoch's value. The fallback to the current value is not a
shrug: a group whose setting never changed has one value at every epoch, which
is the overwhelmingly common case. MDK carries `source_epoch` on its rows for
the same reason.

**`amy` can set it.** `marmot group create --disappearing-secs N`, on both
profiles: component `0x8005` for a current-profile group, and the legacy blob
for `--legacy`, where asking for it bumps the version to 3 because v1/v2
deliberately omit the field to stay byte-compatible with MDK's older parser.

**Interop test 26.** Nothing proved MDK accepts a GroupContext that REQUIRES
`0x8005` with our bytes, and retention has the nastiest encoding in the set:
eight big-endian bytes with no length prefix, unlike almost every other Marmot
field, and MIP-01 spelled it differently. amy creates such a group, wn joins
it, and messaging round-trips.

It asserts acceptance rather than read-back, because wn's CLI `group_json`
does not surface the retention value — that field is on the uniffi struct the
apps consume, not this surface. Acceptance is still the encoding test: a
required component whose bytes wn cannot decode leaves the group unreadable,
so `groups show` returning it at all means the eight bytes parsed. The
direction is one-way for the same reason as push — `wn groups` has no command
that sets retention.

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-10 02:16:25 +00:00
parent 99433d9f5c
commit 37cd798e99
11 changed files with 319 additions and 9 deletions
@@ -112,7 +112,7 @@ class AndroidMarmotMessageStore(
override suspend fun delete(nostrGroupId: String) {
withContext(Dispatchers.IO) {
writeMutex.withLock {
for (file in listOf(messagesFile(nostrGroupId), epochsFile(nostrGroupId), snapshotFile(nostrGroupId), expiriesFile(nostrGroupId))) {
for (file in listOf(messagesFile(nostrGroupId), epochsFile(nostrGroupId), snapshotFile(nostrGroupId), expiriesFile(nostrGroupId), epochRetentionsFile(nostrGroupId))) {
if (file.exists() && !file.delete()) {
Log.w(TAG) { "delete($nostrGroupId): failed to remove ${file.absolutePath}" }
}
@@ -243,6 +243,49 @@ class AndroidMarmotMessageStore(
}
}
private fun epochRetentionsFile(nostrGroupId: String): File = File(groupDir(nostrGroupId), "epoch_retentions")
/**
* What retention this group required at each epoch.
*
* First write wins per epoch: an epoch's required components are fixed the
* moment it exists, so a second answer for the same epoch would be a bug
* rather than an update.
*/
override suspend fun recordEpochRetention(
nostrGroupId: String,
epoch: Long,
retentionSecs: Long,
) = withContext(Dispatchers.IO) {
writeMutex.withLock {
try {
val existing = readAllFrom(epochRetentionsFile(nostrGroupId)).toMutableList()
if (existing.any { it.substringBefore(' ') == epoch.toString() }) return@withLock
existing.add("$epoch $retentionSecs")
writeAllTo(epochRetentionsFile(nostrGroupId), existing)
} catch (e: Exception) {
Log.e(TAG, "recordEpochRetention($nostrGroupId) FAILED: ${e.message}", e)
}
}
}
override suspend fun loadEpochRetentions(nostrGroupId: String): Map<Long, Long> =
withContext(Dispatchers.IO) {
try {
readAllFrom(epochRetentionsFile(nostrGroupId))
.mapNotNull { line ->
val parts = line.trim().split(' ')
if (parts.size != 2) return@mapNotNull null
val epoch = parts[0].toLongOrNull() ?: return@mapNotNull null
val secs = parts[1].toLongOrNull() ?: return@mapNotNull null
epoch to secs
}.toMap()
} catch (e: Exception) {
Log.e(TAG, "loadEpochRetentions($nostrGroupId) FAILED: ${e.message}", e)
emptyMap()
}
}
private fun snapshotFile(nostrGroupId: String): File = File(groupDir(nostrGroupId), "snapshot")
/**
@@ -776,6 +776,7 @@ class GroupEventHandler(
// row is derived from. Deriving here covers OTHER members'
// commits; our own are derived by `commitAndPublish`. Both
// reach the feed through `onSystemRowDerived`.
manager.recordRetentionForCurrentEpoch(result.groupId)
manager.syncGroupSystemRows(result.groupId)
// Epoch just advanced — drain any kind:445 events that
// previously failed as UndecryptableOuterLayer for this
@@ -25,6 +25,7 @@ import com.vitorpamplona.amethyst.cli.Context
import com.vitorpamplona.amethyst.cli.DataDir
import com.vitorpamplona.amethyst.cli.Output
import com.vitorpamplona.quartz.marmot.appComponents.GroupProfileV1
import com.vitorpamplona.quartz.marmot.appComponents.MessageRetentionV1
import com.vitorpamplona.quartz.marmot.mip01Groups.MarmotGroupData
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import com.vitorpamplona.quartz.utils.RandomInstance
@@ -37,6 +38,21 @@ object GroupCreateCommand {
val args = Args(rest)
val name = args.flag("name", "")!!
val legacy = args.bool("legacy")
// Disappearing messages (`marmot.group.message-retention.v1`, 0x8005).
// Fixed at creation: promoting a component to required later needs its
// state installed first, which is a second commit this command does not
// make.
val disappearing = args.flag("disappearing-secs", "")!!
val disappearingSecs =
if (disappearing.isEmpty()) {
null
} else {
disappearing.toULongOrNull()?.takeIf { it > 0uL }
?: return Output.error(
"bad_args",
"--disappearing-secs must be a positive whole number of seconds",
)
}
args.rejectUnknown()
Context.open(dataDir).use { ctx ->
ctx.prepare()
@@ -53,13 +69,23 @@ object GroupCreateCommand {
// so later invitees receive a pre-populated group from the
// welcome and never have to chase an undecryptable bootstrap
// commit that predates their membership.
// The legacy blob only carries `disappearing_message_secs`
// from v3 on, so asking for it bumps the version — v1/v2 stays
// byte-for-byte what MDK's older parser accepts.
val metadata =
MarmotGroupData.bootstrap(
nostrGroupId = gid,
creatorPubKey = ctx.identity.pubKeyHex,
outboxRelays = outboxUrls,
name = name,
)
MarmotGroupData
.bootstrap(
nostrGroupId = gid,
creatorPubKey = ctx.identity.pubKeyHex,
outboxRelays = outboxUrls,
name = name,
).let {
if (disappearingSecs == null) {
it
} else {
it.copy(disappearingMessageSecs = disappearingSecs, version = 3)
}
}
ctx.marmot.createGroup(gid, initialMetadata = metadata)
} else {
// Current profile by default. The difference is what the group
@@ -72,6 +98,7 @@ object GroupCreateCommand {
nostrGroupId = gid,
relays = outboxUrls,
profile = if (name.isEmpty()) null else GroupProfileV1(name, ""),
retention = disappearingSecs?.let { MessageRetentionV1(it) },
)
}
@@ -155,6 +155,7 @@ class FileMarmotMessageStore(
epochFile(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group message epochs")
snapshotFile(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group system-row baseline")
expiryFile(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group message expiries")
epochRetentionFile(nostrGroupId).deleteOrWarn("FileMarmotMessageStore", "group epoch retentions")
}
private fun snapshotFile(id: String) = File(dir, "$id.snapshot")
@@ -219,6 +220,32 @@ class FileMarmotMessageStore(
}
}
private fun epochRetentionFile(id: String) = File(dir, "$id.epoch-retentions")
/** First write wins: an epoch's required components are fixed once it exists. */
override suspend fun recordEpochRetention(
nostrGroupId: String,
epoch: Long,
retentionSecs: Long,
) {
val target = epochRetentionFile(nostrGroupId)
if (target.exists() && target.readLines().any { it.substringBefore(' ') == epoch.toString() }) return
SecureFileIO.appendText(target, "$epoch $retentionSecs\n")
}
override suspend fun loadEpochRetentions(nostrGroupId: String): Map<Long, Long> =
epochRetentionFile(nostrGroupId)
.takeIf { it.exists() }
?.readLines()
?.mapNotNull { line ->
val parts = line.trim().split(' ')
if (parts.size != 2) return@mapNotNull null
val epoch = parts[0].toLongOrNull() ?: return@mapNotNull null
val secs = parts[1].toLongOrNull() ?: return@mapNotNull null
epoch to secs
}?.toMap()
?: emptyMap()
private fun epochFile(id: String) = File(dir, "$id.epochs")
override suspend fun recordEpoch(
@@ -207,6 +207,7 @@ ALL_TESTS=(
test_23_deletion_amy_to_wn
test_24_media_v2_amy_to_wn
test_25_media_v2_wn_to_amy
test_26_retention_amy_to_wn
)
# --tests runs a subset in the order given. Most tests read state a previous
+62
View File
@@ -342,3 +342,65 @@ test_25_media_v2_wn_to_amy() {
record_result "$id" fail "amy decrypted different bytes than wn sent"
fi
}
# --- 26: disappearing messages ------------------------------------------------
# `marmot.group.message-retention.v1` (0x8005) has the nastiest encoding in the
# component set: eight big-endian bytes with NO length prefix, unlike almost
# every other Marmot field, and MIP-01 spelled it differently. Nothing else
# proves MDK accepts a GroupContext that REQUIRES it with our bytes — and if it
# does not, the failure is not cosmetic: wn cannot read the group at all.
#
# One-directional on purpose. `wn` has no command that sets retention
# (`wn groups` is list/create/show/add-members/remove-members/members/admins/
# relays/leave/rename/set-avatar-url), so the reverse direction is untestable
# here. The encode side is the one that can be wrong anyway. What amy DOES with
# the expiry once it holds one is unit-tested in `MarmotRetentionTest`; this is
# purely "does the other implementation accept and read what we wrote".
test_26_retention_amy_to_wn() {
banner "Test 26 — amy creates a group with disappearing messages; wn reads the policy"
local id="26 retention amy->wn"
local want=3600
local out gid mls_gid
out=$(amy_json marmot group create --name "Interop-Retention" --disappearing-secs "$want") || {
record_result "$id" fail "amy group create --disappearing-secs failed"; return
}
gid=$(printf '%s' "$out" | jq -r '.group_id')
mls_gid=$(printf '%s' "$out" | jq -r '.mls_group_id')
if [[ -z "$gid" || "$gid" == "null" ]]; then
record_result "$id" fail "amy reported no group id"; return
fi
amy_json marmot group add "$gid" "$B_NPUB" >/dev/null || {
record_result "$id" fail "amy could not invite wn"; return
}
# The Welcome is the real assertion: a group requiring a component wn cannot
# decode is a group wn refuses to join.
local b_gid
b_gid=$(wait_for_invite B 60) || {
record_result "$id" fail "wn never received a Welcome for a group requiring 0x8005"; return
}
wn_b groups accept "$b_gid" >/dev/null 2>&1 || true
# wn's CLI `group_json` does not surface the retention value — it is on the
# uniffi group struct the apps consume, not this surface — so the assertion
# is acceptance rather than read-back. That is still the encoding test: the
# group REQUIRES 0x8005, and a required component whose bytes wn cannot
# decode makes the group unreadable, so `groups show` returning it at all
# means our eight big-endian bytes parsed.
if ! wn_group_field_becomes "$mls_gid" '.group.group_id // empty' "$mls_gid" 120; then
record_result "$id" fail "wn never surfaced a group that requires 0x8005"; return
fi
# And the group still works: a required component that decodes but breaks
# messaging would pass the check above and still be useless.
wn_b messages send "$mls_gid" "retention round trip" >/dev/null 2>&1 || {
record_result "$id" fail "wn could not send into the retention group"; return
}
if amy_json marmot await message "$gid" --match "retention round trip" --timeout 90 >/dev/null; then
record_result "$id" pass
else
record_result "$id" fail "amy never received wn's message in the retention group"
fi
}
@@ -143,6 +143,7 @@ private suspend fun MarmotManager.ingestGiftWrapUncached(wrap: GiftWrapEvent): M
// writing any: a joiner announcing every existing member as
// newly added would be a timeline full of events that never
// happened.
recordRetentionForCurrentEpoch(result.nostrGroupId)
syncGroupSystemRows(result.nostrGroupId)
MarmotIngestResult.JoinedGroup(
nostrGroupId = result.nostrGroupId,
@@ -180,6 +181,7 @@ private suspend fun MarmotManager.ingestGroupEvent(ge: GroupEvent): MarmotIngest
// kind:1210 row is derived from. Deriving here rather than at
// render time means the rows land in the same log as the messages
// they sit between, in the order they happened.
recordRetentionForCurrentEpoch(result.groupId)
syncGroupSystemRows(result.groupId)
MarmotIngestResult.Commit(result)
}
@@ -759,6 +759,7 @@ class MarmotManager(
// immediately satisfied. Every LATER commit takes the normal
// publish-before-apply path.
publishGate.satisfyEmptyObligation(nostrGroupId)
recordRetentionForCurrentEpoch(nostrGroupId)
inboundProcessor.trackGroup(nostrGroupId)
subscriptionManager.subscribeGroup(nostrGroupId)
Log.d("MarmotManager") { "createGroup($nostrGroupId): persisted and subscribed" }
@@ -800,6 +801,7 @@ class MarmotManager(
// Same empty-obligation exception as [createGroup]: a one-member
// epoch-0 group has no peer that failure to publish could fork.
publishGate.satisfyEmptyObligation(nostrGroupId)
recordRetentionForCurrentEpoch(nostrGroupId)
inboundProcessor.trackGroup(nostrGroupId)
subscriptionManager.subscribeGroup(nostrGroupId)
return nostrGroupId
@@ -891,6 +893,7 @@ class MarmotManager(
// derived rows a peer's commit would. The actor is known here in a
// way it is not for an inbound commit, which is what lets a
// self-removal read as "left" rather than "removed".
recordRetentionForCurrentEpoch(nostrGroupId)
syncGroupSystemRows(nostrGroupId, actor = signer.pubKey)
} else {
Log.w("MarmotManager") {
@@ -1101,7 +1104,7 @@ class MarmotManager(
// Pinned here, at the moment the message enters the log, because
// this is the last point at which the retention of its delivering
// epoch is still the group's current retention.
parsed?.let { pinExpiry(nostrGroupId, it) }
parsed?.let { pinExpiry(nostrGroupId, it, epoch) }
} catch (e: Exception) {
Log.w("MarmotManager", "Failed to persist Marmot message for $nostrGroupId", e)
}
@@ -1285,8 +1288,9 @@ class MarmotManager(
private suspend fun pinExpiry(
nostrGroupId: HexKey,
innerEvent: Event,
sourceEpoch: Long?,
) {
val seconds = retentionSeconds(nostrGroupId)
val seconds = retentionForEpoch(nostrGroupId, sourceEpoch)
if (seconds <= 0L) return
try {
messageStore?.recordExpiry(nostrGroupId, innerEvent.id, innerEvent.createdAt + seconds)
@@ -1295,6 +1299,48 @@ class MarmotManager(
}
}
/**
* The retention that applied at [sourceEpoch] — the epoch that DELIVERED
* the message — falling back to the group's current value.
*
* The distinction only shows up on a message decrypted late: a retained
* candidate, or a replay after a restart, arrives under an epoch the group
* has since moved past. Pinning it to today's setting is exactly what the
* component forbids, so the recorded history wins whenever it has an entry
* for that epoch. The fallback is not a shrug — a group whose setting never
* changed has one value at every epoch, and that is the overwhelmingly
* common case.
*/
private suspend fun retentionForEpoch(
nostrGroupId: HexKey,
sourceEpoch: Long?,
): Long {
if (sourceEpoch != null) {
try {
messageStore?.loadEpochRetentions(nostrGroupId)?.get(sourceEpoch)?.let { return it }
} catch (e: Exception) {
Log.w("MarmotManager", "Failed to read epoch retentions for $nostrGroupId", e)
}
}
return retentionSeconds(nostrGroupId)
}
/**
* Write down what this group's retention is at its current epoch.
*
* Called wherever the epoch may just have advanced, so the history has an
* entry before any message delivered under that epoch needs one.
*/
suspend fun recordRetentionForCurrentEpoch(nostrGroupId: HexKey) {
val store = messageStore ?: return
val epoch = currentEpoch(nostrGroupId) ?: return
try {
store.recordEpochRetention(nostrGroupId, epoch, retentionSeconds(nostrGroupId))
} catch (e: Exception) {
Log.w("MarmotManager", "Failed to record retention at epoch $epoch for $nostrGroupId", e)
}
}
/**
* Delete every message whose pinned expiry has passed.
*
@@ -30,6 +30,7 @@ import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
/**
@@ -235,4 +236,70 @@ class MarmotRetentionTest {
f.manager.pruneExpiredMessages(nostrGroupId, sent.innerEvent.createdAt + 121),
)
}
@Test
fun `a message delivered under an older epoch keeps that epoch's retention`() =
runBlocking {
// The case the fallback used to get wrong. A kind:445 held back as
// a retained candidate is decrypted under an epoch the group has
// since moved past; pinning it to today's setting is exactly what
// the component forbids.
val f = Fixture()
f.createGroup(60uL)
val epochZero = assertNotNull(f.manager.currentEpoch(nostrGroupId))
f.manager.updateGroupMetadata(
nostrGroupId,
MarmotGroupData(
nostrGroupId = nostrGroupId,
name = "retention",
relays = listOf("wss://relay.invalid"),
disappearingMessageSecs = 86_400uL,
version = 3,
),
)
assertEquals(86_400L, f.manager.retentionSeconds(nostrGroupId))
assertTrue(f.manager.currentEpoch(nostrGroupId)!! > epochZero, "the rename must advance the epoch")
// Arrives now, but was delivered by the epoch that still said 60s.
val late =
Event(
id = "2".repeat(64),
pubKey = f.signer.pubKey,
createdAt = TimeUtils.now(),
kind = 9,
tags = emptyArray(),
content = "decrypted late",
sig = "",
)
f.manager.persistDecryptedMessage(nostrGroupId, late.toJson(), epoch = epochZero)
assertEquals(
late.createdAt + 60,
f.messageStore.loadExpiries(nostrGroupId)[late.id],
"a late message must keep its source epoch's retention, not the current one",
)
}
@Test
fun `an unknown epoch falls back to the current retention`() =
runBlocking {
// Not a shrug: a group whose setting never changed has one value at
// every epoch, and that is the overwhelmingly common case.
val f = Fixture()
f.createGroup(60uL)
val orphan =
Event(
id = "3".repeat(64),
pubKey = f.signer.pubKey,
createdAt = TimeUtils.now(),
kind = 9,
tags = emptyArray(),
content = "no history for this epoch",
sig = "",
)
f.manager.persistDecryptedMessage(nostrGroupId, orphan.toJson(), epoch = 9_999L)
assertEquals(orphan.createdAt + 60, f.messageStore.loadExpiries(nostrGroupId)[orphan.id])
}
}
@@ -69,6 +69,7 @@ class SnapshotMessageStore : MarmotMessageStore {
private val messages = mutableMapOf<String, MutableList<String>>()
private val snapshots = mutableMapOf<String, String>()
private val expiries = mutableMapOf<String, MutableMap<String, Long>>()
private val epochRetentions = mutableMapOf<String, MutableMap<Long, Long>>()
override suspend fun appendMessage(
nostrGroupId: String,
@@ -84,6 +85,7 @@ class SnapshotMessageStore : MarmotMessageStore {
messages.remove(nostrGroupId)
snapshots.remove(nostrGroupId)
expiries.remove(nostrGroupId)
epochRetentions.remove(nostrGroupId)
}
override suspend fun recordGroupSnapshot(
@@ -115,6 +117,16 @@ class SnapshotMessageStore : MarmotMessageStore {
messages[nostrGroupId]?.removeAll { json -> Event.fromJsonOrNull(json)?.id in innerEventIds }
expiries[nostrGroupId]?.keys?.removeAll(innerEventIds)
}
override suspend fun recordEpochRetention(
nostrGroupId: String,
epoch: Long,
retentionSecs: Long,
) {
epochRetentions.getOrPut(nostrGroupId) { mutableMapOf() }.putIfAbsent(epoch, retentionSecs)
}
override suspend fun loadEpochRetentions(nostrGroupId: String): Map<Long, Long> = epochRetentions[nostrGroupId]?.toMap() ?: emptyMap()
}
class SnapshotBundleStore : KeyPackageBundleStore {
@@ -139,6 +139,28 @@ interface MarmotMessageStore {
/** Inner event id → its pinned expiry, for what was recorded. */
suspend fun loadExpiries(nostrGroupId: String): Map<String, Long> = emptyMap()
/**
* Remember the retention this group required AT [epoch].
*
* Kept because a message pins the retention of its own source epoch, and
* the source epoch is not always the current one: a kind:445 held back as a
* retained candidate, or replayed after a restart, is decrypted under an
* epoch the group has since moved past. Without this history such a message
* would be pinned to whatever the setting happens to be at decrypt time —
* which is the one thing the component says must not happen.
*
* Small and append-only: one entry per epoch that changed it, not one per
* message.
*/
suspend fun recordEpochRetention(
nostrGroupId: String,
epoch: Long,
retentionSecs: Long,
) = Unit
/** MLS epoch → the retention it required, for what was recorded. */
suspend fun loadEpochRetentions(nostrGroupId: String): Map<Long, Long> = emptyMap()
/**
* Delete these messages, and any expiry recorded for them, permanently.
*