diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/cordn/CordnRuntime.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/cordn/CordnRuntime.kt index d36806dc12..e2ee46afda 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/cordn/CordnRuntime.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/cordn/CordnRuntime.kt @@ -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. */ diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt index 436def3b41..5d0928cd43 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/marmot/AndroidMarmotMessageStore.kt @@ -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 diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt index 32fed3bc1b..abc8b52fa2 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt @@ -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>() + + 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 = 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) } } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupStore.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupStore.kt index 98fcbf83e1..3102e3c744 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupStore.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupStore.kt @@ -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 + + /** + * 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() private val roomStates = mutableMapOf() private val echoStates = mutableMapOf() + private val messages = mutableMapOf>() 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 = 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 = 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, diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cordnGroups/CordnGroupChatroom.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cordnGroups/CordnGroupChatroom.kt index 906643196e..5ea4437112 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cordnGroups/CordnGroupChatroom.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/cordnGroups/CordnGroupChatroom.kt @@ -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, diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cordn/FileCordnStores.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cordn/FileCordnStores.kt index 4f29ec4a43..914fb8c555 100644 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cordn/FileCordnStores.kt +++ b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/cordn/FileCordnStores.kt @@ -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>() + + private fun idsFor(gid: String): MutableSet = + 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 = + 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, diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/marmot/EncryptedAppendLog.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/storage/EncryptedAppendLog.kt similarity index 96% rename from commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/marmot/EncryptedAppendLog.kt rename to commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/storage/EncryptedAppendLog.kt index b8ef668cb8..d81d5f37f3 100644 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/marmot/EncryptedAppendLog.kt +++ b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/storage/EncryptedAppendLog.kt @@ -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. * diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt index 58fc3bfbd9..6dbc3088d2 100644 --- a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt @@ -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(), 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() + + 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 { diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/EncryptedAppendLogTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/storage/EncryptedAppendLogTest.kt similarity index 99% rename from commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/EncryptedAppendLogTest.kt rename to commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/storage/EncryptedAppendLogTest.kt index acf4be7de4..87629543b7 100644 --- a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/marmot/EncryptedAppendLogTest.kt +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/storage/EncryptedAppendLogTest.kt @@ -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 diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnDeliveredMessageCodec.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnDeliveredMessageCodec.kt new file mode 100644 index 0000000000..11879dd57b --- /dev/null +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnDeliveredMessageCodec.kt @@ -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 + } +} diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnEnvelope.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnEnvelope.kt index 41b4ec2293..03eccdeadd 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnEnvelope.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec02Envelopes/CordnEnvelope.kt @@ -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 }