feat(cordn): persist messages, the only copy there will ever be

A cordn room's history lived in a LinkedHashMap and nowhere else, while
the cursor that fetches it was persisted pointing past everything. So a
relaunch opened every room empty and no amount of syncing refilled it:
the one handle to history had already been spent. Rewinding it would not
have helped either — both seal keys are epoch-derived, so the copy the
coordinator still holds cannot be opened once the group has ratcheted.
Ingestion is the only moment the message is in the clear.

**EncryptedAppendLog moves to a neutral package.** It takes encrypt and
decrypt lambdas and stores opaque strings; nothing in it knows what
Marmot is. It only lived under `marmot` because that is where it was
written first, and the independence guard reads that as a dependency, so
cordn could not touch it. Now `commons.storage`, used by both, owned by
neither — which makes the two features more separable, not less. Marmot
changes by one import line.

**The write.** `CordnDeliveredMessageCodec` stores the envelope plus its
cursor: the cursor is the coordinator's sequence, which the room orders
on and the unread divider compares against, and `created_at` is the
sender's clock, which is a claim. `FileCordnGroupStore` appends a segment
per message rather than rewriting the conversation, so receiving costs
the same on message ten thousand as on message one.

**The ordering is the correctness property.** `persist` writes messages
*before* the cursor. Crash between them and the message is on disk while
the cursor still points before it: the next catch-up re-delivers it and
dedup drops it. The other order loses it permanently. That is why the
write lives next to the cursor write rather than at the call sites.

Dedup keys on the envelope id, not the stored entry: a re-delivery
carries the same envelope with a different cursor, so the bytes differ
while naming one message.

**The read, in two parts.** The conversation loads when a room opens,
through `addAll`, so the annotation fold and ordering are rebuilt by the
same code that handles live delivery. A one-entry summary per group
loads at login for the inbox line — which is why every cordn room read
"No messages yet" after a relaunch however much had been said.

Deleting a group takes the log and the summary with it. History that
outlived its MLS state would be the plaintext of an end-to-end encrypted
conversation left behind after the key that read it was destroyed.

Mutation-checked: inverting the write order, dedup'ing on the whole entry
instead of the id, and forgetting to drop history on delete each fail
exactly one test and nothing else.

