mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix(cordn): a client must recognise its own traffic after a restart
Tier B's second finding (interop plan §7.1). A sender's own stream entries
came back as `undecryptable`: its Commit, sealed under the pre-commit epoch
key it has since left, and its message, from a ratchet generation already
consumed.
The record of what is ours lived in memory while the cursor beside it was
persisted. That asymmetry is not obvious until something crosses a process
boundary, because posting deliberately does **not** advance the cursor — a
lower cursor may still hold somebody else's unprocessed message — so the next
run is guaranteed to re-read what this one posted, and needs to be told it is
its own.
`EchoState` is that record, saved beside `GroupCursor` on the same beat and
deleted when there is nothing pending. Own-message cursors at or below the
fetch cursor are pruned, because the stream has delivered them and they can
never come round again. Pending Commits are kept until their echo matches,
however long that takes: `Ingestion.SelfEchoUnapplied` is how a client that
posted a Commit and died before adopting it applies its own Commit, and that
echo is the only copy it will ever be offered.
`amy` exposed this on every run because each verb is a process, but it is not
an `amy` bug. The Android app hits it whenever the OS kills it between sending
and syncing, which is the ordinary case rather than a corner. Two levels of
damage, and the quieter one is worse: a visible gap in the sender's own
conversation, and a recovery path that depended on exactly the record that
dying destroyed.
**The tests that let it through.** `a group survives a restart` already
claimed to cover this, with `assertTrue(delivered.none { it is
Delivery.Message })`. An `Undecryptable` is not a `Message`, so a gap in your
own conversation satisfied the assertion. The lesson is narrow and reusable:
an assertion about what a delivery is *not* passes for every outcome nobody
thought of. Both restart tests now assert what it **is** — `Echo` — and both
fail if either half of the persistence is removed (verified by removing each
half in turn).
Feeding an unrecognised Commit back through the engine, rather than reporting
it, would have been worse than the gap: the epoch advances twice and the
sender desynchronises from every other member, with the failure surfacing
several epochs later as "my messages stopped decrypting" and nothing pointing
at the cause. So the pre-existing `Undecryptable` outcome was the right
failure mode to have — it just should never have been reached.
`cli/tests/cordn/tier-b.sh` now passes end to end against the reference
coordinator: handshake, `kp_publish`, group create, invite, Welcome opened
without joining, join, messages both ways with the sender's own traffic
reported as echoes, and both sides agreeing on epoch 1 and the same two
members.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
+6
-1
@@ -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() {
|
||||
|
||||
+93
@@ -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<PendingEpochOperation>()
|
||||
while (commits.hasRemaining) {
|
||||
val sealed = commits.readOpaque2().decodeToString()
|
||||
pending += PendingEpochOperation(sealed, commits.readUint8() == 1)
|
||||
}
|
||||
|
||||
val cursors = TlsReader(reader.readOpaque2())
|
||||
val own = mutableListOf<Long>()
|
||||
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<String, ByteArray>()
|
||||
private val cursors = mutableMapOf<String, GroupCursor>()
|
||||
private val viaRequest = mutableSetOf<String>()
|
||||
private val roomStates = mutableMapOf<String, CordnRoomState>()
|
||||
private val echoStates = mutableMapOf<String, EchoState>()
|
||||
|
||||
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<String> = 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,
|
||||
|
||||
+27
@@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+50
-5
@@ -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<CordnGroupManager.Delivery>()
|
||||
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<CordnGroupManager.Delivery>()
|
||||
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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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<String, GroupCursor> = 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<String, EchoState> = inboxes.mapValues { it.value.echoes() }
|
||||
|
||||
/**
|
||||
* Drains history for [gids] until every group is current.
|
||||
*
|
||||
|
||||
@@ -127,9 +127,10 @@ data class PendingEpochOperation(
|
||||
*/
|
||||
class GroupInbox(
|
||||
var cursor: GroupCursor = GroupCursor(),
|
||||
echoes: EchoState = EchoState(),
|
||||
) {
|
||||
private val pendingOperations = mutableMapOf<String, PendingEpochOperation>()
|
||||
private val ownMessageCursors = mutableSetOf<Long>()
|
||||
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<PendingEpochOperation> = emptyList(),
|
||||
val ownMessageCursors: List<Long> = emptyList(),
|
||||
) {
|
||||
val isEmpty: Boolean get() = pendingCommits.isEmpty() && ownMessageCursors.isEmpty()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user