mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat(marmot): add several members in one commit, and stop passing vectors vacuously
Batched Adds ------------ MDK — and therefore both White Noise clients — turns a `create_group` with N invitees into ONE Commit: N Add proposals and one Welcome carrying N `EncryptedGroupSecrets`, keyed by KeyPackage reference (RFC 9420 §12.4.3.1). We staged one Add per commit, so the same group landed at epoch N instead of epoch 1 and cost a round trip per invitee. `MlsGroup.addMembers` / `MlsGroupManager.stageAddMembers` propose-N then commit once; `MarmotManager.addMembers` publishes that single commit and fans the same Welcome bytes out to each invitee. The singular entry points delegate, so no caller changes. Scenario vectors: two shapes, one of them unread ------------------------------------------------ The conformance vectors state their expectations two ways — `expected_trace.observations` and `expected_outcomes` — and the parser only read the first. Seven of the nine vectors here use the second, so they parsed to ZERO expectations, replayed their steps and reported green without checking anything. `in_group` was also treated as a leaf, which silently skipped every step nested inside it, and `send_app_message` built a message the runner never queued for delivery. Now parsed and checked: `client_state`, `clients_converged`, `pending_resolution`, `no_pending_work`, inline `assert`/`payload_count` (the forward-secrecy and isolation assertions), `added_members`, `clear_events`, and multi-group clients. `convergence_decision` has no counterpart in our engine, so `convergence-committer-selected` is now REFUSED via `UnsupportedScenarioOutcome` rather than passed on the parts that happen to be modelled. A new guard test fails any vector that parses to nothing to check. three-client-message-exchange and conversation now replay for real. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
+102
-34
@@ -655,6 +655,17 @@ class MarmotManager(
|
||||
return TextMessageBundle(outbound = outbound, innerEvent = innerEvent)
|
||||
}
|
||||
|
||||
/**
|
||||
* A single invitee for [addMemberInvites]: whose KeyPackage is consumed, the bare
|
||||
* KeyPackage bytes as published, and the id of the event that carried them
|
||||
* (the Welcome must reference it so the invitee can retire that KeyPackage).
|
||||
*/
|
||||
data class MemberInvite(
|
||||
val memberPubKey: HexKey,
|
||||
val keyPackageBytes: ByteArray,
|
||||
val keyPackageEventId: HexKey,
|
||||
)
|
||||
|
||||
/**
|
||||
* Add a member to a group by consuming their published [KeyPackageEvent].
|
||||
*
|
||||
@@ -662,21 +673,14 @@ class MarmotManager(
|
||||
* event id into the WelcomeDelivery. Prefer this overload — both the UI's
|
||||
* `Account.addMarmotGroupMember` and the CLI's `group add` command call it.
|
||||
*/
|
||||
@OptIn(kotlin.io.encoding.ExperimentalEncodingApi::class)
|
||||
suspend fun addMember(
|
||||
nostrGroupId: HexKey,
|
||||
keyPackageEvent: KeyPackageEvent,
|
||||
relays: List<NormalizedRelayUrl>,
|
||||
): Pair<OutboundGroupEvent, WelcomeDelivery?> =
|
||||
addMember(
|
||||
nostrGroupId = nostrGroupId,
|
||||
memberPubKey = keyPackageEvent.pubKey,
|
||||
keyPackageBytes =
|
||||
kotlin.io.encoding.Base64
|
||||
.decode(keyPackageEvent.keyPackageBase64()),
|
||||
keyPackageEventId = keyPackageEvent.id,
|
||||
relays = relays,
|
||||
)
|
||||
): Pair<OutboundGroupEvent, WelcomeDelivery?> {
|
||||
val (event, welcomes) = addMembers(nostrGroupId, listOf(keyPackageEvent), relays)
|
||||
return Pair(event, welcomes.firstOrNull())
|
||||
}
|
||||
|
||||
/**
|
||||
* Add a member to a group.
|
||||
@@ -689,18 +693,80 @@ class MarmotManager(
|
||||
keyPackageEventId: HexKey,
|
||||
relays: List<NormalizedRelayUrl>,
|
||||
): Pair<OutboundGroupEvent, WelcomeDelivery?> {
|
||||
// Verify that the KeyPackage credential matches the expected member
|
||||
val (event, welcomes) =
|
||||
addMemberInvites(
|
||||
nostrGroupId,
|
||||
listOf(MemberInvite(memberPubKey, keyPackageBytes, keyPackageEventId)),
|
||||
relays,
|
||||
)
|
||||
return Pair(event, welcomes.firstOrNull())
|
||||
}
|
||||
|
||||
/**
|
||||
* Add several members to a group in a SINGLE commit, from their published
|
||||
* [KeyPackageEvent]s. See [addMemberInvites] for the batching contract.
|
||||
*/
|
||||
@OptIn(kotlin.io.encoding.ExperimentalEncodingApi::class)
|
||||
suspend fun addMembers(
|
||||
nostrGroupId: HexKey,
|
||||
keyPackageEvents: List<KeyPackageEvent>,
|
||||
relays: List<NormalizedRelayUrl>,
|
||||
): Pair<OutboundGroupEvent, List<WelcomeDelivery>> =
|
||||
addMemberInvites(
|
||||
nostrGroupId = nostrGroupId,
|
||||
invites =
|
||||
keyPackageEvents.map {
|
||||
MemberInvite(
|
||||
memberPubKey = it.pubKey,
|
||||
keyPackageBytes =
|
||||
kotlin.io.encoding.Base64
|
||||
.decode(it.keyPackageBase64()),
|
||||
keyPackageEventId = it.id,
|
||||
)
|
||||
},
|
||||
relays = relays,
|
||||
)
|
||||
|
||||
/**
|
||||
* Add several members to a group in a SINGLE commit.
|
||||
*
|
||||
* RFC 9420 §12.4.3.1 lets one Commit carry N Add proposals and one Welcome
|
||||
* holding N `EncryptedGroupSecrets`, one per added member keyed by their
|
||||
* KeyPackage reference. So the group advances by exactly ONE epoch no
|
||||
* matter how many people join, and every invitee receives the SAME Welcome
|
||||
* bytes — each finds its own secrets entry. MDK (and therefore White
|
||||
* Noise) batches this way, so a per-invitee commit loop would diverge from
|
||||
* the reference on epoch numbers and round trips alike.
|
||||
*
|
||||
* Returns the single commit GroupEvent to publish, and one WelcomeDelivery
|
||||
* per invitee — empty if the commit was not confirmed by any relay.
|
||||
*/
|
||||
suspend fun addMemberInvites(
|
||||
nostrGroupId: HexKey,
|
||||
invites: List<MemberInvite>,
|
||||
relays: List<NormalizedRelayUrl>,
|
||||
): Pair<OutboundGroupEvent, List<WelcomeDelivery>> {
|
||||
require(invites.isNotEmpty()) { "addMemberInvites: invites must not be empty" }
|
||||
require(invites.map { it.memberPubKey }.toSet().size == invites.size) {
|
||||
"addMemberInvites: the same member appears twice in one commit"
|
||||
}
|
||||
|
||||
// Verify that each KeyPackage credential matches the expected member
|
||||
// pubkey. Accepts either framing — a peer's published KeyPackage is an
|
||||
// MLSMessage, and bare bytes still arrive from our own pre-fix
|
||||
// publications sitting on relays.
|
||||
val kp = KeyPackageUtils.decodeKeyPackage(keyPackageBytes)
|
||||
val credential = kp.leafNode.credential
|
||||
require(credential is Credential.Basic) {
|
||||
"KeyPackage must use BasicCredential"
|
||||
}
|
||||
require(credential.identity.toHexKey() == memberPubKey) {
|
||||
"KeyPackage credential identity does not match memberPubKey"
|
||||
}
|
||||
val decoded =
|
||||
invites.map { invite ->
|
||||
val kp = KeyPackageUtils.decodeKeyPackage(invite.keyPackageBytes)
|
||||
val credential = kp.leafNode.credential
|
||||
require(credential is Credential.Basic) {
|
||||
"KeyPackage must use BasicCredential"
|
||||
}
|
||||
require(credential.identity.toHexKey() == invite.memberPubKey) {
|
||||
"KeyPackage credential identity does not match memberPubKey"
|
||||
}
|
||||
kp
|
||||
}
|
||||
|
||||
// Per RFC 9420 §12.4 (and MDK), the outbound kind:445 MUST be
|
||||
// outer-encrypted with the pre-commit (epoch-N) exporter secret so
|
||||
@@ -709,29 +775,31 @@ class MarmotManager(
|
||||
// that key.
|
||||
val publication =
|
||||
commitAndPublish(nostrGroupId, relays) {
|
||||
// The BARE KeyPackage, not the bytes as published. Transport
|
||||
// The BARE KeyPackages, not the bytes as published. Transport
|
||||
// framing is the Marmot layer's business; MLS takes the struct.
|
||||
groupManager.stageAddMember(nostrGroupId, kp.toTlsBytes())
|
||||
groupManager.stageAddMembers(nostrGroupId, decoded.map { it.toTlsBytes() })
|
||||
}
|
||||
|
||||
// The Welcome is a SEPARATE, retryable per-invitee delivery obligation
|
||||
// that only exists once the Add is canonical. A Welcome for an epoch
|
||||
// The Welcomes are SEPARATE, retryable per-invitee delivery obligations
|
||||
// that only exist once the Add is canonical. A Welcome for an epoch
|
||||
// no relay accepted would invite someone into a group that does not
|
||||
// exist anywhere else.
|
||||
val welcomeDelivery =
|
||||
val welcomeDeliveries =
|
||||
if (publication.confirmed) {
|
||||
welcomeSender.wrapWelcome(
|
||||
commitResult = publication.commitResult,
|
||||
recipientPubKey = memberPubKey,
|
||||
keyPackageEventId = keyPackageEventId,
|
||||
relays = relays,
|
||||
nostrGroupId = nostrGroupId,
|
||||
)
|
||||
invites.mapNotNull { invite ->
|
||||
welcomeSender.wrapWelcome(
|
||||
commitResult = publication.commitResult,
|
||||
recipientPubKey = invite.memberPubKey,
|
||||
keyPackageEventId = invite.keyPackageEventId,
|
||||
relays = relays,
|
||||
nostrGroupId = nostrGroupId,
|
||||
)
|
||||
}
|
||||
} else {
|
||||
null
|
||||
emptyList()
|
||||
}
|
||||
|
||||
return Pair(publication.event, welcomeDelivery)
|
||||
return Pair(publication.event, welcomeDeliveries)
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+201
-47
@@ -56,11 +56,14 @@ import com.vitorpamplona.quartz.utils.RandomInstance
|
||||
*
|
||||
* ## Refusing rather than skipping
|
||||
*
|
||||
* A step type the runner does not implement throws [UnsupportedScenarioStep].
|
||||
* A step type the runner does not implement throws [UnsupportedScenarioStep],
|
||||
* and an expected outcome it cannot check throws [UnsupportedScenarioOutcome].
|
||||
* Nineteen of the portable vectors need fault injection or group-data steps we
|
||||
* have not modelled, and a runner that quietly ignored those steps would report
|
||||
* a pass for a script it did not execute — worse than no coverage, because it
|
||||
* would look like coverage.
|
||||
* would look like coverage. The same is true one level up: a vector whose
|
||||
* conclusion is a `convergence_decision` we do not model is refused, not passed
|
||||
* on the parts that happen to be checkable.
|
||||
*/
|
||||
class MarmotScenarioRunner(
|
||||
private val vector: ScenarioVector,
|
||||
@@ -70,6 +73,19 @@ class MarmotScenarioRunner(
|
||||
/** Queued outbound events: (sender, event). Drained by `deliver_all`. */
|
||||
private val inFlight = mutableListOf<Pair<String, Event>>()
|
||||
|
||||
/**
|
||||
* The group label the current step runs against.
|
||||
*
|
||||
* `in_group` wraps a step and names a label; everything else runs against
|
||||
* [DEFAULT_GROUP]. A client can be in several groups at once and the
|
||||
* isolation vectors are entirely about what does NOT cross between them,
|
||||
* so a single group per client would test the opposite of the point.
|
||||
*/
|
||||
private var currentGroup = DEFAULT_GROUP
|
||||
|
||||
/** Which label each created group id belongs to, so joiners can be filed. */
|
||||
private val labelByGroupId = mutableMapOf<HexKey, String>()
|
||||
|
||||
private class VectorClient(
|
||||
val name: String,
|
||||
) {
|
||||
@@ -81,7 +97,12 @@ class MarmotScenarioRunner(
|
||||
/** Delivered but not yet processed — `tick` is what processes. */
|
||||
val inbox = mutableListOf<Event>()
|
||||
val received = mutableListOf<String>()
|
||||
var groupId: HexKey? = null
|
||||
|
||||
/** Group label -> nostr group id, for every group this client is in. */
|
||||
val groups = mutableMapOf<String, HexKey>()
|
||||
|
||||
/** Members this client watched join, by pubkey, since the last `clear_events`. */
|
||||
val sawJoin = mutableListOf<HexKey>()
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -94,22 +115,31 @@ class MarmotScenarioRunner(
|
||||
* step. It is the same information in the same order, just consulted when
|
||||
* this implementation needs it.
|
||||
*/
|
||||
private val publishOutcomes = mutableMapOf<String, MutableList<Boolean>>()
|
||||
private val publishOutcomes = mutableMapOf<String, MutableList<Pair<String, Boolean>>>()
|
||||
|
||||
/** Publication label -> whether OUR gate confirmed it, for `pending_resolution`. */
|
||||
private val resolved = mutableMapOf<String, Boolean>()
|
||||
|
||||
private fun preScanPublishOutcomes() {
|
||||
vector.steps.filter { it.type == "acknowledge_outbound" }.forEach { step ->
|
||||
val client = step.string("client") ?: return@forEach
|
||||
val accepted = step.string("outcome") == "accepted"
|
||||
publishOutcomes.getOrPut(client) { mutableListOf() }.add(accepted)
|
||||
publishOutcomes
|
||||
.getOrPut(client) { mutableListOf() }
|
||||
.add(step.string("publication").orEmpty() to accepted)
|
||||
}
|
||||
}
|
||||
|
||||
private fun nextOutcome(client: String): Boolean {
|
||||
val queue = publishOutcomes[client] ?: return true
|
||||
return if (queue.isEmpty()) true else queue.removeAt(0)
|
||||
if (queue.isEmpty()) return true
|
||||
val (label, accepted) = queue.removeAt(0)
|
||||
if (label.isNotEmpty()) resolved[label] = accepted
|
||||
return accepted
|
||||
}
|
||||
|
||||
suspend fun run() {
|
||||
vector.unmodelledOutcomes.firstOrNull()?.let { throw UnsupportedScenarioOutcome(it.type) }
|
||||
preScanPublishOutcomes()
|
||||
clients.values.forEach { client ->
|
||||
client.manager =
|
||||
@@ -138,31 +168,78 @@ class MarmotScenarioRunner(
|
||||
"send_app_message" -> sendAppMessage(step)
|
||||
"deliver_all" -> deliverAll()
|
||||
"tick" -> step.strings("clients").ifEmpty { vector.clients }.forEach { tick(it) }
|
||||
"in_group" -> inGroup(step)
|
||||
"clear_events" -> clearEvents(step)
|
||||
"assert" -> assertPredicate(step)
|
||||
// The publication's outcome was consumed when it was made; the step
|
||||
// itself carries no further state change.
|
||||
"acknowledge_outbound" -> Unit
|
||||
// Assertions the trace re-states; `verify()` checks them from the
|
||||
// expected observations, which is the same information.
|
||||
"observe", "observe_exact", "in_group", "assert", "clear_events", "await_quiescence" -> Unit
|
||||
"observe", "observe_exact", "await_quiescence" -> Unit
|
||||
else -> throw UnsupportedScenarioStep(step.type)
|
||||
}
|
||||
}
|
||||
|
||||
/** Run the wrapped step against the named group label. */
|
||||
private suspend fun inGroup(step: ScenarioVector.Step) {
|
||||
val action = step.step("action") ?: error("in_group without an action")
|
||||
val previous = currentGroup
|
||||
currentGroup = step.string("group") ?: DEFAULT_GROUP
|
||||
try {
|
||||
execute(action)
|
||||
} finally {
|
||||
currentGroup = previous
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reset the observation counters, as the reference's `clear_events` does.
|
||||
*
|
||||
* It is what makes a later `received_payloads` mean "since this point"
|
||||
* rather than "ever", so dropping it would make every post-clear
|
||||
* expectation fail against a list that still holds the setup traffic.
|
||||
*/
|
||||
private fun clearEvents(step: ScenarioVector.Step) {
|
||||
step.strings("clients").ifEmpty { vector.clients }.forEach {
|
||||
client(it).received.clear()
|
||||
client(it).sawJoin.clear()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Check an inline `assert` predicate.
|
||||
*
|
||||
* The only predicate our vectors use is `payload_count`, and every one of
|
||||
* them asserts a count of ZERO: it is how forward secrecy and multigroup
|
||||
* isolation are stated — a payload this client must NOT hold. That makes it
|
||||
* the highest-value assertion in the set and the last one that should be
|
||||
* skipped.
|
||||
*/
|
||||
private fun assertPredicate(step: ScenarioVector.Step) {
|
||||
val assertion = step.obj("assertion") ?: error("assert without an assertion")
|
||||
val predicate = assertion.step("predicate") ?: error("assertion without a predicate")
|
||||
when (predicate.type) {
|
||||
"payload_count" -> {
|
||||
val who = predicate.string("client") ?: error("payload_count without a client")
|
||||
val payload = predicate.string("payload").orEmpty()
|
||||
val want = predicate.int("count") ?: 0
|
||||
val got = client(who).received.count { it == payload }
|
||||
check(got == want) {
|
||||
"vector ${vector.name}: $who holds $got copies of '$payload', expected $want"
|
||||
}
|
||||
}
|
||||
|
||||
else -> throw UnsupportedScenarioStep("assert/${predicate.type}")
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun createGroup(step: ScenarioVector.Step) {
|
||||
val creator = client(step.string("creator") ?: error("create_group without a creator"))
|
||||
val name = step.string("name").orEmpty()
|
||||
val groupId = RandomInstance.bytes(32).toHexKey()
|
||||
val invitees = step.strings("invitees")
|
||||
|
||||
if (invitees.size > 1) {
|
||||
throw ScenarioBatchingDivergence(
|
||||
"create_group names ${invitees.size} invitees and the reference adds them in one " +
|
||||
"commit (epoch 1); MarmotManager.addMember stages one Add per commit, so we " +
|
||||
"would reach epoch ${invitees.size}. Both are valid MLS; the traces cannot match " +
|
||||
"until we can commit several Adds together.",
|
||||
)
|
||||
}
|
||||
|
||||
creator.manager.createCurrentProfileGroup(
|
||||
nostrGroupId = groupId,
|
||||
relays = listOf("wss://vector.invalid"),
|
||||
@@ -172,13 +249,14 @@ class MarmotScenarioRunner(
|
||||
client(admin).signer.pubKey.hexToByteArray()
|
||||
},
|
||||
)
|
||||
creator.groupId = groupId
|
||||
labelByGroupId[groupId] = currentGroup
|
||||
creator.groups[currentGroup] = groupId
|
||||
addMembers(creator, groupId, invitees)
|
||||
}
|
||||
|
||||
private suspend fun inviteMembers(step: ScenarioVector.Step) {
|
||||
val inviter = client(step.string("inviter") ?: error("invite_members without an inviter"))
|
||||
val groupId = inviter.groupId ?: error("${inviter.name} invited before joining a group")
|
||||
val groupId = inviter.groups[currentGroup] ?: error("${inviter.name} invited before joining a group")
|
||||
addMembers(inviter, groupId, step.strings("invitees"))
|
||||
}
|
||||
|
||||
@@ -187,35 +265,54 @@ class MarmotScenarioRunner(
|
||||
groupId: HexKey,
|
||||
invitees: List<String>,
|
||||
) {
|
||||
invitees.forEach { inviteeName ->
|
||||
val invitee = client(inviteeName)
|
||||
// A KeyPackage per invitee, minted on demand: the vector names
|
||||
// members, not key material.
|
||||
val bundle = invitee.manager.generateKeyPackageEvent(relays = emptyList())
|
||||
val (_, delivery) =
|
||||
inviter.manager.addMember(
|
||||
nostrGroupId = groupId,
|
||||
keyPackageEvent = bundle,
|
||||
relays = emptyList(),
|
||||
)
|
||||
if (invitees.isEmpty()) return
|
||||
|
||||
// A KeyPackage per invitee, minted on demand: the vector names
|
||||
// members, not key material.
|
||||
val bundles =
|
||||
invitees.map { inviteeName ->
|
||||
val invitee = client(inviteeName)
|
||||
invitee to invitee.manager.generateKeyPackageEvent(relays = emptyList())
|
||||
}
|
||||
|
||||
// ONE commit for the whole batch, as the reference does — N Adds in a
|
||||
// single Commit and a single Welcome carrying N EncryptedGroupSecrets.
|
||||
// Adding them one at a time would burn an epoch per invitee and the
|
||||
// traces would no longer line up.
|
||||
val (_, deliveries) =
|
||||
inviter.manager.addMembers(
|
||||
nostrGroupId = groupId,
|
||||
keyPackageEvents = bundles.map { it.second },
|
||||
relays = emptyList(),
|
||||
)
|
||||
|
||||
val byPubKey = bundles.associate { (invitee, _) -> invitee.signer.pubKey to invitee }
|
||||
deliveries.forEach { delivery ->
|
||||
// The Welcome goes straight to its recipient's inbox. Gift-wrap
|
||||
// addressing is the transport's job and the interop harness's test.
|
||||
delivery?.let { invitee.inbox.add(it.giftWrapEvent) }
|
||||
invitee.groupId = groupId
|
||||
byPubKey[delivery.recipientPubKey]?.inbox?.add(delivery.giftWrapEvent)
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun sendAppMessage(step: ScenarioVector.Step) {
|
||||
val sender = client(step.string("sender") ?: error("send_app_message without a sender"))
|
||||
val groupId = sender.groupId ?: error("${sender.name} sent before joining a group")
|
||||
val groupId = sender.groups[currentGroup] ?: error("${sender.name} sent before joining a group")
|
||||
val payload = step.string("payload").orEmpty()
|
||||
sender.manager.buildTextMessage(groupId, payload, persistOwn = false)
|
||||
// An application message does NOT go through the publish gate — it
|
||||
// advances nothing and has nothing to roll back — so the runner queues
|
||||
// it itself. Leaving that out meant every `send_app_message` built an
|
||||
// event nobody ever delivered.
|
||||
val bundle = sender.manager.buildTextMessage(groupId, payload, persistOwn = false)
|
||||
inFlight.add(sender.name to bundle.outbound.signedEvent)
|
||||
}
|
||||
|
||||
private fun deliverAll() {
|
||||
val batch = inFlight.toList()
|
||||
inFlight.clear()
|
||||
batch.forEach { (senderName, event) ->
|
||||
// Broadcast to everyone else, including clients who are not in the
|
||||
// sending group. Their engine refusing that traffic is precisely
|
||||
// what multigroup isolation asserts.
|
||||
clients.values.filter { it.name != senderName }.forEach { it.inbox.add(event) }
|
||||
}
|
||||
}
|
||||
@@ -224,9 +321,18 @@ class MarmotScenarioRunner(
|
||||
val client = client(clientName)
|
||||
val batch = client.inbox.toList()
|
||||
client.inbox.clear()
|
||||
|
||||
// Snapshot membership of the groups this client is ALREADY in, so a
|
||||
// commit processed below can be attributed as "saw N join". A group
|
||||
// joined during this tick has no before-state and contributes nothing:
|
||||
// the joiner did not watch anyone join, it arrived to a membership.
|
||||
val before = client.groups.values.associateWith { membersOf(client, it) }
|
||||
|
||||
batch.forEach { event ->
|
||||
when (val result = client.manager.ingest(event)) {
|
||||
is MarmotIngestResult.JoinedGroup -> client.groupId = result.nostrGroupId
|
||||
is MarmotIngestResult.JoinedGroup ->
|
||||
client.groups[labelByGroupId[result.nostrGroupId] ?: DEFAULT_GROUP] = result.nostrGroupId
|
||||
|
||||
is MarmotIngestResult.Message ->
|
||||
Event
|
||||
.fromJsonOrNull(result.inner.innerEventJson)
|
||||
@@ -236,35 +342,80 @@ class MarmotScenarioRunner(
|
||||
else -> Unit
|
||||
}
|
||||
}
|
||||
|
||||
before.forEach { (groupId, was) ->
|
||||
client.sawJoin.addAll(membersOf(client, groupId) - was)
|
||||
}
|
||||
}
|
||||
|
||||
private fun membersOf(
|
||||
client: VectorClient,
|
||||
groupId: HexKey,
|
||||
): Set<HexKey> =
|
||||
client.manager
|
||||
.memberPubkeys(groupId)
|
||||
.map { it.pubkey }
|
||||
.toSet()
|
||||
|
||||
/** Compare every client's end state against the vector's expected trace. */
|
||||
private fun verify() {
|
||||
val failures = mutableListOf<String>()
|
||||
|
||||
vector.pendingResolutions.forEach { expected ->
|
||||
val confirmed = resolved[expected.publication]
|
||||
val want = expected.resolution == "confirmed"
|
||||
if (confirmed == null) {
|
||||
failures.add("publication '${expected.publication}' never happened")
|
||||
} else if (confirmed != want) {
|
||||
failures.add(
|
||||
"publication '${expected.publication}' resolved " +
|
||||
"${if (confirmed) "confirmed" else "rolled_back"}, expected ${expected.resolution}",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
vector.quiescentClients.forEach { name ->
|
||||
val client = clients[name] ?: return@forEach
|
||||
if (client.inbox.isNotEmpty()) failures.add("$name still has ${client.inbox.size} events unprocessed")
|
||||
}
|
||||
if (vector.quiescentClients.isNotEmpty() && inFlight.isNotEmpty()) {
|
||||
failures.add("${inFlight.size} events are still undelivered")
|
||||
}
|
||||
|
||||
vector.observations.forEach { expected ->
|
||||
val client = clients[expected.client] ?: return@forEach
|
||||
val groupId = client.groupId
|
||||
if (groupId == null) {
|
||||
if (client.groups.isEmpty()) {
|
||||
failures.add("${expected.client} is in no group")
|
||||
return@forEach
|
||||
}
|
||||
expected.epoch?.let { want ->
|
||||
val got = client.manager.groupEpoch(groupId)
|
||||
if (got != want) failures.add("${expected.client} epoch $got, expected $want")
|
||||
}
|
||||
expected.memberCount?.let { want ->
|
||||
val got = client.manager.memberCount(groupId)
|
||||
if (got != want) failures.add("${expected.client} has $got members, expected $want")
|
||||
}
|
||||
expected.groupName?.let { want ->
|
||||
val got = client.manager.groupView(groupId)?.name
|
||||
if (got != want) failures.add("${expected.client} group name '$got', expected '$want'")
|
||||
// A per-group fact stated once applies to EVERY group the client
|
||||
// holds — the isolation vectors put a client in several groups and
|
||||
// state one epoch and one member count for all of them.
|
||||
client.groups.forEach { (label, groupId) ->
|
||||
expected.epoch?.let { want ->
|
||||
val got = client.manager.groupEpoch(groupId)
|
||||
if (got != want) failures.add("${expected.client}[$label] epoch $got, expected $want")
|
||||
}
|
||||
expected.memberCount?.let { want ->
|
||||
val got = client.manager.memberCount(groupId)
|
||||
if (got != want) failures.add("${expected.client}[$label] has $got members, expected $want")
|
||||
}
|
||||
expected.groupName?.let { want ->
|
||||
val got = client.manager.groupView(groupId)?.name
|
||||
if (got != want) failures.add("${expected.client}[$label] group name '$got', expected '$want'")
|
||||
}
|
||||
}
|
||||
if (expected.receivedPayloads.isNotEmpty()) {
|
||||
val got = client.received.sorted()
|
||||
val want = expected.receivedPayloads.sorted()
|
||||
if (got != want) failures.add("${expected.client} received $got, expected $want")
|
||||
}
|
||||
expected.addedMembers?.let { want ->
|
||||
val got = client.sawJoin.mapNotNull { pubkey -> clients.values.firstOrNull { it.signer.pubKey == pubkey }?.name }
|
||||
if (got.sorted() != want.sorted()) {
|
||||
failures.add("${expected.client} saw $got join, expected $want")
|
||||
}
|
||||
}
|
||||
}
|
||||
check(failures.isEmpty()) {
|
||||
"vector ${vector.name} diverged from its expected trace:\n " + failures.joinToString("\n ")
|
||||
@@ -275,5 +426,8 @@ class MarmotScenarioRunner(
|
||||
|
||||
private companion object {
|
||||
const val CHAT_KIND = 9
|
||||
|
||||
/** The label for a vector that never says `in_group` — most of them. */
|
||||
const val DEFAULT_GROUP = "default"
|
||||
}
|
||||
}
|
||||
|
||||
+61
-23
@@ -22,6 +22,7 @@ package com.vitorpamplona.amethyst.commons.marmot.scenario
|
||||
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.test.fail
|
||||
@@ -74,32 +75,53 @@ class MarmotScenarioVectorTest {
|
||||
fun invitePublishFail() = replay("invite-publish-fail.v1.json")
|
||||
|
||||
/**
|
||||
* The two vectors we cannot replay, and exactly why.
|
||||
*
|
||||
* Both create a group with several invitees. The reference adds them in one
|
||||
* commit, so the group is at epoch 1; our `addMember` stages one Add per
|
||||
* commit, so we would reach epoch 2. Neither is a protocol error — a commit
|
||||
* per Add is valid MLS and any peer processes it — but the traces cannot
|
||||
* match until we can commit several Adds together, and creating a group
|
||||
* costs us an extra round trip per invitee until then.
|
||||
*
|
||||
* Asserted rather than deleted so the divergence stays visible: the day
|
||||
* batched adds land, this test fails and these two move up to [replay].
|
||||
* The vectors whose `create_group` names several invitees. The reference
|
||||
* adds them all in ONE commit — epoch 1, one Welcome carrying an
|
||||
* EncryptedGroupSecrets per invitee — and so do we now, so the traces line
|
||||
* up. They were refused as a batching divergence until batched Adds landed.
|
||||
*/
|
||||
@Test
|
||||
fun aMultiInviteeCreateStillDivergesOnBatching() {
|
||||
listOf(
|
||||
"three-client-message-exchange.v1.json",
|
||||
"convergence-committer-selected.v1.json",
|
||||
"conversation.v1.json",
|
||||
).forEach { name ->
|
||||
val thrown =
|
||||
assertFailsWith<ScenarioBatchingDivergence>("$name should still diverge on batching") {
|
||||
runBlocking { MarmotScenarioRunner(load(name)).run() }
|
||||
}
|
||||
fun threeClientMessageExchange() = replay("three-client-message-exchange.v1.json")
|
||||
|
||||
@Test
|
||||
fun conversation() = replay("conversation.v1.json")
|
||||
|
||||
/**
|
||||
* `convergence-committer-selected` concludes with a `convergence_decision`
|
||||
* — which tip the client picked, under which rule, and whether the witness
|
||||
* quorum was met. Our convergence engine makes that decision but does not
|
||||
* report it in those terms, so there is nothing to compare against and the
|
||||
* runner refuses the vector rather than passing it on the observations it
|
||||
* can check.
|
||||
*
|
||||
* Asserted rather than deleted so the gap stays visible: the day the engine
|
||||
* exposes its decision, this test fails and the vector moves up to [replay].
|
||||
*/
|
||||
@Test
|
||||
fun aConvergenceDecisionIsStillUnmodelled() {
|
||||
val thrown =
|
||||
assertFailsWith<UnsupportedScenarioOutcome> {
|
||||
runBlocking { MarmotScenarioRunner(load("convergence-committer-selected.v1.json")).run() }
|
||||
}
|
||||
assertEquals("convergence_decision", thrown.outcomeType)
|
||||
}
|
||||
|
||||
/**
|
||||
* Every vector must parse to at least one thing to check.
|
||||
*
|
||||
* Two expectation shapes ship in this set — `expected_trace.observations`
|
||||
* and `expected_outcomes` — and reading only the first left seven of the
|
||||
* nine vectors with nothing to compare against, replaying their steps and
|
||||
* reporting green. A vector that asserts nothing is worse than a missing
|
||||
* vector, so this guards the parser rather than any one scenario.
|
||||
*/
|
||||
@Test
|
||||
fun everyVectorStatesSomethingToCheck() {
|
||||
VECTORS.forEach { name ->
|
||||
val vector = load(name)
|
||||
assertTrue(
|
||||
thrown.message.orEmpty().contains("one commit"),
|
||||
"the divergence must say what it is: ${thrown.message}",
|
||||
vector.observations.isNotEmpty() || vector.unmodelledOutcomes.isNotEmpty(),
|
||||
"$name parsed to zero expectations — it would pass without checking anything",
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -115,4 +137,20 @@ class MarmotScenarioVectorTest {
|
||||
"unexpected conformance version ${vector.conformanceVersion}",
|
||||
)
|
||||
}
|
||||
|
||||
private companion object {
|
||||
/** Every vector this suite ships, so the parser guard covers them all. */
|
||||
val VECTORS =
|
||||
listOf(
|
||||
"invite-member.v1.json",
|
||||
"current-profile-required-set.v1.json",
|
||||
"latecomer-forward-secrecy.v1.json",
|
||||
"multigroup-isolation.v1.json",
|
||||
"publish-fail.v1.json",
|
||||
"invite-publish-fail.v1.json",
|
||||
"three-client-message-exchange.v1.json",
|
||||
"conversation.v1.json",
|
||||
"convergence-committer-selected.v1.json",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
+131
-29
@@ -45,6 +45,9 @@ class ScenarioVector(
|
||||
val clients: List<String>,
|
||||
val steps: List<Step>,
|
||||
val observations: List<Observation>,
|
||||
val pendingResolutions: List<PendingResolution> = emptyList(),
|
||||
val quiescentClients: List<String> = emptyList(),
|
||||
val unmodelledOutcomes: List<UnmodelledOutcome> = emptyList(),
|
||||
) {
|
||||
class Step(
|
||||
val type: String,
|
||||
@@ -52,10 +55,27 @@ class ScenarioVector(
|
||||
) {
|
||||
fun string(key: String): String? = (raw[key] as? JsonPrimitive)?.takeIf { it.isString }?.content
|
||||
|
||||
fun int(key: String): Int? = (raw[key] as? JsonPrimitive)?.content?.toIntOrNull()
|
||||
|
||||
fun strings(key: String): List<String> =
|
||||
(raw[key] as? JsonArray)
|
||||
?.mapNotNull { (it as? JsonPrimitive)?.takeIf { p -> p.isString }?.content }
|
||||
.orEmpty()
|
||||
|
||||
/**
|
||||
* A step nested inside this one, as its own [Step].
|
||||
*
|
||||
* `in_group` is a wrapper: it names a group label and carries the real
|
||||
* step under `action`. Treating it as a leaf silently skipped every
|
||||
* create and every message inside it.
|
||||
*/
|
||||
fun step(key: String): Step? =
|
||||
(raw[key] as? JsonObject)?.let {
|
||||
Step((it["type"] as? JsonPrimitive)?.content.orEmpty(), it)
|
||||
}
|
||||
|
||||
/** A nested object read as a [Step] so the same accessors work on it. */
|
||||
fun obj(key: String): Step? = (raw[key] as? JsonObject)?.let { Step(key, it) }
|
||||
}
|
||||
|
||||
/** What one client's state must look like when the script says to look. */
|
||||
@@ -65,6 +85,31 @@ class ScenarioVector(
|
||||
val memberCount: Int?,
|
||||
val groupName: String?,
|
||||
val receivedPayloads: List<String>,
|
||||
/**
|
||||
* Members this client must have SEEN JOIN, by client name. Null when
|
||||
* the vector does not state it — which is not the same as an empty
|
||||
* list, and an empty list is itself an assertion.
|
||||
*/
|
||||
val addedMembers: List<String>? = null,
|
||||
)
|
||||
|
||||
/**
|
||||
* A publication the script named, and how the reference says it ended.
|
||||
*
|
||||
* `confirmed` means the relay accepted it and the commit became canonical;
|
||||
* `rolled_back` means it did not and the committer stayed where it was.
|
||||
* Checking these is the only thing that proves our publish-before-apply
|
||||
* gate resolved the same way the reference's did.
|
||||
*/
|
||||
class PendingResolution(
|
||||
val client: String,
|
||||
val publication: String,
|
||||
val resolution: String,
|
||||
)
|
||||
|
||||
/** The outcome types this runner has no check for, named so it can refuse. */
|
||||
class UnmodelledOutcome(
|
||||
val type: String,
|
||||
)
|
||||
|
||||
companion object {
|
||||
@@ -78,29 +123,90 @@ class ScenarioVector(
|
||||
val obj = element as JsonObject
|
||||
Step(obj.getValue("type").jsonPrimitive.content, obj)
|
||||
}
|
||||
val observations =
|
||||
// TWO vector shapes ship side by side. The older one nests a
|
||||
// trace under `expected_trace.observations`; the newer one lists
|
||||
// typed entries under `expected_outcomes`. Reading only the first
|
||||
// meant SEVEN of the nine vectors here parsed to zero expectations
|
||||
// and "passed" without checking anything — the exact failure this
|
||||
// runner exists to avoid.
|
||||
val traced =
|
||||
((root["expected_trace"] as? JsonObject)?.get("observations") as? JsonArray)
|
||||
?.map { element ->
|
||||
val obj = element as JsonObject
|
||||
Observation(
|
||||
client = obj.getValue("client").jsonPrimitive.content,
|
||||
epoch = (obj["epoch"] as? JsonPrimitive)?.content?.toLongOrNull(),
|
||||
memberCount = (obj["member_count"] as? JsonPrimitive)?.content?.toIntOrNull(),
|
||||
groupName = (obj["group_name"] as? JsonPrimitive)?.takeIf { it.isString }?.content,
|
||||
receivedPayloads =
|
||||
(obj["received_payloads"] as? JsonArray)
|
||||
?.mapNotNull { (it as? JsonPrimitive)?.content }
|
||||
.orEmpty(),
|
||||
?.map { observationOf(it as JsonObject) }
|
||||
.orEmpty()
|
||||
|
||||
val outcomes = (root["expected_outcomes"] as? JsonArray).orEmpty()
|
||||
val stated = mutableListOf<Observation>()
|
||||
val resolutions = mutableListOf<PendingResolution>()
|
||||
val quiescent = mutableListOf<String>()
|
||||
val unmodelled = mutableListOf<UnmodelledOutcome>()
|
||||
|
||||
outcomes.forEach { element ->
|
||||
val obj = element as JsonObject
|
||||
when ((obj["type"] as? JsonPrimitive)?.content) {
|
||||
"client_state" -> stated.add(observationOf(obj))
|
||||
|
||||
// A converged set states the same per-client facts for
|
||||
// several clients at once.
|
||||
"clients_converged" ->
|
||||
(obj["clients"] as? JsonArray).orEmpty().forEach { name ->
|
||||
stated.add(
|
||||
Observation(
|
||||
client = (name as JsonPrimitive).content,
|
||||
epoch = (obj["epoch"] as? JsonPrimitive)?.content?.toLongOrNull(),
|
||||
memberCount = (obj["member_count"] as? JsonPrimitive)?.content?.toIntOrNull(),
|
||||
groupName = null,
|
||||
receivedPayloads = emptyList(),
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
"pending_resolution" ->
|
||||
resolutions.add(
|
||||
PendingResolution(
|
||||
client = obj.getValue("client").jsonPrimitive.content,
|
||||
publication = obj.getValue("pending").jsonPrimitive.content,
|
||||
resolution = obj.getValue("resolution").jsonPrimitive.content,
|
||||
),
|
||||
)
|
||||
}.orEmpty()
|
||||
|
||||
"no_pending_work" ->
|
||||
(obj["clients"] as? JsonArray).orEmpty().forEach {
|
||||
quiescent.add((it as JsonPrimitive).content)
|
||||
}
|
||||
|
||||
else ->
|
||||
unmodelled.add(
|
||||
UnmodelledOutcome((obj["type"] as? JsonPrimitive)?.content.orEmpty()),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
return ScenarioVector(
|
||||
name = (root["scenario_name"] as JsonPrimitive).content,
|
||||
conformanceVersion = (root["conformance_version"] as? JsonPrimitive)?.content.orEmpty(),
|
||||
clients = (scenario["clients"] as JsonArray).map { (it as JsonPrimitive).content },
|
||||
steps = steps,
|
||||
observations = observations,
|
||||
observations = traced + stated,
|
||||
pendingResolutions = resolutions,
|
||||
quiescentClients = quiescent,
|
||||
unmodelledOutcomes = unmodelled,
|
||||
)
|
||||
}
|
||||
|
||||
private fun observationOf(obj: JsonObject) =
|
||||
Observation(
|
||||
client = obj.getValue("client").jsonPrimitive.content,
|
||||
epoch = (obj["epoch"] as? JsonPrimitive)?.content?.toLongOrNull(),
|
||||
memberCount = (obj["member_count"] as? JsonPrimitive)?.content?.toIntOrNull(),
|
||||
groupName = (obj["group_name"] as? JsonPrimitive)?.takeIf { it.isString }?.content,
|
||||
receivedPayloads =
|
||||
(obj["received_payloads"] as? JsonArray)
|
||||
?.mapNotNull { (it as? JsonPrimitive)?.content }
|
||||
.orEmpty(),
|
||||
addedMembers =
|
||||
(obj["added_members"] as? JsonArray)
|
||||
?.mapNotNull { (it as? JsonPrimitive)?.content },
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -113,20 +219,16 @@ class UnsupportedScenarioStep(
|
||||
)
|
||||
|
||||
/**
|
||||
* Raised where our engine cannot reach the vector's trace because it batches
|
||||
* differently, not because either side is wrong.
|
||||
* Raised when a vector states an expected outcome this runner cannot check.
|
||||
*
|
||||
* The one case today: the reference adds every invitee named by `create_group`
|
||||
* in a SINGLE commit, so a group created with two invitees is at epoch 1. Our
|
||||
* `MarmotManager.addMember` stages one Add per commit, so the same group
|
||||
* reaches epoch 2. Both are valid MLS — a commit per Add is not a protocol
|
||||
* error, and a peer processes either — but the epoch numbers differ, and so
|
||||
* does the round-trip cost of creating a group.
|
||||
*
|
||||
* Kept as its own signal rather than folded into a trace mismatch: a divergence
|
||||
* we understand and have chosen not to fix yet should not read like a bug we
|
||||
* have not noticed.
|
||||
* Same contract as [UnsupportedScenarioStep], one level up: a vector whose
|
||||
* conclusion we cannot evaluate has not been conformed to, however cleanly its
|
||||
* steps replayed. Silently dropping the outcome would turn the vector into an
|
||||
* expensive no-op that reports green.
|
||||
*/
|
||||
class ScenarioBatchingDivergence(
|
||||
message: String,
|
||||
) : IllegalStateException(message)
|
||||
class UnsupportedScenarioOutcome(
|
||||
val outcomeType: String,
|
||||
) : IllegalStateException(
|
||||
"expected outcome '$outcomeType' has no check in this runner — the vector is refused " +
|
||||
"rather than passed on the outcomes that happen to be modelled",
|
||||
)
|
||||
|
||||
@@ -4167,8 +4167,24 @@ class MlsGroup private constructor(
|
||||
* [CommitResult.preCommitExporterSecret] is the key the outer kind:445
|
||||
* MUST be encrypted with (RFC 9420 §12.4 + MDK parity).
|
||||
*/
|
||||
fun addMember(keyPackageBytes: ByteArray): CommitResult {
|
||||
proposeAdd(keyPackageBytes)
|
||||
fun addMember(keyPackageBytes: ByteArray): CommitResult = addMembers(listOf(keyPackageBytes))
|
||||
|
||||
/**
|
||||
* Add several members in ONE commit.
|
||||
*
|
||||
* Not a convenience wrapper over [addMember] — a commit per Add costs an
|
||||
* epoch and a publish round trip each, and every existing member processes
|
||||
* each one. The Welcome already carries a separate `EncryptedGroupSecrets`
|
||||
* per added member, keyed by KeyPackage reference (RFC 9420 §12.4.3.1), so
|
||||
* one commit serves all of them and each joiner finds its own secrets.
|
||||
*
|
||||
* The reference implementation adds every invitee named at group creation
|
||||
* this way, which is why a group it creates with two invitees sits at epoch
|
||||
* 1 while ours used to reach epoch 2.
|
||||
*/
|
||||
fun addMembers(keyPackagesBytes: List<ByteArray>): CommitResult {
|
||||
require(keyPackagesBytes.isNotEmpty()) { "addMembers needs at least one KeyPackage" }
|
||||
keyPackagesBytes.forEach { proposeAdd(it) }
|
||||
return commit()
|
||||
}
|
||||
|
||||
|
||||
+7
-1
@@ -397,7 +397,13 @@ class MlsGroupManager(
|
||||
suspend fun stageAddMember(
|
||||
nostrGroupId: HexKey,
|
||||
keyPackageBytes: ByteArray,
|
||||
): StagedCommit = stage(nostrGroupId) { it.addMember(keyPackageBytes) }
|
||||
): StagedCommit = stageAddMembers(nostrGroupId, listOf(keyPackageBytes))
|
||||
|
||||
/** Stage several Adds as ONE commit. See [MlsGroup.addMembers]. */
|
||||
suspend fun stageAddMembers(
|
||||
nostrGroupId: HexKey,
|
||||
keyPackagesBytes: List<ByteArray>,
|
||||
): StagedCommit = stage(nostrGroupId) { it.addMembers(keyPackagesBytes) }
|
||||
|
||||
/** Stage a Remove. See [StagedCommit] for why this does not apply. */
|
||||
suspend fun stageRemoveMember(
|
||||
|
||||
Reference in New Issue
Block a user