Not yet carried into backup or device migration — the plan
(commons/plans/2026-09-24-cordn-message-store.md) recommends migration
yes, backup no, and that is still an open call.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
Claude
2026-09-24 15:08:10 +00:00
parent 1d64e65a95
commit e4dc1c52d4
11 changed files with 562 additions and 13 deletions
@@ -181,7 +181,13 @@ class CordnRuntime(
// opened: an inbox that shows every room as unread until it is
// visited is worse than one with no unread state at all.
val saved = session.manager.roomState(it)
groups.get(config.pubKey, it)?.restoreState(saved.draft, saved.lastReadCursor)
val room = groups.get(config.pubKey, it)
room?.restoreState(saved.draft, saved.lastReadCursor)
// The last message, for the inbox line. One small key per group
// rather than the whole log: the conversation itself loads when
// a room is opened. Without it every cordn room read "No
// messages yet" after a relaunch, however much had been said.
session.manager.storedMessageSummary(it)?.let { summary -> room?.restorePreview(summary.newest) }
}
remember()
maintainKeyPackages(session)
@@ -725,6 +731,11 @@ class CordnRuntime(
val room = groups.get(coordinatorPubKey, gid) ?: return
val saved = session.manager.roomState(gid)
room.restoreState(saved.draft, saved.lastReadCursor)
// The conversation itself. Through addAll, so the annotation fold, the
// ordering and `newest` are all rebuilt by exactly the code that
// handles live delivery — a room restored down a second path would be
// a second set of rules to keep in step.
room.addAll(session.manager.storedMessages(gid))
}
/** Persists [gid]'s draft and read position. */
@@ -20,7 +20,7 @@
*/
package com.vitorpamplona.amethyst.model.marmot
import com.vitorpamplona.amethyst.commons.marmot.EncryptedAppendLog
import com.vitorpamplona.amethyst.commons.storage.EncryptedAppendLog
import com.vitorpamplona.amethyst.model.preferences.KeyStoreEncryption
import com.vitorpamplona.quartz.marmot.groups.MarmotMessageStore
import com.vitorpamplona.quartz.nip01Core.core.Event
@@ -583,6 +583,10 @@ class CordnGroupManager(
)
val sealed = CordnApplicationMessage.seal(group, accountPubKey, envelope)
val posted = call { sync.postMessage(gid, sealed) }
// Our own message never arrives as a Delivery.Message — it comes back
// as an Echo, which by design carries nothing — so the only place it
// can be recorded is here, where we still hold the plaintext.
queueForStore(gid, CordnDeliveredMessage(envelope, posted.cursor))
persist(gid)
return CordnDeliveredMessage(envelope, posted.cursor)
}
@@ -659,8 +663,11 @@ class CordnGroupManager(
else -> {
val decrypted = group.decrypt(opened)
when (decrypted.contentType) {
ContentType.APPLICATION ->
Delivery.Message(gid, ingestion.cursor, CordnApplicationMessage.open(decrypted))
ContentType.APPLICATION -> {
val received = CordnApplicationMessage.open(decrypted)
queueForStore(gid, CordnDeliveredMessage(received.envelope, ingestion.cursor))
Delivery.Message(gid, ingestion.cursor, received)
}
// decrypt() applies a Commit, so the epoch has moved already.
else -> Delivery.EpochAdvanced(gid, ingestion.cursor, group.epoch)
}
@@ -747,6 +754,20 @@ class CordnGroupManager(
private suspend fun persist(gid: String) {
val group = groups[gid] ?: return
store.saveGroup(gid, group.saveState().encodeTls())
// Messages BEFORE the cursor, and deliberately so. Crash between the
// two and the message is on disk while the cursor still points before
// it: the next catch-up re-delivers it and appendMessage's dedup drops
// it, costing nothing. The other order loses the message permanently —
// both cordn seal keys are epoch-derived, so the copy the coordinator
// still holds can no longer be opened, and the cursor would not offer
// it again anyway.
//
// Living here rather than at the call sites is the point: the ordering
// is the correctness property, so it sits next to the cursor write
// where it cannot be forgotten.
pendingMessages.remove(gid)?.forEach { store.appendMessage(gid, it) }
sync.cursors()[gid]?.let { store.saveCursor(gid, it) }
// On the same beat as the cursor, because the two are only meaningful
// together: a cursor that outlives the process while the record of
@@ -755,6 +776,31 @@ class CordnGroupManager(
sync.echoes()[gid]?.let { store.saveEchoState(gid, it) }
}
/**
* Messages delivered but not yet written, per group.
*
* [ingest] is not a suspend function — it runs inside the sync's own
* delivery callback — so it cannot write. It queues here instead, and
* [persist] drains the queue immediately before the cursor. Held in memory
* only for as long as the cursor is: neither is on disk until a run ends,
* so a crash mid-run loses both together and the next catch-up refetches
* from the last cursor that *was* written.
*/
private val pendingMessages = mutableMapOf<String, MutableList<CordnDeliveredMessage>>()
private fun queueForStore(
gid: String,
message: CordnDeliveredMessage,
) {
pendingMessages.getOrPut(gid) { mutableListOf() }.add(message)
}
/** Every message this group has stored, oldest first. */
suspend fun storedMessages(gid: String): List<CordnDeliveredMessage> = store.loadMessages(gid)
/** The newest stored message and the count, without reading the log. */
suspend fun storedMessageSummary(gid: String): CordnMessageSummary? = store.loadMessageSummary(gid)
private suspend fun persistAll() {
groups.keys.forEach { persist(it) }
}
@@ -20,6 +20,8 @@
*/
package com.vitorpamplona.amethyst.commons.cordn
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessage
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessageCodec
import com.vitorpamplona.quartz.cordn.sync.EchoState
import com.vitorpamplona.quartz.cordn.sync.GroupCursor
import com.vitorpamplona.quartz.cordn.sync.PendingEpochOperation
@@ -116,8 +118,53 @@ interface CordnGroupStore {
/** The saved room state for [gid], or a blank one. */
suspend fun loadRoomState(gid: String): CordnRoomState
/**
* Records one delivered message, as the only copy there will ever be.
*
* Not a cache. Both of a cordn payload's seal keys are derived per epoch,
* so once the group ratchets forward the ciphertext the coordinator still
* holds cannot be opened again — and the cursor has advanced past it in any
* case, so it will not even be offered. Ingestion is the one moment the
* message is in the clear.
*
* MUST be idempotent on the envelope's id. A client that dies between this
* call and [saveCursor] re-fetches the message on the next catch-up, which
* is the crash window the append order is chosen to land in — see
* `commons/plans/2026-09-24-cordn-message-store.md`.
*/
suspend fun appendMessage(
gid: String,
message: CordnDeliveredMessage,
)
/** Every stored message for [gid], oldest first. Empty when there are none. */
suspend fun loadMessages(gid: String): List<CordnDeliveredMessage>
/**
* The newest message and the count, without reading the whole log.
*
* The inbox needs a preview line for every room before any room is opened,
* and reading every group's history at login is what
* `CordnRuntime.restoreRoomState` deliberately avoids. This is the small
* read that makes the preview possible: one key per group, written on the
* same beat as [appendMessage].
*/
suspend fun loadMessageSummary(gid: String): CordnMessageSummary?
}
/**
* What the inbox needs to know about a room without opening it.
*
* @param newest the last message delivered, for the preview line.
* @param count how many are stored, for anything that wants to say "empty" and
* mean it rather than meaning "not loaded yet".
*/
data class CordnMessageSummary(
val newest: CordnDeliveredMessage,
val count: Int,
)
/**
* What a cordn room remembers between visits.
*
@@ -134,6 +181,43 @@ data class CordnRoomState(
val isBlank: Boolean get() = draft.isEmpty() && lastReadCursor == 0L
}
/**
* The on-disk form of a [CordnMessageSummary].
*
* The message itself rides as its own JSON entry rather than being taken apart
* into TLS fields, so there is exactly one definition of what a stored message
* looks like and the summary cannot drift from the log it summarises.
*/
object CordnMessageSummaryCodec {
const val VERSION = 1
fun encode(
newest: CordnDeliveredMessage,
count: Int,
): ByteArray {
val writer = TlsWriter()
writer.putUint16(VERSION)
writer.putUint32(count.toLong())
writer.putOpaque2(CordnDeliveredMessageCodec.encode(newest).encodeToByteArray())
return writer.toByteArray()
}
/** Null on anything unreadable: a preview line is not worth failing a login over. */
fun decode(bytes: ByteArray): CordnMessageSummary? =
try {
val reader = TlsReader(bytes)
if (reader.readUint16() != VERSION) {
null
} else {
val count = reader.readUint32().toInt()
val newest = CordnDeliveredMessageCodec.decode(reader.readOpaque2().decodeToString())
CordnMessageSummary(newest, count)
}
} catch (e: Exception) {
null
}
}
/** The on-disk layout of a [CordnRoomState]. */
object CordnRoomStateCodec {
const val VERSION = 1
@@ -228,6 +312,7 @@ class InMemoryCordnGroupStore : CordnGroupStore {
private val viaRequest = mutableSetOf<String>()
private val roomStates = mutableMapOf<String, CordnRoomState>()
private val echoStates = mutableMapOf<String, EchoState>()
private val messages = mutableMapOf<String, MutableList<CordnDeliveredMessage>>()
override suspend fun saveGroup(
gid: String,
@@ -244,6 +329,10 @@ class InMemoryCordnGroupStore : CordnGroupStore {
viaRequest.remove(gid)
roomStates.remove(gid)
echoStates.remove(gid)
// Leaving a group that left its history behind would keep the plaintext
// of an end-to-end encrypted conversation on disk after the one thing
// that could read it was thrown away.
messages.remove(gid)
}
override suspend fun listGroups(): List<String> = groups.keys.toList()
@@ -263,6 +352,18 @@ class InMemoryCordnGroupStore : CordnGroupStore {
override suspend fun loadRoomState(gid: String): CordnRoomState = roomStates[gid] ?: CordnRoomState()
override suspend fun appendMessage(
gid: String,
message: CordnDeliveredMessage,
) {
val log = messages.getOrPut(gid) { mutableListOf() }
if (log.none { it.envelope.id == message.envelope.id }) log.add(message)
}
override suspend fun loadMessages(gid: String): List<CordnDeliveredMessage> = messages[gid]?.toList() ?: emptyList()
override suspend fun loadMessageSummary(gid: String): CordnMessageSummary? = messages[gid]?.lastOrNull()?.let { CordnMessageSummary(it, messages[gid]?.size ?: 0) }
override suspend fun saveEchoState(
gid: String,
state: EchoState,
@@ -172,6 +172,21 @@ class CordnGroupChatroom(
recountUnread()
}
/**
* Seeds the inbox's preview line without loading the conversation.
*
* The inbox needs a last message for every room before any room is opened,
* and reading every group's whole history at login to get one would be
* paying for the rooms nobody visits. The store keeps a one-entry summary
* for exactly this.
*
* Ignored once the room holds anything, so a summary read cannot overwrite
* a newer message that live delivery already put here.
*/
fun restorePreview(newest: CordnDeliveredMessage) {
if (byId.isEmpty()) _newest.value = newest
}
/** Restores what the store remembered for this room. */
fun restoreState(
draft: String,
@@ -20,11 +20,16 @@
*/
package com.vitorpamplona.amethyst.commons.cordn
import com.vitorpamplona.amethyst.commons.storage.EncryptedAppendLog
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessage
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessageCodec
import com.vitorpamplona.quartz.cordn.sync.EchoState
import com.vitorpamplona.quartz.cordn.sync.GroupCursor
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
import java.nio.ByteBuffer
@@ -168,6 +173,46 @@ class FileCordnGroupStore(
private fun echoStateFile(gid: String) = File(groupDir(gid), "echoes")
/**
* Inside [groupDir] so [deleteGroup]'s recursive delete already covers it:
* a group that left its history behind would keep the plaintext of an
* end-to-end encrypted conversation after the key that read it was gone.
*/
private fun messagesFile(gid: String) = File(groupDir(gid), "messages")
private fun summaryFile(gid: String) = File(groupDir(gid), "newest")
/**
* Appending a segment rather than rewriting the conversation, so the cost
* of receiving a message does not grow with how much has been said.
*/
private val messageLog =
EncryptedAppendLog(
encrypt = cipher::encrypt,
// The log treats null as "this segment is unreadable" and carries on
// with the rest, which is what one corrupt segment should cost.
decrypt = { runCatching { cipher.decrypt(it) }.getOrNull() },
)
private val messageLock = Mutex()
/**
* Envelope ids already in a group's log.
*
* Dedup cannot be the log's own whole-entry comparison: the same message
* re-delivered after a crash carries the same envelope but not necessarily
* the same cursor, so the entries differ as strings while naming one
* message. Built once per group from the log the first time it is touched.
*/
private val seenIds = mutableMapOf<String, MutableSet<HexKey>>()
private fun idsFor(gid: String): MutableSet<HexKey> =
seenIds.getOrPut(gid) {
messageLog
.readAll(messagesFile(gid))
.mapNotNullTo(mutableSetOf()) { CordnDeliveredMessageCodec.decodeOrNull(it)?.envelope?.id }
}
override suspend fun saveGroup(
gid: String,
state: ByteArray,
@@ -183,6 +228,13 @@ class FileCordnGroupStore(
override suspend fun deleteGroup(gid: String) {
withContext(Dispatchers.IO) {
messageLock.withLock {
// The log caches a file's entries by path, so dropping the
// directory alone would leave a re-join of the same gid reading
// the previous membership's messages out of memory.
messageLog.forget(messagesFile(gid))
seenIds.remove(gid)
}
// The cursor goes with it. Leaving one behind would mean a later
// re-join of the same gid resumes from a cursor belonging to a
// group it is no longer in, skipping everything before it.
@@ -253,6 +305,49 @@ class FileCordnGroupStore(
}
}
override suspend fun appendMessage(
gid: String,
message: CordnDeliveredMessage,
) = withContext(Dispatchers.IO) {
messageLock.withLock {
val ids = idsFor(gid)
// Idempotent on the envelope id. The crash window between this and
// saveCursor is deliberate — see the store interface — and it is
// this check that makes re-delivery free rather than duplicating.
if (!ids.add(message.envelope.id)) return@withContext
groupDir(gid).mkdirs()
messageLog.append(messagesFile(gid), CordnDeliveredMessageCodec.encode(message))
// Written on the same beat, so the inbox preview cannot disagree
// with the room. Whole-blob rather than appended: it is one entry
// that is always overwritten.
atomicWrite(summaryFile(gid), cipher.encrypt(CordnMessageSummaryCodec.encode(message, ids.size)))
}
}
override suspend fun loadMessages(gid: String): List<CordnDeliveredMessage> =
withContext(Dispatchers.IO) {
messageLock.withLock {
// One unreadable entry costs that message, not the conversation
// behind it in the file.
messageLog.readAll(messagesFile(gid)).mapNotNull { CordnDeliveredMessageCodec.decodeOrNull(it) }
}
}
override suspend fun loadMessageSummary(gid: String): CordnMessageSummary? =
withContext(Dispatchers.IO) {
val file = summaryFile(gid)
if (!file.exists()) return@withContext null
try {
CordnMessageSummaryCodec.decode(cipher.decrypt(file.readBytes()))
} catch (e: Exception) {
// A summary is a derived convenience; losing one costs a preview
// line until the next message, not the history it summarises.
Log.w("FileCordnGroupStore", "unreadable message summary for $gid: ${e.message}", e)
null
}
}
override suspend fun saveEchoState(
gid: String,
state: EchoState,
@@ -18,7 +18,7 @@
* 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
package com.vitorpamplona.amethyst.commons.storage
import java.io.File
import java.io.FileOutputStream
@@ -34,10 +34,17 @@ import java.io.RandomAccessFile
* plain := uint32 count, (uint32 len, byte[len])*
* ```
*
* Protocol-neutral on purpose: it takes [encrypt]/[decrypt] lambdas and stores
* opaque strings, so both group-chat implementations keep their own ciphers and
* their own formats on top of one file layout. It used to live in the `marmot`
* package, which made it unreachable from cordn — the independence guard
* forbids either feature importing the other by name — for no reason other
* than where it happened to be written first.
*
* Segments are the whole point. The format this replaces was a single blob
* covering the entire history, so recording one line meant pushing every line
* ever written back through the cipher and out to disk again — work that grew
* with the log and, for Marmot's message log, was paid on the send path. A
* with the log and was paid on the send path. A
* conversation a few thousand messages long was moving hundreds of KB through
* a hardware-backed cipher to append a couple of hundred bytes.
*
@@ -25,6 +25,10 @@ import com.vitorpamplona.quartz.cordn.groups.CordnGroupPolicy
import com.vitorpamplona.quartz.cordn.spec00Coordinator.GroupMessage
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ICoordinator
import com.vitorpamplona.quartz.cordn.spec01GroupMetadata.CordnGroupMetadata
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnAnnotationIndex
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessage
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnEnvelope
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnMessageReferences
import com.vitorpamplona.quartz.cordn.sync.GroupCursor
import com.vitorpamplona.quartz.mls.group.MlsGroup
import com.vitorpamplona.quartz.mls.messages.KeyPackageBundle
@@ -298,6 +302,148 @@ class CordnGroupManagerTest {
)
}
@Test
fun `a message is stored before the cursor that would skip it`() =
runTest {
// The ordering IS the correctness property. Write the cursor first
// and a crash in between loses the message permanently: both cordn
// seal keys are epoch-derived, so the copy the coordinator still
// holds can no longer be opened, and a cursor past it means it is
// never offered again. Write the message first and the same crash
// costs nothing — the message is on disk, the cursor still points
// before it, and the re-delivery dedups.
val coordinator = FakeCoordinator(callerPubKey = alice)
val store = OrderRecordingStore()
val m = manager(alice, coordinator, store)
m.createGroup(gid, CordnGroupMetadata(name = "Order"))
m.send(gid, "one")
val append = store.calls.indexOf("appendMessage")
val cursor = store.calls.indexOf("saveCursor")
assertTrue(append >= 0, "the message was never stored at all: ${store.calls}")
assertTrue(cursor >= 0, "the cursor was never stored at all: ${store.calls}")
assertTrue(
append < cursor,
"the cursor was written before the message, so a crash between them loses it: ${store.calls}",
)
}
@Test
fun `the same message delivered twice is stored once`() =
runTest {
// The crash window above lands here: a re-fetch after a restart
// re-delivers a message that is already on disk, and it arrives
// carrying a cursor of its own. Dedup on the whole stored entry
// would not catch that — same envelope, different cursor, different
// bytes — so it has to key on the envelope id.
val store = InMemoryCordnGroupStore()
val envelope =
CordnEnvelope.build(
pubKey = alice,
createdAt = 1_757_000_000L,
kind = 9,
content = "once",
)
store.appendMessage(gid, CordnDeliveredMessage(envelope, cursor = 7))
store.appendMessage(gid, CordnDeliveredMessage(envelope, cursor = 9))
assertEquals(1, store.loadMessages(gid).size, "the same envelope was stored twice")
}
@Test
fun `a stored conversation reloads as the one that was ingested`() =
runTest {
// Equality of the raw list is not enough: what the room shows is
// the annotation fold — edits applied, deletions withdrawn,
// reactions counted. A reload that produced the same messages but a
// different fold would look right in a list assertion and wrong on
// screen.
val coordinator = FakeCoordinator(callerPubKey = alice)
val store = InMemoryCordnGroupStore()
val m = manager(alice, coordinator, store)
m.createGroup(gid, CordnGroupMetadata(name = "Fold"))
val first = m.send(gid, "original")
m.post(
gid = gid,
content = "edited",
editTo =
CordnMessageReferences.Target(
id = first.envelope.id,
pubKey = first.envelope.pubKey,
kind = first.envelope.kind,
tags = first.envelope.tags,
),
)
m.post(gid, "\uD83D\uDC4D", reactionTo = CordnMessageReferences.Target(first.envelope.id, first.envelope.pubKey, first.envelope.kind, first.envelope.tags))
val live = CordnAnnotationIndex.of(listOf(first) + m.storedMessages(gid).drop(1))
val reloaded = CordnAnnotationIndex.of(m.storedMessages(gid))
assertEquals("edited", reloaded.contentOf(first.envelope.id), "the edit did not survive the reload")
assertTrue(reloaded.isEdited(first.envelope.id))
assertEquals(
live.reactions[first.envelope.id]?.keys,
reloaded.reactions[first.envelope.id]?.keys,
"the reaction fold differs after a reload",
)
}
@Test
fun `the summary names the newest message without loading the log`() =
runTest {
val coordinator = FakeCoordinator(callerPubKey = alice)
val m = manager(alice, coordinator)
m.createGroup(gid, CordnGroupMetadata(name = "Summary"))
m.send(gid, "first")
val last = m.send(gid, "last")
val summary = m.storedMessageSummary(gid)
assertEquals(last.envelope.id, summary?.newest?.envelope?.id)
assertEquals(2, summary?.count)
}
@Test
fun `leaving a group takes its messages with it`() =
runTest {
// A group whose history outlived it would keep the plaintext of an
// end-to-end encrypted conversation on disk after the MLS state
// that could read it was thrown away.
val store = InMemoryCordnGroupStore()
val envelope = CordnEnvelope.build(alice, 1_757_000_000L, 9, content = "secret")
store.saveGroup(gid, ByteArray(1))
store.appendMessage(gid, CordnDeliveredMessage(envelope, cursor = 1))
store.deleteGroup(gid)
assertEquals(emptyList<CordnDeliveredMessage>(), store.loadMessages(gid))
assertEquals(null, store.loadMessageSummary(gid))
}
/** Records the order writes arrive in, so an ordering invariant can be asserted. */
private class OrderRecordingStore(
private val inner: CordnGroupStore = InMemoryCordnGroupStore(),
) : CordnGroupStore by inner {
val calls = mutableListOf<String>()
override suspend fun appendMessage(
gid: String,
message: CordnDeliveredMessage,
) {
calls += "appendMessage"
inner.appendMessage(gid, message)
}
override suspend fun saveCursor(
gid: String,
cursor: GroupCursor,
) {
calls += "saveCursor"
inner.saveCursor(gid, cursor)
}
}
@Test
fun `a posted commit is still recognised as ours after a restart`() =
runTest {
@@ -18,7 +18,7 @@
* 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
package com.vitorpamplona.amethyst.commons.storage
import java.io.File
import java.nio.file.Files
@@ -0,0 +1,110 @@
/*
* 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.cordn.spec02Envelopes
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.buildJsonObject
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
import kotlinx.serialization.json.long
import kotlinx.serialization.json.put
/**
* The on-disk form of one delivered message.
*
* ## Why the plaintext is what gets stored
*
* A cordn payload is sealed twice — the MLS application message, then
* [com.vitorpamplona.quartz.cordn.spec03Payloads.SealedPayload] under an
* exporter key — and both keys are derived per epoch. Once the group ratchets
* forward, the ciphertext the coordinator still holds is unreadable to us, and
* the cursor has in any case advanced past it. Ingestion is the only moment the
* message is in the clear, so this is not a cache in front of a source of
* truth: it is the only copy.
*
* ## Why the cursor rides along
*
* A room orders on the cursor and the unread divider compares against it, and
* it is not derivable from the envelope — the envelope's `created_at` is the
* sender's clock, which is a claim, while the cursor is the coordinator's
* sequence. Storing the envelope alone would lose the ordering the room is
* built on.
*/
object CordnDeliveredMessageCodec {
const val VERSION = 1
private const val V = "v"
private const val CURSOR = "c"
private const val ENVELOPE = "e"
fun encode(message: CordnDeliveredMessage): String =
Json.encodeToString(
JsonObject.serializer(),
buildJsonObject {
put(V, VERSION)
put(CURSOR, message.cursor)
put(ENVELOPE, message.envelope.toJsonObject())
},
)
/**
* Reads one entry back.
*
* Throws on anything malformed rather than returning null, so a caller
* reading a log decides for itself whether one bad entry drops the line or
* the conversation. [decodeOrNull] is the "drop the line" form.
*/
fun decode(entry: String): CordnDeliveredMessage {
val json =
Json.parseToJsonElement(entry) as? JsonObject
?: throw IllegalArgumentException("a stored cordn message must be a JSON object")
val version =
json[V]?.jsonPrimitive?.content?.toIntOrNull()
?: throw IllegalArgumentException("a stored cordn message is missing `$V`")
require(version == VERSION) { "unknown stored cordn message version: $version" }
val cursor =
json[CURSOR]?.jsonPrimitive?.long
?: throw IllegalArgumentException("a stored cordn message is missing `$CURSOR`")
val envelope =
json[ENVELOPE]?.jsonObject
?: throw IllegalArgumentException("a stored cordn message is missing `$ENVELOPE`")
return CordnDeliveredMessage(CordnEnvelope.fromJsonObject(envelope), cursor)
}
/**
* [decode], or null when the entry cannot be read.
*
* Loading a room uses this: one corrupted entry — a half-written segment, a
* format from a version that did not ship — should cost that message, not
* every message behind it in the log.
*/
fun decodeOrNull(entry: String): CordnDeliveredMessage? =
try {
decode(entry)
} catch (e: Exception) {
null
}
}
@@ -137,6 +137,30 @@ data class CordnEnvelope(
require(SIG !in json) { "cordn envelope must not carry a `sig` (spec/02.md §2)" }
val envelope = fromJsonObject(json)
// Both checks are MUSTs, and each covers a different lie: the id
// check catches a rewritten body, the pubkey check catches a member
// posting under someone else's name. fromJsonObject does the first,
// because a body that does not hash to its id is malformed wherever
// it came from; only the second needs the MLS sender.
require(envelope.pubKey == senderIdentity) {
"cordn envelope claims pubkey ${envelope.pubKey} but the MLS sender is $senderIdentity"
}
return envelope
}
/**
* Reads an envelope out of its JSON form, checking that the body hashes
* to the id it carries.
*
* Separate from [decode] because a re-read from our own storage has no
* MLS sender to compare against — the sender was checked when the
* message was first ingested, and the plaintext has been ours since.
* What is still worth checking on the way back in is integrity, which
* the id covers.
*/
fun fromJsonObject(json: JsonObject): CordnEnvelope {
val envelope =
CordnEnvelope(
id = json.str(ID),
@@ -151,15 +175,9 @@ data class CordnEnvelope(
content = json.str(CONTENT),
)
// Both checks are MUSTs, and each covers a different lie: the id
// check catches a rewritten body, the pubkey check catches a member
// posting under someone else's name.
require(envelope.id == envelope.computedId()) {
"cordn envelope id ${envelope.id} does not match its contents (${envelope.computedId()})"
}
require(envelope.pubKey == senderIdentity) {
"cordn envelope claims pubkey ${envelope.pubKey} but the MLS sender is $senderIdentity"
}
return envelope
}