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 b16b4e152c..a4512da5a0 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 @@ -156,7 +156,7 @@ class CordnGroupManager( store.listGroups().forEach { gid -> val blob = store.loadGroup(gid) ?: return@forEach groups[gid] = MlsGroup.restore(MlsGroupState.decodeTls(blob), CordnGroupPolicy) - store.loadCursor(gid)?.let { sync.restore(gid, it) } + store.loadCursor(gid)?.let { sync.restore(gid, it, store.loadEchoState(gid)) } } _gids.value = groups.keys.toSet() } @@ -707,6 +707,11 @@ class CordnGroupManager( val group = groups[gid] ?: return store.saveGroup(gid, group.saveState().encodeTls()) 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 + // what is ours does not leaves the next run re-reading its own + // traffic with no way to recognise it. See EchoState. + sync.echoes()[gid]?.let { store.saveEchoState(gid, it) } } private suspend fun persistAll() { 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 a94f86fa0d..98fcbf83e1 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,7 +20,9 @@ */ package com.vitorpamplona.amethyst.commons.cordn +import com.vitorpamplona.quartz.cordn.sync.EchoState import com.vitorpamplona.quartz.cordn.sync.GroupCursor +import com.vitorpamplona.quartz.cordn.sync.PendingEpochOperation import com.vitorpamplona.quartz.mls.codec.TlsReader import com.vitorpamplona.quartz.mls.codec.TlsWriter @@ -62,6 +64,25 @@ interface CordnGroupStore { suspend fun loadCursor(gid: String): GroupCursor? + /** + * Saves the record of which stream entries are this client's own. + * + * Beside the cursor because it is only meaningful with it, and durable for + * the same reason: posting does not advance the cursor (a lower one may + * still hold somebody else's unprocessed message), so the next run + * re-reads what this one posted. Without this it cannot tell that it is + * its own — a Commit sealed under an epoch key it has since left, a + * message from a ratchet generation already consumed — and reports a gap + * in its own conversation. See `EchoState`. + */ + suspend fun saveEchoState( + gid: String, + state: EchoState, + ) + + /** The saved echo bookkeeping for [gid], or an empty one. */ + suspend fun loadEchoState(gid: String): EchoState + /** * Records that admission to [gid] went through `join_request_store`. * @@ -139,12 +160,74 @@ object CordnRoomStateCodec { } } +/** + * The on-disk layout of an [EchoState]. + * + * ``` + * version:u16 + * pending_commits: vector2 of { sealed:opaque2, applied:u8 } + * own_cursors: vector2 of u64 + * ``` + * + * Empty on anything unreadable, which is the safe direction: losing the record + * costs a client one run of reporting its own traffic as a gap, while + * accepting a half-decoded one could let it skip somebody else's message as + * though it were its own. + */ +object EchoStateCodec { + const val VERSION = 1 + + fun encode(state: EchoState): ByteArray { + val writer = TlsWriter() + writer.putUint16(VERSION) + + val commits = TlsWriter() + state.pendingCommits.forEach { + commits.putOpaque2(it.sealedBase64.encodeToByteArray()) + commits.putUint8(if (it.localStateApplied) 1 else 0) + } + writer.putOpaque2(commits.toByteArray()) + + val cursors = TlsWriter() + state.ownMessageCursors.forEach { cursors.putUint64(it) } + writer.putOpaque2(cursors.toByteArray()) + + return writer.toByteArray() + } + + fun decode(bytes: ByteArray): EchoState = + try { + val reader = TlsReader(bytes) + if (reader.readUint16() != VERSION) { + EchoState() + } else { + val commits = TlsReader(reader.readOpaque2()) + val pending = mutableListOf() + while (commits.hasRemaining) { + val sealed = commits.readOpaque2().decodeToString() + pending += PendingEpochOperation(sealed, commits.readUint8() == 1) + } + + val cursors = TlsReader(reader.readOpaque2()) + val own = mutableListOf() + while (cursors.hasRemaining) { + own += cursors.readUint64() + } + + EchoState(pending, own) + } + } catch (e: Exception) { + EchoState() + } +} + /** A [CordnGroupStore] that keeps everything in memory. Tests, and nothing else. */ class InMemoryCordnGroupStore : CordnGroupStore { private val groups = mutableMapOf() private val cursors = mutableMapOf() private val viaRequest = mutableSetOf() private val roomStates = mutableMapOf() + private val echoStates = mutableMapOf() override suspend fun saveGroup( gid: String, @@ -160,6 +243,7 @@ class InMemoryCordnGroupStore : CordnGroupStore { cursors.remove(gid) viaRequest.remove(gid) roomStates.remove(gid) + echoStates.remove(gid) } override suspend fun listGroups(): List = groups.keys.toList() @@ -179,6 +263,15 @@ class InMemoryCordnGroupStore : CordnGroupStore { override suspend fun loadRoomState(gid: String): CordnRoomState = roomStates[gid] ?: CordnRoomState() + override suspend fun saveEchoState( + gid: String, + state: EchoState, + ) { + echoStates[gid] = state + } + + override suspend fun loadEchoState(gid: String): EchoState = echoStates[gid] ?: EchoState() + override suspend fun saveCursor( gid: String, cursor: GroupCursor, 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 aa106920e8..681849eb3d 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,6 +20,7 @@ */ package com.vitorpamplona.amethyst.commons.cordn +import com.vitorpamplona.quartz.cordn.sync.EchoState import com.vitorpamplona.quartz.cordn.sync.GroupCursor import com.vitorpamplona.quartz.nip01Core.core.HexKey import kotlinx.coroutines.Dispatchers @@ -162,6 +163,8 @@ class FileCordnGroupStore( private fun roomStateFile(gid: String) = File(groupDir(gid), "room") + private fun echoStateFile(gid: String) = File(groupDir(gid), "echoes") + override suspend fun saveGroup( gid: String, state: ByteArray, @@ -246,6 +249,30 @@ class FileCordnGroupStore( CordnRoomState() } } + + override suspend fun saveEchoState( + gid: String, + state: EchoState, + ) = withContext(Dispatchers.IO) { + // Deleted when there is nothing pending, so the common steady state is + // no file rather than an empty one. + if (state.isEmpty) { + echoStateFile(gid).delete() + return@withContext + } + atomicWrite(echoStateFile(gid), cipher.encrypt(EchoStateCodec.encode(state))) + } + + override suspend fun loadEchoState(gid: String): EchoState = + withContext(Dispatchers.IO) { + val file = echoStateFile(gid) + if (!file.exists()) return@withContext EchoState() + try { + EchoStateCodec.decode(cipher.decrypt(file.readBytes())) + } catch (e: Exception) { + EchoState() + } + } } /** 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 19f36dedbe..fe15fdf48c 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 @@ -271,16 +271,61 @@ class CordnGroupManagerTest { assertEquals(setOf(gid), restarted.gids.value) assertEquals("Persisted", CordnGroupMetadata.fromExtensions(restarted.group(gid)!!.extensions)?.name) - // The cursor came back too, so our own message is not re-delivered - // as somebody else's on the next catch-up. + // Our own message comes back, because posting deliberately does + // not advance the cursor — a lower cursor may still hold somebody + // else's unprocessed message. So it must come back recognised. + // + // This assertion used to be `none { it is Message }`, and it + // passed while the delivery was `Undecryptable`: not a Message, so + // not a failure, so nobody noticed that a sender reported a gap in + // its own conversation. Only a run against a real coordinator + // made it visible. Assert what it IS, not what it is not. val delivered = mutableListOf() restarted.catchUp { delivered += it } - assertTrue( - delivered.none { it is CordnGroupManager.Delivery.Message }, - "a restored cursor must not replay our own history: ${delivered.map { it::class.simpleName }}", + assertEquals( + listOf("Echo"), + delivered.map { it::class.simpleName }, + "our own message must come back as an echo, not as a gap", ) } + @Test + fun `a posted commit is still recognised as ours after a restart`() = + runTest { + // The same bug, one layer worse. A Commit is sealed under the + // PRE-commit epoch key (spec/03.md §5), which the poster no longer + // has once it has advanced — so when the echo arrives it cannot be + // opened at all, and without the record that it is ours it is + // indistinguishable from a payload from an epoch we never had. + // + // Feeding it back through the engine instead would be worse than a + // gap: it would advance the epoch twice and desynchronise us from + // every other member, with the failure surfacing epochs later. + val coordinator = FakeCoordinator(callerPubKey = alice) + val store = InMemoryCordnGroupStore() + val (bundle, stored) = bobsPublication() + coordinator.seedKeyPackage(stored) + + val first = manager(alice, coordinator, store) + first.createGroup(gid, CordnGroupMetadata(name = "Restart")) + first.invite(gid, bob, stored.keyPackageRef) + assertEquals(1, first.group(gid)!!.epoch) + + val restarted = manager(alice, coordinator, store) + restarted.restore() + + val delivered = mutableListOf() + restarted.catchUp { delivered += it } + + assertEquals( + listOf("Echo"), + delivered.map { it::class.simpleName }, + "a restarted client must recognise the Commit it posted: ${delivered.map { it::class.simpleName }}", + ) + assertEquals(1, restarted.group(gid)!!.epoch, "and must not apply it a second time") + assertTrue(bundle.keyPackage.toTlsBytes().isNotEmpty()) + } + @Test fun `inviting someone the coordinator has no KeyPackage for fails loudly`() = runTest { diff --git a/quartz/plans/2026-09-17-cordn-interop.md b/quartz/plans/2026-09-17-cordn-interop.md index d1f1e577a4..1cb76d3488 100644 --- a/quartz/plans/2026-09-17-cordn-interop.md +++ b/quartz/plans/2026-09-17-cordn-interop.md @@ -658,7 +658,21 @@ visible damage is a gap in the sender's own conversation. The latent damage is w `Ingestion.SelfEchoUnapplied` is how a client that died mid-post applies its own Commit, and that recovery path depends on exactly the record that dying destroys. -Fixed by persisting the bookkeeping next to the cursor — see the commit that follows this one. +Fixed by persisting the bookkeeping next to the cursor: `EchoState` beside `GroupCursor`, saved on +the same beat, deleted when there is nothing pending. Own-message cursors at or below the fetch +cursor are pruned — the stream has delivered them and they can never come round again — while +pending Commits are kept until their echo matches, because that echo is the only copy that will +ever be offered. + +There was a test for this too, and it also passed: `a group survives a restart` asserted +`delivered.none { it is Delivery.Message }`. An `Undecryptable` is not a `Message`, so a gap in +the sender's own conversation satisfied it. Both restart tests now assert what the delivery **is** +— `Echo` — and both fail if either half of the persistence is removed. + +With both fixed, `cli/tests/cordn/tier-b.sh` passes end to end: handshake, `kp_publish`, group +create, invite, Welcome opened without joining, join, messages both ways with the sender's own +traffic reported as echoes rather than gaps, and both sides agreeing on epoch 1 and the same two +members. ### What Tier B is, and is not diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/CordnGroupSync.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/CordnGroupSync.kt index bae7d8f8f1..8c3af2af25 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/CordnGroupSync.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/CordnGroupSync.kt @@ -54,13 +54,24 @@ class CordnGroupSync( fun restore( gid: String, cursor: GroupCursor, + echoes: EchoState = EchoState(), ) { - inbox(gid).cursor = cursor + inboxes[gid] = GroupInbox(cursor, echoes) } /** The cursors to persist. */ fun cursors(): Map = inboxes.mapValues { it.value.cursor } + /** + * The echo bookkeeping to persist, per group. + * + * Saved on the same beat as [cursors] and for the same reason: a cursor + * that outlives the process while the record of what is ours does not + * leaves a client re-reading its own traffic with no way to recognise it. + * See [EchoState]. + */ + fun echoes(): Map = inboxes.mapValues { it.value.echoes() } + /** * Drains history for [gids] until every group is current. * diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/GroupCursor.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/GroupCursor.kt index cd236ff72f..19f6fc71ce 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/GroupCursor.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/sync/GroupCursor.kt @@ -127,9 +127,10 @@ data class PendingEpochOperation( */ class GroupInbox( var cursor: GroupCursor = GroupCursor(), + echoes: EchoState = EchoState(), ) { - private val pendingOperations = mutableMapOf() - private val ownMessageCursors = mutableSetOf() + private val pendingOperations = echoes.pendingCommits.associateByTo(mutableMapOf()) { it.sealedBase64 } + private val ownMessageCursors = echoes.ownMessageCursors.toMutableSet() /** Records a Commit we just posted, so its echo is recognised. */ fun expectEcho(operation: PendingEpochOperation) { @@ -184,4 +185,49 @@ class GroupInbox( fun skipUnprocessable(cursor: Long) { this.cursor = this.cursor.advancedTo(cursor) } + + /** + * The bookkeeping to persist beside [cursor]. + * + * Own-message cursors at or below the fetch cursor are dropped: the stream + * has already delivered them, so they can never come round again, and + * keeping them would grow the record for the life of the group. + * + * Pending Commits are kept until their echo matches, however long that + * takes, because the echo is the only copy that will ever be offered and + * [Ingestion.SelfEchoUnapplied] is how a client that died mid-post applies + * its own Commit. + */ + fun echoes(): EchoState = + EchoState( + pendingCommits = pendingOperations.values.toList(), + ownMessageCursors = ownMessageCursors.filter { it > cursor.fetchCursor }.sorted(), + ) +} + +/** + * The part of a [GroupInbox] that has to outlive the process. + * + * Found by running against a real coordinator, not by reading. Posting + * deliberately does not advance the cursor — a lower cursor may still hold + * somebody else's unprocessed message — so a client that exits between + * posting and ingesting comes back with a cursor that will re-read its own + * traffic. Without this record it cannot tell that it is its own: its Commit + * was sealed under an epoch key it has since left, and its message came from a + * ratchet generation already consumed, so both arrive as + * [Ingestion.Process] and fail to open. + * + * The visible damage is a gap in the sender's own conversation. The worse, + * narrower damage is that [Ingestion.SelfEchoUnapplied] — the recovery path + * for a client that posted a Commit and died before adopting it — depends on + * exactly the record that dying used to destroy. + * + * Every client needs this, not only a process-per-command one: a phone killed + * between sending a message and syncing is the ordinary case, not a corner. + */ +data class EchoState( + val pendingCommits: List = emptyList(), + val ownMessageCursors: List = emptyList(), +) { + val isEmpty: Boolean get() = pendingCommits.isEmpty() && ownMessageCursors.isEmpty() }