mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 11:48:24 +00:00
feat(cordn): headless app layer -- group manager, coordinator config, §8 exposure, amy verbs
Stage 4 of the interop plan, minus the UI. `commons/.../cordn/` now holds what every front end needs: which coordinator we talk to, whether it is answering, what it learns about us, and the group lifecycle on top of it. `CordnGroupManager` is the cordn counterpart of MarmotManager and shares no code with it, as Stage 1 predicted: Marmot's is keyed on the Nostr group id throughout, and kinds 443/444/445 plus h-tag subscriptions have no cordn analogue -- here there is one coordinator, one ordered stream per gid, and a cursor. What the two share is the RFC 9420 engine underneath. Its store is separate from MlsGroupStateStore for the same reason: `nostrGroupId` and `gid` look alike and mean different things. `ICoordinator` (quartz) makes the eleven tools a contract so the cursor loop and the manager can be driven without a relay. It adds no behaviour and deliberately exposes no identity parameter -- which key signs which call is fixed by CoordinatorMethod per spec/00.md §8, and a substitute implementation cannot widen that. `CordnExposure` turns §8 into a model instead of a comment. The plan requires that exposure be "surfaced, not buried" because a cordn group and a Marmot group are not privacy-equivalent and a user cannot infer that from either one looking like a group chat. GroupExposure computes it from real state -- notably §8.2's linked-group count, which is the part nobody guesses -- so a UI renders facts rather than prose. Nothing renders it yet; that is the remaining Stage 4 work. Four bugs the tests found, each documented in the plan: - CommitResult.commitBytes is the bare RFC 9420 Commit struct, not an MLSMessage. Posting it produces a payload no receiver can parse. - The pre-commit epoch rule (spec/03.md §5) is invisible to a two-party test: sealing a Commit under the epoch it creates still passes create -> invite -> join -> read, because the joiner's Welcome cursor skips the only Commit in the stream. It takes a third member to catch. Confirmed by mutation -- the new three-member test is the only one that dies. - A Welcome does not carry the gid. spec/03.md §2 decouples it from the MLS group_id; the reference client closes the gap with `group_id = utf8(gid)`, which staircase's fixtures confirm. We follow the convention and treat it as an observation: a non-UTF-8 group id yields a named skip, not a mojibake key that fetches nothing forever. MlsGroup.create gained an optional groupId. - Public-framed handshake messages interoperate. Checked in their code rather than assumed: mlsMessages.ts:68 accepts wireformat 1 and 2 and hands both to ts-mls's processMessage. `amy cordn ref encode/decode` and `amy cordn exposure` ship; ref encode produces byte-identical output to cordn's own vector. The coordinator-driving verbs do not, and the reason is in the file: Tier B is blocked on licensing (§7), so there is nothing to exercise them against, and unexercised coordinator verbs are a guess with a command-line interface. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
@@ -27,6 +27,7 @@ import com.vitorpamplona.amethyst.cli.commands.Bolt12Commands
|
||||
import com.vitorpamplona.amethyst.cli.commands.BunkerCommand
|
||||
import com.vitorpamplona.amethyst.cli.commands.BuzzCommands
|
||||
import com.vitorpamplona.amethyst.cli.commands.ConcordCommands
|
||||
import com.vitorpamplona.amethyst.cli.commands.CordnCommands
|
||||
import com.vitorpamplona.amethyst.cli.commands.CountCommand
|
||||
import com.vitorpamplona.amethyst.cli.commands.CreateCommand
|
||||
import com.vitorpamplona.amethyst.cli.commands.DebitCommands
|
||||
@@ -333,6 +334,7 @@ private suspend fun dispatch(argv: Array<String>): Int {
|
||||
FofCommand.dispatch(dataDir, tail)
|
||||
}
|
||||
"concord" -> ConcordCommands.dispatch(dataDir, tail)
|
||||
"cordn" -> CordnCommands.dispatch(tail)
|
||||
else -> {
|
||||
Output.error("bad_args", "unknown subcommand: $head")
|
||||
printVerbList()
|
||||
@@ -355,7 +357,7 @@ private fun printVerbList() {
|
||||
| primitives: decode encode verify key filter nip kind pow namecoin
|
||||
| events: event publish fetch subscribe count sync encrypt decrypt gift
|
||||
| social: notes profile follow unfollow search zap dm outbox
|
||||
| groups: marmot relaygroup concord geochat
|
||||
| groups: marmot relaygroup concord cordn geochat
|
||||
| relays: relay admin serve store
|
||||
| trust: graperank fof
|
||||
| media/sites: blossom nsite napplet podcast podcast20 git
|
||||
@@ -878,6 +880,12 @@ private fun printUsage() {
|
||||
| concord revoke COMMUNITY TOKEN|URL retire a link you minted (vsk=9 tombstone)
|
||||
| concord join URL redeem an invite link and save the community
|
||||
|
|
||||
| cordn ref encode --gid GID [--coordinator PK] [--relay URL[,URL]]
|
||||
| build a cordn1… group reference
|
||||
| cordn ref decode REF read one back
|
||||
| cordn exposure --coordinator PK [--groups N] [--published]
|
||||
| what that coordinator would learn (spec/00.md §8)
|
||||
|
|
||||
|Local event store (shared, under `<data-dir>/shared/`):
|
||||
| Backend selected by AMY_STORE: sqlite (default; `shared/events.db`)
|
||||
| or fs (`AMY_STORE=fs`; the `shared/events-store/` tree). SQLite is
|
||||
|
||||
@@ -0,0 +1,207 @@
|
||||
/*
|
||||
* 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.amethyst.cli.commands
|
||||
|
||||
import com.vitorpamplona.amethyst.cli.Args
|
||||
import com.vitorpamplona.amethyst.cli.Output
|
||||
import com.vitorpamplona.amethyst.commons.cordn.CoordinatorConfig
|
||||
import com.vitorpamplona.amethyst.commons.cordn.ExposureNote
|
||||
import com.vitorpamplona.amethyst.commons.cordn.GroupExposure
|
||||
import com.vitorpamplona.quartz.cordn.appGroupRef.CordnGroupRef
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
|
||||
/**
|
||||
* `amy cordn …` — the cordn (MLS-over-an-MCP-coordinator) surface.
|
||||
*
|
||||
* What is here is what works without a coordinator: the `cordn1…` group ref
|
||||
* codec, and the §8 metadata-exposure model. Both are what an interop script
|
||||
* actually needs — a ref is the one cordn artifact a human copies by hand, and
|
||||
* `exposure` is how a script (or a person) checks what a coordinator would
|
||||
* learn before joining anything.
|
||||
*
|
||||
* The verbs that drive a coordinator (`publish`, `invite`, `send`, `sync`) are
|
||||
* not here yet, and the reason is worth stating: there is nothing to test them
|
||||
* against. The reference coordinator is unlicensed — see
|
||||
* `quartz/plans/2026-09-17-cordn-interop.md` §7 — so Tier B cannot be built on
|
||||
* it, and shipping unexercised coordinator verbs would be shipping a guess.
|
||||
* [com.vitorpamplona.amethyst.commons.cordn.CordnGroupManager] holds that logic
|
||||
* and is covered against an in-memory coordinator in `commons`.
|
||||
*/
|
||||
object CordnCommands {
|
||||
val USAGE: String =
|
||||
"""
|
||||
|cordn (MLS group chat over an MCP coordinator):
|
||||
| cordn ref encode --gid GID [--coordinator PK] build a cordn1… group reference
|
||||
| [--relay URL[,URL…]]
|
||||
| cordn ref decode REF read one back
|
||||
| cordn exposure --coordinator PK what that coordinator would learn
|
||||
| [--groups N] [--published] (spec/00.md §8)
|
||||
|
|
||||
|A group ref is a locator, not an invitation: holding one lets you ASK to
|
||||
|join, it does not make you a member. Relays say where to reach the
|
||||
|coordinator and are meaningless without --coordinator.
|
||||
""".trimMargin()
|
||||
|
||||
suspend fun dispatch(tail: Array<String>): Int =
|
||||
route(
|
||||
"cordn",
|
||||
tail,
|
||||
"cordn <ref|exposure>",
|
||||
help = USAGE,
|
||||
routes =
|
||||
mapOf(
|
||||
"ref" to { rest -> ref(rest) },
|
||||
"exposure" to { rest -> exposure(rest) },
|
||||
),
|
||||
)
|
||||
|
||||
private suspend fun ref(tail: Array<String>): Int =
|
||||
route(
|
||||
"cordn ref",
|
||||
tail,
|
||||
"cordn ref <encode|decode>",
|
||||
help = USAGE,
|
||||
routes =
|
||||
mapOf(
|
||||
"encode" to { rest -> encode(rest) },
|
||||
"decode" to { rest -> decode(rest) },
|
||||
),
|
||||
)
|
||||
|
||||
private fun encode(tail: Array<String>): Int {
|
||||
val args = Args(tail)
|
||||
val gid = args.flag("gid") ?: return Output.error("bad_args", "cordn ref encode needs --gid")
|
||||
val coordinator = args.flag("coordinator")
|
||||
// The flag map collapses repeats, so several relays arrive comma-separated.
|
||||
val relays =
|
||||
args
|
||||
.flag("relay")
|
||||
?.split(",")
|
||||
?.map { it.trim() }
|
||||
?.filter { it.isNotEmpty() }
|
||||
.orEmpty()
|
||||
args.rejectUnknown()
|
||||
|
||||
val normalized =
|
||||
relays.map { url ->
|
||||
RelayUrlNormalizer.normalizeOrNull(url)
|
||||
?: return Output.error("bad_args", "not a relay URL: $url")
|
||||
}
|
||||
|
||||
val ref =
|
||||
try {
|
||||
CordnGroupRef(gid, coordinator, normalized.map { it.url })
|
||||
} catch (e: IllegalArgumentException) {
|
||||
return Output.error("bad_args", e.message ?: "invalid group reference")
|
||||
}
|
||||
|
||||
Output.emit(
|
||||
mapOf(
|
||||
"ref" to ref.encode(),
|
||||
"gid" to ref.gid,
|
||||
"coordinator" to ref.coordinatorPubKey,
|
||||
"relays" to ref.relays,
|
||||
),
|
||||
)
|
||||
return 0
|
||||
}
|
||||
|
||||
private fun decode(tail: Array<String>): Int {
|
||||
val args = Args(tail)
|
||||
val encoded = args.positional.firstOrNull() ?: return Output.error("bad_args", "cordn ref decode needs a cordn1… reference")
|
||||
args.rejectUnknown()
|
||||
|
||||
val ref =
|
||||
try {
|
||||
CordnGroupRef.decode(encoded)
|
||||
} catch (e: IllegalArgumentException) {
|
||||
return Output.error("bad_ref", e.message ?: "not a cordn group reference")
|
||||
}
|
||||
|
||||
Output.emit(
|
||||
mapOf(
|
||||
"gid" to ref.gid,
|
||||
"coordinator" to ref.coordinatorPubKey,
|
||||
"relays" to ref.relays,
|
||||
// A ref carrying no coordinator is legal (spec §2) and means the
|
||||
// recipient must already know who serves this group.
|
||||
"reachable" to (ref.coordinatorPubKey != null && ref.relays.isNotEmpty()),
|
||||
),
|
||||
)
|
||||
return 0
|
||||
}
|
||||
|
||||
private fun exposure(tail: Array<String>): Int {
|
||||
val args = Args(tail)
|
||||
val coordinator = args.flag("coordinator") ?: return Output.error("bad_args", "cordn exposure needs --coordinator")
|
||||
val groups = args.flag("groups")?.toIntOrNull() ?: 1
|
||||
val published = "published" in args.booleans
|
||||
val fromLink = "from-link" in args.booleans
|
||||
args.rejectUnknown("published", "from-link")
|
||||
|
||||
if (groups < 1) return Output.error("bad_args", "--groups must be at least 1")
|
||||
|
||||
val exposure =
|
||||
GroupExposure(
|
||||
coordinator = coordinator,
|
||||
linkedGroupCount = groups,
|
||||
joinedFromShareLink = fromLink,
|
||||
publishedKeyPackage = published,
|
||||
// Our transport pins CEP-4 encryption to REQUIRED and fails
|
||||
// closed (§8.6), so this is a property of the client, not a
|
||||
// setting a coordinator can talk us out of.
|
||||
encryptionPinned = true,
|
||||
)
|
||||
|
||||
Output.emit(
|
||||
mapOf(
|
||||
"coordinator" to coordinator,
|
||||
"content" to exposure.content.name,
|
||||
"membership" to exposure.membership.name,
|
||||
"messaging" to exposure.messaging.name,
|
||||
"linked_groups" to exposure.linkedGroupCount,
|
||||
"differs_from_marmot" to exposure.differsFromMarmot(),
|
||||
"notes" to exposure.notes().map { it.name to describe(it) }.toMap(),
|
||||
),
|
||||
)
|
||||
return 0
|
||||
}
|
||||
|
||||
/** One line per note, for the human-readable half of the dual output. */
|
||||
private fun describe(note: ExposureNote): String =
|
||||
when (note) {
|
||||
ExposureNote.MEMBERSHIP_IS_IDENTIFIED ->
|
||||
"admission names real npubs on both ends, so the coordinator sees who is in this group"
|
||||
ExposureNote.GROUPS_LINKED_BY_SESSION ->
|
||||
"one throwaway key posts and fetches for every group you have here, linking them to each other"
|
||||
ExposureNote.SINGLE_OPERATOR_HOLDS_HISTORY ->
|
||||
"one operator holds the complete ordered history of every group it serves"
|
||||
ExposureNote.PUBLICATION_IS_A_SIGNED_RECORD ->
|
||||
"your published KeyPackage is a signed, re-servable record that this account uses cordn"
|
||||
ExposureNote.MESSAGE_SIZES_UNPADDED ->
|
||||
"sealed payloads are not padded, so message sizes are visible"
|
||||
ExposureNote.ENCRYPTION_NOT_PINNED ->
|
||||
"BUG: this client is not pinning transport encryption; report it"
|
||||
}
|
||||
|
||||
/** The coordinator a ref points at, for callers wiring one up. */
|
||||
fun coordinatorOf(ref: CordnGroupRef): CoordinatorConfig? = CoordinatorConfig.from(ref)
|
||||
}
|
||||
+142
@@ -0,0 +1,142 @@
|
||||
/*
|
||||
* 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.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.appGroupRef.CordnGroupRef
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
|
||||
/**
|
||||
* One coordinator this account talks to.
|
||||
*
|
||||
* A coordinator is not a relay and the difference matters to the user: relays
|
||||
* are interchangeable and redundant, a coordinator is the single authority for
|
||||
* the groups it serves. Losing it loses the ordering; a second one does not
|
||||
* mirror the first. So this is a named, first-class thing a user chooses,
|
||||
* never a URL buried in settings.
|
||||
*
|
||||
* [relays] is where the coordinator is *reachable* — the ContextVM kind-25910
|
||||
* traffic goes over Nostr, so the coordinator has no address of its own beyond
|
||||
* its pubkey (§8.5: the relay sees the traffic pattern, the coordinator never
|
||||
* sees an IP).
|
||||
*/
|
||||
data class CoordinatorConfig(
|
||||
val pubKey: HexKey,
|
||||
val relays: List<NormalizedRelayUrl>,
|
||||
/** How this account came to know about this coordinator. */
|
||||
val origin: Origin = Origin.MANUAL,
|
||||
/** What the user calls it. Never a claim — a coordinator cannot prove a name. */
|
||||
val label: String? = null,
|
||||
) {
|
||||
init {
|
||||
require(pubKey.length == PUBKEY_HEX_LENGTH) { "a coordinator pubkey is 32 bytes of hex" }
|
||||
require(relays.isNotEmpty()) { "a coordinator with no relays cannot be reached" }
|
||||
}
|
||||
|
||||
/** Where a coordinator came from, because it changes how much to trust it. */
|
||||
enum class Origin {
|
||||
/** Typed or pasted by the user. */
|
||||
MANUAL,
|
||||
|
||||
/** Read out of a `cordn1…` group ref someone shared. */
|
||||
GROUP_REF,
|
||||
|
||||
/** The application default. */
|
||||
DEFAULT,
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val PUBKEY_HEX_LENGTH = 64
|
||||
|
||||
/**
|
||||
* The coordinator a shared `cordn1…` ref points at, or null when the
|
||||
* ref carries only a `gid`.
|
||||
*
|
||||
* A ref without a coordinator is not broken — §2 makes both optional —
|
||||
* it just means the recipient has to already know which coordinator
|
||||
* serves that group.
|
||||
*/
|
||||
fun from(ref: CordnGroupRef): CoordinatorConfig? {
|
||||
val pubKey = ref.coordinatorPubKey ?: return null
|
||||
val relays = ref.relays.mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) }
|
||||
if (relays.isEmpty()) return null
|
||||
return CoordinatorConfig(pubKey, relays, Origin.GROUP_REF)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether a coordinator is answering, kept per coordinator.
|
||||
*
|
||||
* Deliberately thin. This is not a health *check* — nothing here polls, because
|
||||
* a poll is a call, and every call to a coordinator is metadata (§8). It
|
||||
* records what the calls the app was making anyway have observed.
|
||||
*/
|
||||
class CoordinatorHealth {
|
||||
private val _state = MutableStateFlow(State())
|
||||
|
||||
val state: StateFlow<State> = _state.asStateFlow()
|
||||
|
||||
data class State(
|
||||
val lastSuccessAt: Long? = null,
|
||||
val lastFailureAt: Long? = null,
|
||||
val lastFailure: String? = null,
|
||||
/** Failures since the last success. Resets on any success. */
|
||||
val consecutiveFailures: Int = 0,
|
||||
) {
|
||||
/**
|
||||
* Nothing has worked since the last success, repeatedly.
|
||||
*
|
||||
* A single failure is a network blip and worth no UI at all; the
|
||||
* threshold is what separates "retrying" from "tell the user their
|
||||
* groups are not syncing".
|
||||
*/
|
||||
val isDown: Boolean get() = consecutiveFailures >= DOWN_AFTER
|
||||
|
||||
/** True before the first call of the session — not the same as down. */
|
||||
val isUnknown: Boolean get() = lastSuccessAt == null && lastFailureAt == null
|
||||
}
|
||||
|
||||
fun recordSuccess(atSeconds: Long) {
|
||||
_state.value = _state.value.copy(lastSuccessAt = atSeconds, consecutiveFailures = 0, lastFailure = null)
|
||||
}
|
||||
|
||||
fun recordFailure(
|
||||
atSeconds: Long,
|
||||
reason: String?,
|
||||
) {
|
||||
val previous = _state.value
|
||||
_state.value =
|
||||
previous.copy(
|
||||
lastFailureAt = atSeconds,
|
||||
lastFailure = reason,
|
||||
consecutiveFailures = previous.consecutiveFailures + 1,
|
||||
)
|
||||
}
|
||||
|
||||
companion object {
|
||||
const val DOWN_AFTER = 3
|
||||
}
|
||||
}
|
||||
+146
@@ -0,0 +1,146 @@
|
||||
/*
|
||||
* 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.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
|
||||
/**
|
||||
* What the delivery operator learns, as data rather than as documentation.
|
||||
*
|
||||
* §8 of `quartz/plans/2026-09-17-cordn-interop.md` ends with a requirement:
|
||||
* *"This section should be surfaced in the UI if we ship this, not buried. A
|
||||
* Marmot group and a cordn group have materially different metadata exposure
|
||||
* and users cannot infer that from either one looking like a group chat."*
|
||||
*
|
||||
* A comment cannot satisfy that; a screen can, and a screen needs a model. So
|
||||
* the analysis lives here as [GroupExposure], computed from the group's actual
|
||||
* state rather than written out per group — a group joined from a share link
|
||||
* really does expose more than one whose members were added directly, and the
|
||||
* difference is mechanical, not editorial.
|
||||
*
|
||||
* The one thing this deliberately does **not** do is rank the two bindings.
|
||||
* cordn is weaker against the delivery operator and stronger against the
|
||||
* network (§8.5: the coordinator never sees an IP). Which trade is right
|
||||
* depends on who runs the coordinator, which is the user's call and not ours.
|
||||
*/
|
||||
enum class ExposureLevel {
|
||||
/** The operator cannot learn this at all. */
|
||||
NONE,
|
||||
|
||||
/** Learned, but tied to a throwaway key rather than to the account. */
|
||||
PSEUDONYMOUS,
|
||||
|
||||
/** Learned and tied to the real npub. */
|
||||
IDENTIFIED,
|
||||
}
|
||||
|
||||
/** One thing worth telling the user about, with the rule it follows from. */
|
||||
enum class ExposureNote {
|
||||
/**
|
||||
* §8.1 — admission names both ends. `join_request_store` rides the stable
|
||||
* identity and carries the `gid`; `welcome_store` names the target's real
|
||||
* pubkey. For a group joined from a share link the coordinator observes
|
||||
* real-identity membership directly. Structural, not a bug.
|
||||
*/
|
||||
MEMBERSHIP_IS_IDENTIFIED,
|
||||
|
||||
/**
|
||||
* §8.2 — the ephemeral identity is per-session, not per-message, so one
|
||||
* pseudonym posts and fetches across every group on that coordinator. The
|
||||
* set of `gid`s it touches links those groups together and leaks how many
|
||||
* there are. This is the one the user cannot guess and the reason a plain
|
||||
* "messages are pseudonymous" would be a lie by omission.
|
||||
*/
|
||||
GROUPS_LINKED_BY_SESSION,
|
||||
|
||||
/**
|
||||
* §8.3 — one operator holds the complete ordered history of every group it
|
||||
* serves, with timestamps and payload sizes, in one SQLite file. Relays
|
||||
* are redundant and partitioned; a coordinator is neither.
|
||||
*/
|
||||
SINGLE_OPERATOR_HOLDS_HISTORY,
|
||||
|
||||
/**
|
||||
* §8.4 — `kp_publish` rides the stable identity and the coordinator keeps
|
||||
* the signed event, by design, to re-serve. That is a verifiable record
|
||||
* that this account uses cordn, rotation cadence included.
|
||||
*/
|
||||
PUBLICATION_IS_A_SIGNED_RECORD,
|
||||
|
||||
/**
|
||||
* §8.3 — nothing in `spec/03.md` pads the sealed payload, so message sizes
|
||||
* are visible to the operator.
|
||||
*/
|
||||
MESSAGE_SIZES_UNPADDED,
|
||||
|
||||
/**
|
||||
* §8.6 — the transport is NOT pinned to required encryption, so a
|
||||
* coordinator that declines to announce `support_encryption` would receive
|
||||
* plaintext JSON-RPC on public relays. Should be impossible in our client;
|
||||
* if this ever appears, it is a bug, not a disclosure.
|
||||
*/
|
||||
ENCRYPTION_NOT_PINNED,
|
||||
}
|
||||
|
||||
/**
|
||||
* The exposure of one cordn group on one coordinator.
|
||||
*
|
||||
* @param linkedGroupCount how many groups share this coordinator's ephemeral
|
||||
* session (§8.2). One is already a link between that group and the account's
|
||||
* traffic; more is a graph.
|
||||
*/
|
||||
data class GroupExposure(
|
||||
val coordinator: HexKey,
|
||||
val linkedGroupCount: Int,
|
||||
val joinedFromShareLink: Boolean,
|
||||
val publishedKeyPackage: Boolean,
|
||||
val encryptionPinned: Boolean,
|
||||
) {
|
||||
/** §8: double-sealed, and the coordinator is forbidden from parsing. Always. */
|
||||
val content: ExposureLevel = ExposureLevel.NONE
|
||||
|
||||
/** §8.1. Admission names real keys on both ends; there is no other way in. */
|
||||
val membership: ExposureLevel = ExposureLevel.IDENTIFIED
|
||||
|
||||
/** §8.2. The ephemeral identity covers the message path — and only that. */
|
||||
val messaging: ExposureLevel = ExposureLevel.PSEUDONYMOUS
|
||||
|
||||
fun notes(): List<ExposureNote> =
|
||||
buildList {
|
||||
add(ExposureNote.MEMBERSHIP_IS_IDENTIFIED)
|
||||
add(ExposureNote.SINGLE_OPERATOR_HOLDS_HISTORY)
|
||||
add(ExposureNote.MESSAGE_SIZES_UNPADDED)
|
||||
// Only worth saying once there is actually something to link to.
|
||||
if (linkedGroupCount > 1) add(ExposureNote.GROUPS_LINKED_BY_SESSION)
|
||||
if (publishedKeyPackage) add(ExposureNote.PUBLICATION_IS_A_SIGNED_RECORD)
|
||||
if (!encryptionPinned) add(ExposureNote.ENCRYPTION_NOT_PINNED)
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether this group's exposure differs from a Marmot group's in a way the
|
||||
* user should be told before they treat the two the same.
|
||||
*
|
||||
* Always true today, and written as a function anyway: the interesting case
|
||||
* is a self-hosted coordinator, where the answer changes without any of
|
||||
* this code changing.
|
||||
*/
|
||||
fun differsFromMarmot(): Boolean = true
|
||||
}
|
||||
+516
@@ -0,0 +1,516 @@
|
||||
/*
|
||||
* 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.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.appGroupRef.CordnGroupRef
|
||||
import com.vitorpamplona.quartz.cordn.groups.CordnCredential
|
||||
import com.vitorpamplona.quartz.cordn.groups.CordnGroupPolicy
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ICoordinator
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.KeyPackagePublication
|
||||
import com.vitorpamplona.quartz.cordn.spec01GroupMetadata.CordnGroupMetadata
|
||||
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnApplicationMessage
|
||||
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnEnvelope
|
||||
import com.vitorpamplona.quartz.cordn.spec02Envelopes.ReceivedMessage
|
||||
import com.vitorpamplona.quartz.cordn.spec03Payloads.SealedPayload
|
||||
import com.vitorpamplona.quartz.cordn.sync.CordnGroupSync
|
||||
import com.vitorpamplona.quartz.cordn.sync.Ingestion
|
||||
import com.vitorpamplona.quartz.mls.codec.TlsReader
|
||||
import com.vitorpamplona.quartz.mls.framing.ContentType
|
||||
import com.vitorpamplona.quartz.mls.framing.MlsMessage
|
||||
import com.vitorpamplona.quartz.mls.framing.WireFormat
|
||||
import com.vitorpamplona.quartz.mls.group.MlsGroup
|
||||
import com.vitorpamplona.quartz.mls.group.MlsGroupState
|
||||
import com.vitorpamplona.quartz.mls.messages.KeyPackageBundle
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import kotlin.io.encoding.Base64
|
||||
import kotlin.io.encoding.ExperimentalEncodingApi
|
||||
|
||||
/**
|
||||
* The cordn groups this account holds on one coordinator.
|
||||
*
|
||||
* The cordn counterpart of `MarmotManager`, and deliberately **not** a reuse of
|
||||
* it. Marmot's manager is keyed on the Nostr group id throughout and its whole
|
||||
* delivery model — publish kinds 443/444/445 to relays, subscribe by `h` tag,
|
||||
* carry publish obligations — has no cordn analogue: here there is one
|
||||
* coordinator, one ordered stream per `gid`, and a cursor. What the two share
|
||||
* is the RFC 9420 engine underneath, which is what Stage 1 of the interop plan
|
||||
* made possible.
|
||||
*
|
||||
* ## One manager per coordinator
|
||||
*
|
||||
* Not per account. A `gid` is scoped to the coordinator that issued it
|
||||
* (`spec/00.md` §4-5), so the same string can name two unrelated groups on two
|
||||
* coordinators, and a cursor from one is meaningless to the other. Keying
|
||||
* groups by `gid` alone is only safe because the coordinator is fixed here.
|
||||
*
|
||||
* It also matches the privacy model: §8.2's ephemeral identity is per session
|
||||
* per coordinator, so the set of groups one manager touches is exactly the set
|
||||
* that coordinator can link together. [exposure] reports that from the real
|
||||
* count rather than from a guess.
|
||||
*/
|
||||
@OptIn(ExperimentalEncodingApi::class)
|
||||
class CordnGroupManager(
|
||||
/** The account these groups belong to, as lowercase hex. */
|
||||
val accountPubKey: HexKey,
|
||||
val config: CoordinatorConfig,
|
||||
private val coordinator: ICoordinator,
|
||||
private val store: CordnGroupStore,
|
||||
val health: CoordinatorHealth = CoordinatorHealth(),
|
||||
/** Seconds. Injected so tests are not at the mercy of the wall clock. */
|
||||
private val clock: () -> Long = { TimeUtils.now() },
|
||||
) {
|
||||
private val groups = mutableMapOf<String, MlsGroup>()
|
||||
private val sync = CordnGroupSync(coordinator)
|
||||
|
||||
private val _gids = MutableStateFlow<Set<String>>(emptySet())
|
||||
|
||||
/** The groups this manager holds, for a UI to observe. */
|
||||
val gids: StateFlow<Set<String>> = _gids.asStateFlow()
|
||||
|
||||
private var publishedKeyPackage = false
|
||||
|
||||
/** What happened to one delivered payload. */
|
||||
sealed interface Delivery {
|
||||
/** An application message that passed every §5 check. */
|
||||
data class Message(
|
||||
val gid: String,
|
||||
val cursor: Long,
|
||||
val received: ReceivedMessage,
|
||||
) : Delivery
|
||||
|
||||
/** A handshake message; the group has advanced to [epoch]. */
|
||||
data class EpochAdvanced(
|
||||
val gid: String,
|
||||
val cursor: Long,
|
||||
val epoch: Long,
|
||||
) : Delivery
|
||||
|
||||
/** Our own message coming back, already accounted for. */
|
||||
data class Echo(
|
||||
val gid: String,
|
||||
val cursor: Long,
|
||||
) : Delivery
|
||||
|
||||
/**
|
||||
* A payload this group could not open.
|
||||
*
|
||||
* Not fatal and not silent: the cursor moves past it (the alternative
|
||||
* is a stream that never advances), but a UI that shows nothing here is
|
||||
* hiding a gap in a conversation.
|
||||
*/
|
||||
data class Undecryptable(
|
||||
val gid: String,
|
||||
val cursor: Long,
|
||||
val reason: String,
|
||||
) : Delivery
|
||||
}
|
||||
|
||||
// ---- membership ------------------------------------------------------
|
||||
|
||||
/** The group behind [gid], if this manager holds it. */
|
||||
fun group(gid: String): MlsGroup? = groups[gid]
|
||||
|
||||
/** Restores every group and cursor the store holds. Call once at startup. */
|
||||
suspend fun restore() {
|
||||
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) }
|
||||
}
|
||||
_gids.value = groups.keys.toSet()
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a group this account administers.
|
||||
*
|
||||
* [gid] is the caller's to choose and the coordinator never interprets it
|
||||
* (`spec/00.md` §4). The reference client uses a random UUID; anything
|
||||
* unique on this coordinator works, and it must not be derived from the MLS
|
||||
* `group_id`, which is secret.
|
||||
*/
|
||||
suspend fun createGroup(
|
||||
gid: String,
|
||||
metadata: CordnGroupMetadata,
|
||||
): MlsGroup {
|
||||
require(gid !in groups) { "already in a group with gid $gid" }
|
||||
val group =
|
||||
MlsGroup.create(
|
||||
identity = CordnCredential.of(accountPubKey).identity,
|
||||
policy = CordnGroupPolicy,
|
||||
initialExtensions = listOf(metadata.toExtension()),
|
||||
// See [gidFrom]: this is what lets a joiner -- ours or theirs --
|
||||
// learn the delivery id from the Welcome alone.
|
||||
groupId = gid.encodeToByteArray(),
|
||||
)
|
||||
groups[gid] = group
|
||||
persist(gid)
|
||||
_gids.value = groups.keys.toSet()
|
||||
return group
|
||||
}
|
||||
|
||||
/**
|
||||
* Adds [targetPubKey] to [gid], taking their KeyPackage from the
|
||||
* coordinator and leaving them a Welcome.
|
||||
*
|
||||
* The order is load-bearing:
|
||||
*
|
||||
* 1. Take and **verify** the publication payload. §9/§10 require the client
|
||||
* to check it even though §8 makes the coordinator check too — a
|
||||
* coordinator that skipped the check could otherwise hand us any
|
||||
* account's name over any account's key material, and we would invite
|
||||
* the wrong person.
|
||||
* 2. Seal the Commit under the **pre-commit** epoch key (`spec/03.md` §5).
|
||||
* The engine hands that key back on [CommitResult.preCommitExporterSecret]
|
||||
* precisely because it is unobtainable afterwards: sealing under the
|
||||
* epoch the Commit *creates* produces a payload only the sender can
|
||||
* open, and every other member silently stops advancing.
|
||||
* 3. Post the Commit, then store the Welcome with the cursor it landed at,
|
||||
* so the joiner starts from the epoch they can actually decrypt rather
|
||||
* than replaying history that predates them.
|
||||
*
|
||||
* ## Wire format
|
||||
*
|
||||
* The Commit goes out **public-framed** — `MlsMessage(PublicMessage)`,
|
||||
* which is what this engine produces. cordn's reference client emits
|
||||
* private-framed handshake messages instead, but accepts either: its
|
||||
* `processMessageBase64` admits wireformat 1 and 2 and hands both to
|
||||
* ts-mls's `processMessage`. Nothing is weakened by the choice, because the
|
||||
* coordinator sees only the outer ChaCha seal either way (§8) — the framing
|
||||
* is visible only to members, who can already tell a Commit from a message.
|
||||
* [ingest] reads both.
|
||||
*/
|
||||
suspend fun invite(
|
||||
gid: String,
|
||||
targetPubKey: HexKey,
|
||||
): InviteResult {
|
||||
val group = requireGroup(gid)
|
||||
|
||||
val taken =
|
||||
call { coordinator.takeKeyPackage(targetPubKey) }
|
||||
?: throw CordnGroupException("the coordinator holds no KeyPackage for $targetPubKey")
|
||||
val verified = KeyPackagePublication.verify(taken.publicationEvent)
|
||||
if (verified.pubKey != targetPubKey) {
|
||||
throw CordnGroupException("the KeyPackage served for $targetPubKey belongs to ${verified.pubKey}")
|
||||
}
|
||||
|
||||
val result = group.addMember(verified.bytes)
|
||||
val welcome = result.welcomeBytes ?: throw CordnGroupException("adding a member produced no Welcome")
|
||||
|
||||
// framedCommitBytes, not commitBytes: the latter is the bare RFC 9420
|
||||
// Commit struct with no MLSMessage around it, which no receiver can
|
||||
// parse. preCommitExporterSecret is the epoch key the Commit leaves.
|
||||
val posted =
|
||||
call { sync.postCommit(gid, SealedPayload.seal(result.framedCommitBytes, result.preCommitExporterSecret)) }
|
||||
val welcomeAt =
|
||||
call {
|
||||
coordinator.storeWelcome(
|
||||
targetPubKey = targetPubKey,
|
||||
keyPackageRef = taken.keyPackageRef,
|
||||
welcomeBase64 = Base64.encode(welcome),
|
||||
after = posted.cursor,
|
||||
)
|
||||
}
|
||||
|
||||
persist(gid)
|
||||
return InviteResult(gid, targetPubKey, posted.cursor, welcomeAt)
|
||||
}
|
||||
|
||||
/**
|
||||
* Joins every group we have been invited to and can open.
|
||||
*
|
||||
* [bundleFor] resolves a `kp_ref` to the KeyPackageBundle we published
|
||||
* under it. A Welcome we have no private half for is left alone rather
|
||||
* than acknowledged, because acknowledging it retires it forever.
|
||||
*
|
||||
* [gidFor] overrides where the delivery id comes from — see [gidFrom] for
|
||||
* why it can be missing and what to pass when it is.
|
||||
*/
|
||||
suspend fun joinPendingWelcomes(
|
||||
bundleFor: (String) -> KeyPackageBundle?,
|
||||
gidFor: (MlsGroup) -> String? = ::gidFrom,
|
||||
): JoinResults {
|
||||
val pending = call { coordinator.takeWelcomes() }
|
||||
val joined = mutableListOf<String>()
|
||||
val skipped = mutableListOf<SkippedWelcome>()
|
||||
|
||||
pending.forEach { welcome ->
|
||||
val bundle = bundleFor(welcome.keyPackageRef)
|
||||
if (bundle == null) {
|
||||
skipped += SkippedWelcome(welcome.keyPackageRef, "no KeyPackage private half for this ref")
|
||||
return@forEach
|
||||
}
|
||||
val group =
|
||||
try {
|
||||
MlsGroup.processWelcome(Base64.decode(welcome.welcomeBase64), bundle, CordnGroupPolicy)
|
||||
} catch (e: Exception) {
|
||||
skipped += SkippedWelcome(welcome.keyPackageRef, e.message ?: "the Welcome did not open")
|
||||
return@forEach
|
||||
}
|
||||
val gid = gidFor(group)
|
||||
if (gid == null) {
|
||||
skipped += SkippedWelcome(welcome.keyPackageRef, UNKNOWN_GID)
|
||||
return@forEach
|
||||
}
|
||||
|
||||
groups[gid] = group
|
||||
// `after` is the inviter saying where this member's history starts.
|
||||
// Without it a joiner replays epochs from before it existed and
|
||||
// every one lands as Undecryptable.
|
||||
welcome.after?.let { sync.restore(gid, sync.inbox(gid).cursor.advancedTo(it)) }
|
||||
persist(gid)
|
||||
joined += gid
|
||||
}
|
||||
|
||||
_gids.value = groups.keys.toSet()
|
||||
return JoinResults(joined, skipped)
|
||||
}
|
||||
|
||||
// ---- messages --------------------------------------------------------
|
||||
|
||||
/** Sends [content] to [gid] as a cordn application message. */
|
||||
suspend fun send(
|
||||
gid: String,
|
||||
content: String,
|
||||
kind: Int = CHAT_KIND,
|
||||
tags: Array<Array<String>> = emptyArray(),
|
||||
): CordnEnvelope {
|
||||
val group = requireGroup(gid)
|
||||
val envelope =
|
||||
CordnEnvelope.build(
|
||||
pubKey = accountPubKey,
|
||||
createdAt = clock(),
|
||||
kind = kind,
|
||||
tags = tags,
|
||||
content = content,
|
||||
)
|
||||
val sealed = CordnApplicationMessage.seal(group, accountPubKey, envelope)
|
||||
call { sync.postMessage(gid, sealed) }
|
||||
persist(gid)
|
||||
return envelope
|
||||
}
|
||||
|
||||
/** Drains history for every group this manager holds. */
|
||||
suspend fun catchUp(onDelivery: (Delivery) -> Unit): Int {
|
||||
if (groups.isEmpty()) return 0
|
||||
return call { sync.catchUp(groups.keys.toList()) { gid, ingestion -> onDelivery(ingest(gid, ingestion)) } }
|
||||
.also { persistAll() }
|
||||
}
|
||||
|
||||
/**
|
||||
* Subscribes to live delivery for every group. Suspends until the
|
||||
* coordinator closes the stream, so give it its own coroutine.
|
||||
*/
|
||||
suspend fun subscribe(
|
||||
timeoutMs: Long,
|
||||
onDelivery: (Delivery) -> Unit,
|
||||
) {
|
||||
if (groups.isEmpty()) return
|
||||
call { sync.subscribe(groups.keys.toList(), timeoutMs) { gid, ingestion -> onDelivery(ingest(gid, ingestion)) } }
|
||||
}
|
||||
|
||||
private fun ingest(
|
||||
gid: String,
|
||||
ingestion: Ingestion,
|
||||
): Delivery {
|
||||
val group = groups[gid] ?: return Delivery.Undecryptable(gid, ingestion.cursor, "no such group")
|
||||
|
||||
val sealed =
|
||||
when (ingestion) {
|
||||
is Ingestion.OwnMessage -> return Delivery.Echo(gid, ingestion.cursor)
|
||||
// Our own Commit, already applied locally when we posted it.
|
||||
is Ingestion.SelfEchoConfirmed -> return Delivery.Echo(gid, ingestion.cursor)
|
||||
// Our own Commit that we posted but had not applied -- the echo
|
||||
// is the instruction to apply it, which is how a client that
|
||||
// died mid-post recovers.
|
||||
is Ingestion.SelfEchoUnapplied -> ingestion.sealedBase64
|
||||
is Ingestion.Process -> ingestion.sealedBase64
|
||||
}
|
||||
|
||||
return try {
|
||||
// Our current epoch key opens both an application message sent at
|
||||
// this epoch and a Commit leaving it, which is the same key by
|
||||
// construction (spec/03.md §5).
|
||||
val opened = SealedPayload.open(sealed, SealedPayload.applicationKey(group))
|
||||
when (MlsMessage.decodeTls(TlsReader(opened)).wireFormat) {
|
||||
// Public-framed handshake: what this engine emits, and what
|
||||
// Marmot uses throughout.
|
||||
WireFormat.PUBLIC_MESSAGE -> {
|
||||
group.processFramedCommit(opened)
|
||||
Delivery.EpochAdvanced(gid, ingestion.cursor, group.epoch)
|
||||
}
|
||||
// Private-framed: what cordn's reference client emits, for both
|
||||
// application messages and handshake traffic.
|
||||
else -> {
|
||||
val decrypted = group.decrypt(opened)
|
||||
when (decrypted.contentType) {
|
||||
ContentType.APPLICATION ->
|
||||
Delivery.Message(gid, ingestion.cursor, CordnApplicationMessage.open(decrypted))
|
||||
// decrypt() applies a Commit, so the epoch has moved already.
|
||||
else -> Delivery.EpochAdvanced(gid, ingestion.cursor, group.epoch)
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
// The cursor has already advanced past it. Reporting rather than
|
||||
// throwing is what keeps one bad payload from stalling every other
|
||||
// group in the same page.
|
||||
Delivery.Undecryptable(gid, ingestion.cursor, e.message ?: "could not open payload")
|
||||
}
|
||||
}
|
||||
|
||||
// ---- sharing and disclosure -----------------------------------------
|
||||
|
||||
/** The `cordn1…` ref to share this group, pointing at this coordinator. */
|
||||
fun shareRef(gid: String): CordnGroupRef {
|
||||
requireGroup(gid)
|
||||
return CordnGroupRef(
|
||||
gid = gid,
|
||||
coordinatorPubKey = config.pubKey,
|
||||
relays = config.relays.map { it.url },
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* What this coordinator learns about [gid]. See [GroupExposure].
|
||||
*
|
||||
* @param joinedFromShareLink whether admission went through
|
||||
* `join_request_store`, which names the asker's real npub (§8.1).
|
||||
*/
|
||||
fun exposure(
|
||||
gid: String,
|
||||
joinedFromShareLink: Boolean = false,
|
||||
): GroupExposure {
|
||||
requireGroup(gid)
|
||||
return GroupExposure(
|
||||
coordinator = config.pubKey,
|
||||
linkedGroupCount = groups.size,
|
||||
joinedFromShareLink = joinedFromShareLink,
|
||||
publishedKeyPackage = publishedKeyPackage,
|
||||
// :contextvm's CvmGiftWrap defaults to REQUIRED and CoordinatorClient
|
||||
// does not undo it, so this is a fact about our transport, not a
|
||||
// setting. It is reported rather than assumed so that a future
|
||||
// configurable transport cannot quietly make it false (§8.6).
|
||||
encryptionPinned = true,
|
||||
)
|
||||
}
|
||||
|
||||
/** Publishes a KeyPackage so others can add this account. §8.4 applies. */
|
||||
suspend fun publishKeyPackage(
|
||||
keyPackageRef: String,
|
||||
keyPackageBase64: String,
|
||||
) = call { coordinator.publishKeyPackage(keyPackageRef, keyPackageBase64) }
|
||||
.also { publishedKeyPackage = true }
|
||||
|
||||
// ---- plumbing --------------------------------------------------------
|
||||
|
||||
private fun requireGroup(gid: String) = groups[gid] ?: throw CordnGroupException("not a member of $gid")
|
||||
|
||||
private suspend fun persist(gid: String) {
|
||||
val group = groups[gid] ?: return
|
||||
store.saveGroup(gid, group.saveState().encodeTls())
|
||||
sync.cursors()[gid]?.let { store.saveCursor(gid, it) }
|
||||
}
|
||||
|
||||
private suspend fun persistAll() {
|
||||
groups.keys.forEach { persist(it) }
|
||||
}
|
||||
|
||||
/** Runs a coordinator call, recording the outcome against [health]. */
|
||||
private suspend fun <T> call(block: suspend () -> T): T =
|
||||
try {
|
||||
block().also { health.recordSuccess(clock()) }
|
||||
} catch (e: Exception) {
|
||||
health.recordFailure(clock(), e.message)
|
||||
throw e
|
||||
}
|
||||
|
||||
companion object {
|
||||
/** `spec/02.md` §6: a cordn chat message is a NIP-C7 kind 9. */
|
||||
const val CHAT_KIND = 9
|
||||
|
||||
/** Why a Welcome was left in the inbox. Actionable, so it is a constant. */
|
||||
const val UNKNOWN_GID = "cannot tell which delivery group this Welcome is for"
|
||||
|
||||
/**
|
||||
* The delivery `gid` a Welcome implies, or null when it implies none.
|
||||
*
|
||||
* **A Welcome does not carry the `gid`.** It carries the MLS
|
||||
* `group_id`, and `spec/03.md` §2 is explicit that the two are
|
||||
* decoupled — the coordinator's delivery id is not the MLS group id and
|
||||
* must not be assumed to be. So in principle a joiner has no way to
|
||||
* learn where to fetch from, and needs the `cordn1…` ref out of band.
|
||||
*
|
||||
* In practice the reference client closes that gap by convention:
|
||||
* cordn-web sets `group_id = utf8(gid)` when it creates a group and
|
||||
* reads the `gid` straight back out of the group context on join
|
||||
* (`chatGroupLifecycle.svelte.ts`). Verified against staircase's
|
||||
* fixtures, whose `group_id` decodes to exactly their published `gid`.
|
||||
*
|
||||
* This follows that convention — [createGroup] writes it and this reads
|
||||
* it — while treating it as what it is: an observation about one
|
||||
* implementation, not a guarantee. A group id that is not valid UTF-8
|
||||
* cannot be a `gid`, and a conformant peer is free to produce one, so
|
||||
* the answer is null and the caller is told rather than handed a
|
||||
* mojibake key that would quietly fetch nothing forever. Pass `gidFor`
|
||||
* to supply the id from a share ref instead.
|
||||
*/
|
||||
fun gidFrom(group: MlsGroup): String? =
|
||||
try {
|
||||
group.groupId.decodeToString(throwOnInvalidSequence = true).takeIf { it.isNotEmpty() }
|
||||
} catch (e: CharacterCodingException) {
|
||||
null
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Something went wrong that is this manager's to explain, not MLS's. */
|
||||
class CordnGroupException(
|
||||
message: String,
|
||||
) : IllegalStateException(message)
|
||||
|
||||
/** What [CordnGroupManager.joinPendingWelcomes] did with each pending Welcome. */
|
||||
data class JoinResults(
|
||||
val joined: List<String>,
|
||||
/**
|
||||
* Welcomes left pending, with why. Not an error list: a Welcome for a
|
||||
* KeyPackage this device never held belongs to another device of the same
|
||||
* account, and draining it here would destroy it.
|
||||
*/
|
||||
val skipped: List<SkippedWelcome>,
|
||||
)
|
||||
|
||||
/** One Welcome that stayed in the inbox, and the reason. */
|
||||
data class SkippedWelcome(
|
||||
val keyPackageRef: String,
|
||||
val reason: String,
|
||||
)
|
||||
|
||||
/** The result of adding a member: where the Commit and the Welcome landed. */
|
||||
data class InviteResult(
|
||||
val gid: String,
|
||||
val invited: HexKey,
|
||||
val commitCursor: Long,
|
||||
val welcomeAt: Long,
|
||||
)
|
||||
+93
@@ -0,0 +1,93 @@
|
||||
/*
|
||||
* 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.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.sync.GroupCursor
|
||||
|
||||
/**
|
||||
* Local storage for cordn group state, keyed by the delivery `gid`.
|
||||
*
|
||||
* **Implementations MUST encrypt at rest.** The blobs are
|
||||
* `MlsGroupState.encodeTls()` output: private keys and epoch secrets.
|
||||
*
|
||||
* Separate from Marmot's `MlsGroupStateStore` rather than shared, and the key
|
||||
* is the reason. That interface is keyed on the Nostr group id — the MIP-01 `h`
|
||||
* tag — which cordn has no equivalent of; cordn's key is the coordinator's
|
||||
* delivery `gid`, which is opaque, is not the MLS `group_id`, and is not even
|
||||
* unique across coordinators. The two look alike and mean different things, so
|
||||
* one interface serving both would be an invitation to hand the wrong id to the
|
||||
* wrong store and silently find no group.
|
||||
*
|
||||
* The cursor lives here too because it is worthless apart from the state it
|
||||
* belongs to: a cursor restored against a group at a different epoch replays
|
||||
* messages the group can no longer decrypt.
|
||||
*/
|
||||
interface CordnGroupStore {
|
||||
suspend fun saveGroup(
|
||||
gid: String,
|
||||
state: ByteArray,
|
||||
)
|
||||
|
||||
suspend fun loadGroup(gid: String): ByteArray?
|
||||
|
||||
suspend fun deleteGroup(gid: String)
|
||||
|
||||
/** Every `gid` with saved state, for restoring memberships at startup. */
|
||||
suspend fun listGroups(): List<String>
|
||||
|
||||
suspend fun saveCursor(
|
||||
gid: String,
|
||||
cursor: GroupCursor,
|
||||
)
|
||||
|
||||
suspend fun loadCursor(gid: String): GroupCursor?
|
||||
}
|
||||
|
||||
/** 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>()
|
||||
|
||||
override suspend fun saveGroup(
|
||||
gid: String,
|
||||
state: ByteArray,
|
||||
) {
|
||||
groups[gid] = state
|
||||
}
|
||||
|
||||
override suspend fun loadGroup(gid: String): ByteArray? = groups[gid]
|
||||
|
||||
override suspend fun deleteGroup(gid: String) {
|
||||
groups.remove(gid)
|
||||
cursors.remove(gid)
|
||||
}
|
||||
|
||||
override suspend fun listGroups(): List<String> = groups.keys.toList()
|
||||
|
||||
override suspend fun saveCursor(
|
||||
gid: String,
|
||||
cursor: GroupCursor,
|
||||
) {
|
||||
cursors[gid] = cursor
|
||||
}
|
||||
|
||||
override suspend fun loadCursor(gid: String): GroupCursor? = cursors[gid]
|
||||
}
|
||||
+377
@@ -0,0 +1,377 @@
|
||||
/*
|
||||
* 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.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.groups.CordnCredential
|
||||
import com.vitorpamplona.quartz.cordn.groups.CordnGroupPolicy
|
||||
import com.vitorpamplona.quartz.cordn.spec01GroupMetadata.CordnGroupMetadata
|
||||
import com.vitorpamplona.quartz.mls.group.MlsGroup
|
||||
import com.vitorpamplona.quartz.mls.messages.KeyPackageBundle
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.EventTemplate
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlin.io.encoding.Base64
|
||||
import kotlin.io.encoding.ExperimentalEncodingApi
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* [CordnGroupManager] over a full two-party lifecycle.
|
||||
*
|
||||
* Alice creates a group, takes Bob's published KeyPackage off the coordinator,
|
||||
* adds him, and sends a message; Bob joins from the Welcome the coordinator
|
||||
* held for him, catches up, and reads it. Both managers run against one
|
||||
* [FakeCoordinator], so the only thing joining them is the wire — which is the
|
||||
* point: a cursor or epoch mistake shows up as Bob failing to read, not as an
|
||||
* assertion about internals.
|
||||
*/
|
||||
@OptIn(ExperimentalEncodingApi::class)
|
||||
class CordnGroupManagerTest {
|
||||
private val gid = "6d1f0f6a-2a3e-4f2c-9a1d-7c6b5e4d3a21"
|
||||
private val coordinatorKey = "cc".repeat(32)
|
||||
|
||||
private val aliceSigner = NostrSignerSync(KeyPair())
|
||||
private val bobSigner = NostrSignerSync(KeyPair())
|
||||
private val alice: HexKey get() = aliceSigner.pubKey
|
||||
private val bob: HexKey get() = bobSigner.pubKey
|
||||
|
||||
private val config =
|
||||
CoordinatorConfig(
|
||||
pubKey = coordinatorKey,
|
||||
relays = listOf(RelayUrlNormalizer.normalizeOrNull("wss://relay.example.com")!!),
|
||||
origin = CoordinatorConfig.Origin.DEFAULT,
|
||||
)
|
||||
|
||||
private fun manager(
|
||||
account: HexKey,
|
||||
coordinator: FakeCoordinator,
|
||||
store: CordnGroupStore = InMemoryCordnGroupStore(),
|
||||
) = CordnGroupManager(
|
||||
accountPubKey = account,
|
||||
config = config,
|
||||
coordinator = coordinator,
|
||||
store = store,
|
||||
clock = { 1_757_000_000L },
|
||||
)
|
||||
|
||||
/**
|
||||
* Bob's KeyPackage plus the signed `kp_publish` request event that binds it
|
||||
* to him — `spec/00.md` §7's "signed publication payload", which for cordn
|
||||
* *is* the ContextVM request event. Built here rather than faked because
|
||||
* `invite` verifies it, and a fake would verify nothing.
|
||||
*/
|
||||
private fun bobsPublication(): Pair<KeyPackageBundle, FakeCoordinator.StoredKeyPackage> {
|
||||
val scratch = MlsGroup.create(CordnCredential.of(bob).identity, policy = CordnGroupPolicy)
|
||||
val bundle = scratch.createKeyPackage(CordnCredential.of(bob).identity, ByteArray(0))
|
||||
val bytes = bundle.keyPackage.toTlsBytes()
|
||||
val base64 = Base64.encode(bytes)
|
||||
val ref =
|
||||
bundle.keyPackage
|
||||
.toTlsBytes()
|
||||
.take(16)
|
||||
.joinToString("") { "%02x".format(it) }
|
||||
|
||||
val content =
|
||||
"""{"jsonrpc":"2.0","id":1,"method":"tools/call",""" +
|
||||
""""params":{"name":"kp_publish","arguments":{"kp_ref":"$ref","kp_64":"$base64"}}}"""
|
||||
val event: Event =
|
||||
bobSigner.sign(
|
||||
EventTemplate(
|
||||
createdAt = 1_757_000_000L,
|
||||
kind = 25910,
|
||||
tags = arrayOf(arrayOf("p", coordinatorKey)),
|
||||
content = content,
|
||||
),
|
||||
)
|
||||
|
||||
return bundle to FakeCoordinator.StoredKeyPackage(bob, ref, base64, event)
|
||||
}
|
||||
|
||||
/** As [bobsPublication], for a third member. */
|
||||
private fun carolsPublication(): Pair<KeyPackageBundle, FakeCoordinator.StoredKeyPackage> {
|
||||
val carolSigner = NostrSignerSync(KeyPair())
|
||||
val carol = carolSigner.pubKey
|
||||
val scratch = MlsGroup.create(CordnCredential.of(carol).identity, policy = CordnGroupPolicy)
|
||||
val bundle = scratch.createKeyPackage(CordnCredential.of(carol).identity, ByteArray(0))
|
||||
val base64 = Base64.encode(bundle.keyPackage.toTlsBytes())
|
||||
val ref = "ca".repeat(16)
|
||||
val content =
|
||||
"""{"jsonrpc":"2.0","id":1,"method":"tools/call",""" +
|
||||
""""params":{"name":"kp_publish","arguments":{"kp_ref":"$ref","kp_64":"$base64"}}}"""
|
||||
val event: Event =
|
||||
carolSigner.sign(
|
||||
EventTemplate(
|
||||
createdAt = 1_757_000_000L,
|
||||
kind = 25910,
|
||||
tags = arrayOf(arrayOf("p", coordinatorKey)),
|
||||
content = content,
|
||||
),
|
||||
)
|
||||
return bundle to FakeCoordinator.StoredKeyPackage(carol, ref, base64, event)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `alice creates, invites, sends, and bob joins and reads`() =
|
||||
runTest {
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
val (bobBundle, stored) = bobsPublication()
|
||||
coordinator.seedKeyPackage(stored)
|
||||
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Stage 4", adminPubkeys = listOf(alice)))
|
||||
assertEquals(setOf(gid), aliceManager.gids.value)
|
||||
|
||||
val invite = aliceManager.invite(gid, bob)
|
||||
assertEquals(bob, invite.invited)
|
||||
assertEquals(1L, aliceManager.group(gid)!!.epoch, "adding a member advances the epoch")
|
||||
|
||||
aliceManager.send(gid, "hello bob")
|
||||
|
||||
// Bob's side, from nothing but what the coordinator holds.
|
||||
val bobCoordinator = FakeCoordinator(callerPubKey = bob)
|
||||
bobCoordinator.welcomes.putAll(coordinator.welcomes)
|
||||
bobCoordinator.streams.putAll(coordinator.streams)
|
||||
val bobManager = manager(bob, bobCoordinator)
|
||||
|
||||
val results = bobManager.joinPendingWelcomes({ ref -> bobBundle.takeIf { ref == stored.keyPackageRef } })
|
||||
assertEquals(listOf(gid), results.joined, "Bob recovers the gid from the group id")
|
||||
assertTrue(results.skipped.isEmpty())
|
||||
|
||||
val delivered = mutableListOf<CordnGroupManager.Delivery>()
|
||||
bobManager.catchUp { delivered += it }
|
||||
|
||||
val messages = delivered.filterIsInstance<CordnGroupManager.Delivery.Message>()
|
||||
assertEquals(1, messages.size, "got ${delivered.map { it::class.simpleName }}")
|
||||
assertEquals("hello bob", messages[0].received.envelope.content)
|
||||
assertEquals(alice, messages[0].received.sender, "the sender is MLS-authenticated, not claimed")
|
||||
assertEquals(CordnGroupManager.CHAT_KIND, messages[0].received.envelope.kind)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `the welcome cursor keeps bob from replaying epochs he cannot read`() =
|
||||
runTest {
|
||||
// Without `after`, catch-up starts at 0 and hands Bob the Commit
|
||||
// that created him plus everything before it -- all Undecryptable,
|
||||
// and all indistinguishable from real loss.
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
val (bobBundle, stored) = bobsPublication()
|
||||
coordinator.seedKeyPackage(stored)
|
||||
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Stage 4"))
|
||||
aliceManager.send(gid, "before bob existed")
|
||||
aliceManager.invite(gid, bob)
|
||||
aliceManager.send(gid, "after bob joined")
|
||||
|
||||
val bobCoordinator = FakeCoordinator(callerPubKey = bob)
|
||||
bobCoordinator.welcomes.putAll(coordinator.welcomes)
|
||||
bobCoordinator.streams.putAll(coordinator.streams)
|
||||
val bobManager = manager(bob, bobCoordinator)
|
||||
bobManager.joinPendingWelcomes({ bobBundle.takeIf { _ -> true } })
|
||||
|
||||
val delivered = mutableListOf<CordnGroupManager.Delivery>()
|
||||
bobManager.catchUp { delivered += it }
|
||||
|
||||
val messages = delivered.filterIsInstance<CordnGroupManager.Delivery.Message>()
|
||||
assertEquals(listOf("after bob joined"), messages.map { it.received.envelope.content })
|
||||
assertTrue(
|
||||
delivered.none { it is CordnGroupManager.Delivery.Undecryptable },
|
||||
"nothing from before Bob's epoch should have been offered: ${delivered.map { it::class.simpleName }}",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `an existing member applies a later commit, sealed under the pre-commit epoch`() =
|
||||
runTest {
|
||||
// `spec/03.md` §5: a Commit is sealed with the exporter of the
|
||||
// epoch it LEAVES, because that is the only key its recipients
|
||||
// have. Sealing under the epoch it creates produces a payload the
|
||||
// sender can read and nobody else can -- and it cannot be caught
|
||||
// by the join path, where the Welcome's cursor skips the Commit
|
||||
// that created the joiner. It needs a member already in the group
|
||||
// when a later Commit lands, which is what this is.
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
val (bobBundle, stored) = bobsPublication()
|
||||
coordinator.seedKeyPackage(stored)
|
||||
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Three"))
|
||||
aliceManager.invite(gid, bob)
|
||||
|
||||
val bobCoordinator = FakeCoordinator(callerPubKey = bob)
|
||||
bobCoordinator.welcomes.putAll(coordinator.welcomes)
|
||||
bobCoordinator.streams.putAll(coordinator.streams)
|
||||
val bobManager = manager(bob, bobCoordinator)
|
||||
bobManager.joinPendingWelcomes({ bobBundle })
|
||||
bobManager.catchUp { }
|
||||
assertEquals(1L, bobManager.group(gid)!!.epoch)
|
||||
|
||||
// Carol arrives. Bob was already here, so the Commit is his to apply.
|
||||
val (_, carolStored) = carolsPublication()
|
||||
coordinator.seedKeyPackage(carolStored)
|
||||
aliceManager.invite(gid, carolStored.pubKey)
|
||||
aliceManager.send(gid, "carol is in")
|
||||
|
||||
bobCoordinator.streams.clear()
|
||||
bobCoordinator.streams.putAll(coordinator.streams)
|
||||
val delivered = mutableListOf<CordnGroupManager.Delivery>()
|
||||
bobManager.catchUp { delivered += it }
|
||||
|
||||
assertTrue(
|
||||
delivered.none { it is CordnGroupManager.Delivery.Undecryptable },
|
||||
"Bob could not open a Commit addressed to his own epoch: " +
|
||||
delivered.filterIsInstance<CordnGroupManager.Delivery.Undecryptable>().map { it.reason },
|
||||
)
|
||||
assertEquals(2L, bobManager.group(gid)!!.epoch, "the Commit must have advanced Bob's epoch")
|
||||
assertEquals(
|
||||
listOf("carol is in"),
|
||||
delivered.filterIsInstance<CordnGroupManager.Delivery.Message>().map { it.received.envelope.content },
|
||||
"and the message sent at the new epoch must still open",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a group survives a restart`() =
|
||||
runTest {
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val store = InMemoryCordnGroupStore()
|
||||
val first = manager(alice, coordinator, store)
|
||||
first.createGroup(gid, CordnGroupMetadata(name = "Persisted"))
|
||||
first.send(gid, "one")
|
||||
|
||||
val restarted = manager(alice, coordinator, store)
|
||||
restarted.restore()
|
||||
|
||||
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.
|
||||
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 }}",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `inviting someone the coordinator has no KeyPackage for fails loudly`() =
|
||||
runTest {
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Stage 4"))
|
||||
|
||||
val failure = assertFailsWith<CordnGroupException> { aliceManager.invite(gid, bob) }
|
||||
assertTrue(failure.message!!.contains("holds no KeyPackage"))
|
||||
assertEquals(0L, aliceManager.group(gid)!!.epoch, "a failed invite must not advance the group")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a KeyPackage published under someone else's name is refused`() =
|
||||
runTest {
|
||||
// §9/§10: the coordinator checks this too, and we check it anyway.
|
||||
// A coordinator that skipped the check could otherwise hand us any
|
||||
// account's name over any account's key material.
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Stage 4"))
|
||||
|
||||
val (_, stored) = bobsPublication()
|
||||
val impostor = "ee".repeat(32)
|
||||
coordinator.seedKeyPackage(
|
||||
FakeCoordinator.StoredKeyPackage(impostor, stored.keyPackageRef, stored.base64, stored.publicationEvent),
|
||||
)
|
||||
|
||||
assertFailsWith<Exception> { aliceManager.invite(gid, impostor) }
|
||||
assertEquals(0L, aliceManager.group(gid)!!.epoch)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `health tracks the coordinator without polling it`() =
|
||||
runTest {
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Stage 4"))
|
||||
|
||||
assertTrue(aliceManager.health.state.value.isUnknown, "no call yet is not the same as down")
|
||||
|
||||
aliceManager.send(gid, "works")
|
||||
assertEquals(0, aliceManager.health.state.value.consecutiveFailures)
|
||||
|
||||
coordinator.failNext = CoordinatorHealth.DOWN_AFTER
|
||||
repeat(CoordinatorHealth.DOWN_AFTER) {
|
||||
runCatching { aliceManager.send(gid, "fails") }
|
||||
}
|
||||
assertTrue(aliceManager.health.state.value.isDown)
|
||||
|
||||
aliceManager.send(gid, "works again")
|
||||
assertEquals(0, aliceManager.health.state.value.consecutiveFailures, "a success clears the streak")
|
||||
assertTrue(!aliceManager.health.state.value.isDown)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `exposure reports the real linked-group count`() =
|
||||
runTest {
|
||||
// Section 8.2 is the one a user cannot guess: one ephemeral session
|
||||
// per coordinator means every group on it is linked to the others.
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "One"))
|
||||
|
||||
val alone = aliceManager.exposure(gid)
|
||||
assertEquals(1, alone.linkedGroupCount)
|
||||
assertTrue(ExposureNote.GROUPS_LINKED_BY_SESSION !in alone.notes())
|
||||
assertEquals(ExposureLevel.NONE, alone.content)
|
||||
assertEquals(ExposureLevel.IDENTIFIED, alone.membership)
|
||||
assertTrue(ExposureNote.ENCRYPTION_NOT_PINNED !in alone.notes(), "our transport pins REQUIRED (§8.6)")
|
||||
|
||||
aliceManager.createGroup("second-gid", CordnGroupMetadata(name = "Two"))
|
||||
val linked = aliceManager.exposure(gid)
|
||||
assertEquals(2, linked.linkedGroupCount)
|
||||
assertTrue(ExposureNote.GROUPS_LINKED_BY_SESSION in linked.notes())
|
||||
|
||||
assertTrue(ExposureNote.PUBLICATION_IS_A_SIGNED_RECORD !in linked.notes())
|
||||
aliceManager.publishKeyPackage("ref", "a2s=")
|
||||
assertTrue(ExposureNote.PUBLICATION_IS_A_SIGNED_RECORD in aliceManager.exposure(gid).notes(), "§8.4")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a share ref round-trips through this coordinator`() =
|
||||
runTest {
|
||||
val coordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, coordinator)
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Shared"))
|
||||
|
||||
val ref = aliceManager.shareRef(gid)
|
||||
assertEquals(gid, ref.gid)
|
||||
assertEquals(coordinatorKey, ref.coordinatorPubKey)
|
||||
assertEquals(CoordinatorConfig.from(ref)?.pubKey, coordinatorKey)
|
||||
assertNull(aliceManager.group("nope"))
|
||||
}
|
||||
}
|
||||
+177
@@ -0,0 +1,177 @@
|
||||
/*
|
||||
* 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.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.AvailableKeyPackage
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ConsumedJoinRequestRef
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ConsumedWelcomeRef
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.GroupMessage
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ICoordinator
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.JoinRequest
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.PendingWelcome
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.PostedMessage
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.PublishedKeyPackage
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.TakenKeyPackage
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
|
||||
/**
|
||||
* A coordinator that behaves, in memory.
|
||||
*
|
||||
* It models the two things `spec/00.md` makes the coordinator responsible for
|
||||
* and that the manager's correctness depends on: a **monotonic per-group
|
||||
* cursor** and **per-account Welcome inboxes**. Everything else it stores
|
||||
* verbatim, which is also what a real one does — §8: it cannot parse payloads
|
||||
* and does not try.
|
||||
*
|
||||
* Wire shapes are not this fake's job. Those are pinned against cordn's own zod
|
||||
* schemas in `quartz`'s `CoordinatorContractVectorTest`; here the point is the
|
||||
* manager's behaviour on top of them.
|
||||
*/
|
||||
class FakeCoordinator(
|
||||
/** The account the manager under test is calling as, for Welcome routing. */
|
||||
private val callerPubKey: HexKey,
|
||||
) : ICoordinator {
|
||||
class StoredKeyPackage(
|
||||
val pubKey: HexKey,
|
||||
val keyPackageRef: String,
|
||||
val base64: String,
|
||||
val publicationEvent: Event,
|
||||
val lastResort: Boolean = false,
|
||||
)
|
||||
|
||||
val keyPackages = mutableMapOf<String, StoredKeyPackage>()
|
||||
val welcomes = mutableMapOf<HexKey, MutableList<PendingWelcome>>()
|
||||
val streams = mutableMapOf<String, MutableList<GroupMessage>>()
|
||||
|
||||
/** Every method name called, in order. */
|
||||
val calls = mutableListOf<String>()
|
||||
|
||||
private var cursor = 0L
|
||||
private var clock = 1_757_000_000L
|
||||
|
||||
/** Fails the next [failNext] calls, to exercise the health surface. */
|
||||
var failNext = 0
|
||||
|
||||
/** Pre-loads a KeyPackage as if its owner had published it. */
|
||||
fun seedKeyPackage(stored: StoredKeyPackage) {
|
||||
keyPackages[stored.keyPackageRef] = stored
|
||||
}
|
||||
|
||||
/** Everything posted to [gid], oldest first. */
|
||||
fun posted(gid: String): List<String> = streams[gid].orEmpty().map { it.sealedBase64 }
|
||||
|
||||
private fun record(method: String) {
|
||||
calls += method
|
||||
if (failNext > 0) {
|
||||
failNext--
|
||||
throw IllegalStateException("$method failed")
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun publishKeyPackage(
|
||||
keyPackageRef: String,
|
||||
keyPackageBase64: String,
|
||||
): PublishedKeyPackage {
|
||||
record("kp_publish")
|
||||
return PublishedKeyPackage(keyPackageRef, false, clock++)
|
||||
}
|
||||
|
||||
override suspend fun removeKeyPackages(keyPackageRefs: List<String>): List<String> {
|
||||
record("kp_remove")
|
||||
return keyPackageRefs.filter { keyPackages.remove(it) != null }
|
||||
}
|
||||
|
||||
override suspend fun listKeyPackages(): List<AvailableKeyPackage> {
|
||||
record("kp_list")
|
||||
return keyPackages.values.map { AvailableKeyPackage(it.pubKey, it.keyPackageRef, it.lastResort, clock) }
|
||||
}
|
||||
|
||||
override suspend fun takeKeyPackage(id: String): TakenKeyPackage? {
|
||||
record("kp_take")
|
||||
// `id` accepts a ref or an account hex, like the real one.
|
||||
val stored = keyPackages[id] ?: keyPackages.values.firstOrNull { it.pubKey == id } ?: return null
|
||||
return TakenKeyPackage(stored.pubKey, stored.keyPackageRef, stored.lastResort, clock++, stored.publicationEvent)
|
||||
}
|
||||
|
||||
override suspend fun storeWelcome(
|
||||
targetPubKey: HexKey,
|
||||
keyPackageRef: String,
|
||||
welcomeBase64: String,
|
||||
after: Long?,
|
||||
): Long {
|
||||
record("welcome_store")
|
||||
val at = clock++
|
||||
welcomes.getOrPut(targetPubKey) { mutableListOf() } += PendingWelcome(keyPackageRef, welcomeBase64, at, after)
|
||||
return at
|
||||
}
|
||||
|
||||
override suspend fun takeWelcomes(consumed: List<ConsumedWelcomeRef>): List<PendingWelcome> {
|
||||
record("welcome_take")
|
||||
val mine = welcomes[callerPubKey].orEmpty()
|
||||
consumed.forEach { ack -> welcomes[callerPubKey]?.removeAll { it.keyPackageRef == ack.keyPackageRef && it.at == ack.at } }
|
||||
return mine.toList()
|
||||
}
|
||||
|
||||
override suspend fun storeJoinRequest(
|
||||
gid: String,
|
||||
keyPackageRef: String,
|
||||
): Long {
|
||||
record("join_request_store")
|
||||
return clock++
|
||||
}
|
||||
|
||||
override suspend fun takeJoinRequests(
|
||||
gids: List<String>,
|
||||
consumed: List<ConsumedJoinRequestRef>,
|
||||
): List<JoinRequest> {
|
||||
record("join_request_take_many")
|
||||
return emptyList()
|
||||
}
|
||||
|
||||
override suspend fun postMessage(
|
||||
gid: String,
|
||||
sealedBase64: String,
|
||||
): PostedMessage {
|
||||
record("msg_post")
|
||||
// Monotonic across the whole coordinator, which is stricter than
|
||||
// spec/00.md §4 requires (per group) and so a safe stand-in.
|
||||
val assigned = ++cursor
|
||||
streams.getOrPut(gid) { mutableListOf() } += GroupMessage(gid, assigned, sealedBase64, clock++)
|
||||
return PostedMessage(gid, assigned, clock)
|
||||
}
|
||||
|
||||
override suspend fun fetchMessages(cursors: Map<String, Long?>): List<GroupMessage> {
|
||||
record("msg_fetch_many")
|
||||
return cursors.entries
|
||||
.flatMap { (gid, after) -> streams[gid].orEmpty().filter { it.cursor > (after ?: 0L) } }
|
||||
.sortedBy { it.cursor }
|
||||
}
|
||||
|
||||
override suspend fun subscribeMessages(
|
||||
cursors: Map<String, Long?>,
|
||||
timeoutMs: Long,
|
||||
onMessage: (GroupMessage) -> Unit,
|
||||
) {
|
||||
record("msg_sub_many")
|
||||
fetchMessages(cursors).forEach(onMessage)
|
||||
}
|
||||
}
|
||||
@@ -1022,7 +1022,62 @@ Still open in Stage 3:
|
||||
non-goal (§4.6).
|
||||
- **Tier B**, live against `ghcr.io/cordn-msg/cordn:latest`.
|
||||
|
||||
### Stage 4 — App integration
|
||||
### Stage 4 — App integration — HEADLESS LAYER LANDED, UI OPEN
|
||||
|
||||
Landed in `commons/…/cordn/` (2026-09-19):
|
||||
|
||||
| What | Where |
|
||||
| ---- | ----- |
|
||||
| Coordinator identity, relays, provenance | `CoordinatorConfig` (+ `from(CordnGroupRef)`) |
|
||||
| Per-coordinator health, recorded not polled | `CoordinatorHealth` |
|
||||
| §8 as a model a UI renders | `CordnExposure` (`GroupExposure`, `ExposureNote`) |
|
||||
| Group lifecycle keyed on `gid` | `CordnGroupManager` |
|
||||
| Encrypted-at-rest state + cursors | `CordnGroupStore` (+ an in-memory one for tests) |
|
||||
| The 11 tools as a contract | `ICoordinator` in `quartz`, implemented by `CoordinatorClient` |
|
||||
| `amy cordn ref encode/decode`, `amy cordn exposure` | `cli/…/CordnCommands` |
|
||||
|
||||
`CordnGroupManager` is the cordn counterpart of `MarmotManager` and shares no code with it, as
|
||||
Stage 1 predicted: Marmot's is keyed on the Nostr group id 147 times and its delivery model —
|
||||
kinds 443/444/445, `h`-tag subscriptions, publish obligations — has no cordn analogue. What the
|
||||
two share is the engine underneath. Its store is separate from `MlsGroupStateStore` for the same
|
||||
reason: the keys look alike (`nostrGroupId` vs `gid`) and mean different things, so one interface
|
||||
serving both would invite handing the wrong id to the wrong store.
|
||||
|
||||
Four findings, each of which was a bug until the test that found it:
|
||||
|
||||
1. **`CommitResult.commitBytes` is the bare RFC 9420 `Commit` struct**, not an `MLSMessage`.
|
||||
Posting it produces a payload no receiver can parse. `framedCommitBytes` is the wire form.
|
||||
2. **The pre-commit epoch rule (`spec/03.md` §5) is invisible to a two-party test.** Sealing a
|
||||
Commit under the epoch it *creates* still passes create→invite→join→read, because the
|
||||
joiner's Welcome cursor skips the only Commit in the stream. It takes a third member — one
|
||||
already in the group when a later Commit lands — to catch it. `CommitResult` already hands
|
||||
back `preCommitExporterSecret` precisely because the key is unobtainable afterwards.
|
||||
3. **A Welcome does not carry the `gid`.** `spec/03.md` §2 decouples the delivery id from the MLS
|
||||
`group_id`, so in principle a joiner cannot know where to fetch from. The reference client
|
||||
closes the gap by convention — `group_id = utf8(gid)` — which staircase's fixtures confirm
|
||||
(their `group_id` decodes to exactly their published `gid`). We follow the convention and
|
||||
treat it as an observation, not a guarantee: a non-UTF-8 group id yields a named skip rather
|
||||
than a mojibake key that fetches nothing forever. `MlsGroup.create` gained an optional
|
||||
`groupId` so our groups carry it too.
|
||||
4. **Public-framed handshake messages interoperate.** Our engine frames Commits as
|
||||
`MLSMessage(PublicMessage)`; cordn's client emits private framing. Checked rather than
|
||||
assumed: `packages/cli/src/utils/mlsMessages.ts:68` accepts wireformat 1 *and* 2 and hands
|
||||
both to ts-mls's `processMessage`. Nothing is weakened — the coordinator sees only the outer
|
||||
seal either way — so the manager emits public framing and reads both.
|
||||
|
||||
Still open:
|
||||
|
||||
- **The UI.** `GroupExposure` exists so that the §8 requirement ("surfaced, not buried") has
|
||||
something to render, but nothing renders it yet. That is the remaining Stage 4 work, along
|
||||
with cordn chatroom/feed models beside `model/marmotGroups/` and a ViewModel.
|
||||
- **Coordinator-driving `amy` verbs** (`publish`, `invite`, `send`, `sync`). Deliberately not
|
||||
shipped: with Tier B blocked (§7) there is nothing to exercise them against, and unexercised
|
||||
coordinator verbs are a guess with a command-line interface. The logic they would call is in
|
||||
`CordnGroupManager` and is covered against an in-memory coordinator.
|
||||
- Key-package rotation and a published-KeyPackage lifecycle, which cordn has no event kind for
|
||||
(§4.2) and which therefore lives entirely in coordinator calls.
|
||||
|
||||
The original scope, for reference:
|
||||
|
||||
- `commons/` state holders and ViewModels for cordn groups, alongside the Marmot ones.
|
||||
- Coordinator configuration and health surface (their client tracks per-coordinator health).
|
||||
|
||||
+14
-14
@@ -64,7 +64,7 @@ class CoordinatorClient(
|
||||
private val mcp: CvmMcpClient,
|
||||
/** Applied to every call. Coordinator work is storage, not computation. */
|
||||
private val timeoutMs: Long = CvmTransport.DEFAULT_TIMEOUT_MS,
|
||||
) {
|
||||
) : ICoordinator {
|
||||
/** Performs the MCP handshake. Optional, but it is where CEP-35 tags ride. */
|
||||
suspend fun initialize() = mcp.initialize()
|
||||
|
||||
@@ -78,7 +78,7 @@ class CoordinatorClient(
|
||||
* Nothing extra is signed here, and nothing else can be: the transport owns
|
||||
* the event.
|
||||
*/
|
||||
suspend fun publishKeyPackage(
|
||||
override suspend fun publishKeyPackage(
|
||||
keyPackageRef: String,
|
||||
keyPackageBase64: String,
|
||||
): PublishedKeyPackage =
|
||||
@@ -91,7 +91,7 @@ class CoordinatorClient(
|
||||
}
|
||||
|
||||
/** Withdraws published KeyPackages. Authorized by us being their owner (§12). */
|
||||
suspend fun removeKeyPackages(keyPackageRefs: List<String>): List<String> {
|
||||
override suspend fun removeKeyPackages(keyPackageRefs: List<String>): List<String> {
|
||||
require(keyPackageRefs.isNotEmpty()) { "kp_remove needs at least one ref" }
|
||||
val args =
|
||||
buildJsonObject {
|
||||
@@ -110,7 +110,7 @@ class CoordinatorClient(
|
||||
* A last-resort KeyPackage can back several Welcomes, which is why a record
|
||||
* is identified by `(kp_ref, at)` rather than by `kp_ref` alone.
|
||||
*/
|
||||
suspend fun takeWelcomes(consumed: List<ConsumedWelcomeRef> = emptyList()): List<PendingWelcome> {
|
||||
override suspend fun takeWelcomes(consumed: List<ConsumedWelcomeRef>): List<PendingWelcome> {
|
||||
val args =
|
||||
buildJsonObject {
|
||||
if (consumed.isNotEmpty()) {
|
||||
@@ -145,7 +145,7 @@ class CoordinatorClient(
|
||||
* Rides the stable identity because the request is "npub X wants into group
|
||||
* G" — there is no version of this call that does not name the asker.
|
||||
*/
|
||||
suspend fun storeJoinRequest(
|
||||
override suspend fun storeJoinRequest(
|
||||
gid: String,
|
||||
keyPackageRef: String,
|
||||
): Long =
|
||||
@@ -160,7 +160,7 @@ class CoordinatorClient(
|
||||
// ---- ephemeral identity ----------------------------------------------
|
||||
|
||||
/** Every KeyPackage this coordinator holds. */
|
||||
suspend fun listKeyPackages(): List<AvailableKeyPackage> =
|
||||
override suspend fun listKeyPackages(): List<AvailableKeyPackage> =
|
||||
call(CoordinatorMethod.KP_LIST, buildJsonObject {}).list(CoordinatorFields.KEY_PACKAGES) {
|
||||
AvailableKeyPackage(
|
||||
pubKey = it.str(CoordinatorFields.PK),
|
||||
@@ -177,7 +177,7 @@ class CoordinatorClient(
|
||||
* on its publication event. Null means the coordinator holds nothing
|
||||
* matching.
|
||||
*/
|
||||
suspend fun takeKeyPackage(id: String): TakenKeyPackage? {
|
||||
override suspend fun takeKeyPackage(id: String): TakenKeyPackage? {
|
||||
val result = call(CoordinatorMethod.KP_TAKE, buildJsonObject { put(CoordinatorFields.ID, id) })
|
||||
val entry = result[CoordinatorFields.KEY_PACKAGE]?.takeIf { it is JsonObject }?.jsonObject ?: return null
|
||||
val eventJson =
|
||||
@@ -198,11 +198,11 @@ class CoordinatorClient(
|
||||
* [after] tells the joiner which cursor to start their history from, so
|
||||
* they do not replay epochs they cannot decrypt.
|
||||
*/
|
||||
suspend fun storeWelcome(
|
||||
override suspend fun storeWelcome(
|
||||
targetPubKey: HexKey,
|
||||
keyPackageRef: String,
|
||||
welcomeBase64: String,
|
||||
after: Long? = null,
|
||||
after: Long?,
|
||||
): Long =
|
||||
call(
|
||||
CoordinatorMethod.WELCOME_STORE,
|
||||
@@ -215,9 +215,9 @@ class CoordinatorClient(
|
||||
).num(CoordinatorFields.AT)
|
||||
|
||||
/** Drains join requests for the groups we administer, acknowledging handled ones. */
|
||||
suspend fun takeJoinRequests(
|
||||
override suspend fun takeJoinRequests(
|
||||
gids: List<String>,
|
||||
consumed: List<ConsumedJoinRequestRef> = emptyList(),
|
||||
consumed: List<ConsumedJoinRequestRef>,
|
||||
): List<JoinRequest> {
|
||||
require(gids.isNotEmpty()) { "join_request_take_many needs at least one group" }
|
||||
val args =
|
||||
@@ -254,7 +254,7 @@ class CoordinatorClient(
|
||||
}
|
||||
|
||||
/** Posts one sealed payload to a group's stream. */
|
||||
suspend fun postMessage(
|
||||
override suspend fun postMessage(
|
||||
gid: String,
|
||||
sealedBase64: String,
|
||||
): PostedMessage =
|
||||
@@ -279,7 +279,7 @@ class CoordinatorClient(
|
||||
* catching up loops until a page comes back empty — see
|
||||
* [com.vitorpamplona.quartz.cordn.sync.GroupSync].
|
||||
*/
|
||||
suspend fun fetchMessages(cursors: Map<String, Long?>): List<GroupMessage> {
|
||||
override suspend fun fetchMessages(cursors: Map<String, Long?>): List<GroupMessage> {
|
||||
require(cursors.isNotEmpty()) { "msg_fetch_many needs at least one group" }
|
||||
val args = buildJsonObject { put(CoordinatorFields.GROUPS, groupsArray(cursors)) }
|
||||
return call(CoordinatorMethod.MSG_FETCH_MANY, args).list(CoordinatorFields.MESSAGES, ::groupMessage)
|
||||
@@ -294,7 +294,7 @@ class CoordinatorClient(
|
||||
* does not complete the request, so the returned list is the whole run's
|
||||
* traffic, not a partial view.
|
||||
*/
|
||||
suspend fun subscribeMessages(
|
||||
override suspend fun subscribeMessages(
|
||||
cursors: Map<String, Long?>,
|
||||
timeoutMs: Long,
|
||||
onMessage: (GroupMessage) -> Unit,
|
||||
|
||||
+81
@@ -0,0 +1,81 @@
|
||||
/*
|
||||
* 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.spec00Coordinator
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
|
||||
/**
|
||||
* The eleven coordinator tools, as a contract.
|
||||
*
|
||||
* [CoordinatorClient] is the only real implementation and this interface adds
|
||||
* nothing to it — no defaults, no behaviour. It exists so the layers above
|
||||
* (`CordnGroupSync`, and the group manager in `commons`) can be exercised
|
||||
* without standing up a relay, a transport and a server, which is otherwise the
|
||||
* price of testing a cursor loop.
|
||||
*
|
||||
* What it deliberately does **not** expose is identity. Which key signs which
|
||||
* call is fixed by [CoordinatorMethod] per `spec/00.md` §8, and a substitute
|
||||
* implementation cannot widen that, because there is no parameter to widen.
|
||||
*/
|
||||
interface ICoordinator {
|
||||
suspend fun publishKeyPackage(
|
||||
keyPackageRef: String,
|
||||
keyPackageBase64: String,
|
||||
): PublishedKeyPackage
|
||||
|
||||
suspend fun removeKeyPackages(keyPackageRefs: List<String>): List<String>
|
||||
|
||||
suspend fun listKeyPackages(): List<AvailableKeyPackage>
|
||||
|
||||
suspend fun takeKeyPackage(id: String): TakenKeyPackage?
|
||||
|
||||
suspend fun storeWelcome(
|
||||
targetPubKey: HexKey,
|
||||
keyPackageRef: String,
|
||||
welcomeBase64: String,
|
||||
after: Long? = null,
|
||||
): Long
|
||||
|
||||
suspend fun takeWelcomes(consumed: List<ConsumedWelcomeRef> = emptyList()): List<PendingWelcome>
|
||||
|
||||
suspend fun storeJoinRequest(
|
||||
gid: String,
|
||||
keyPackageRef: String,
|
||||
): Long
|
||||
|
||||
suspend fun takeJoinRequests(
|
||||
gids: List<String>,
|
||||
consumed: List<ConsumedJoinRequestRef> = emptyList(),
|
||||
): List<JoinRequest>
|
||||
|
||||
suspend fun postMessage(
|
||||
gid: String,
|
||||
sealedBase64: String,
|
||||
): PostedMessage
|
||||
|
||||
suspend fun fetchMessages(cursors: Map<String, Long?>): List<GroupMessage>
|
||||
|
||||
suspend fun subscribeMessages(
|
||||
cursors: Map<String, Long?>,
|
||||
timeoutMs: Long,
|
||||
onMessage: (GroupMessage) -> Unit,
|
||||
)
|
||||
}
|
||||
@@ -20,8 +20,8 @@
|
||||
*/
|
||||
package com.vitorpamplona.quartz.cordn.sync
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.CoordinatorClient
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.GroupMessage
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ICoordinator
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.PostedMessage
|
||||
|
||||
/**
|
||||
@@ -43,7 +43,7 @@ import com.vitorpamplona.quartz.cordn.spec00Coordinator.PostedMessage
|
||||
* either time.
|
||||
*/
|
||||
class CordnGroupSync(
|
||||
private val client: CoordinatorClient,
|
||||
private val client: ICoordinator,
|
||||
/** Inbox per `gid`. The caller owns them, because they outlive a sync run. */
|
||||
private val inboxes: MutableMap<String, GroupInbox> = mutableMapOf(),
|
||||
) {
|
||||
|
||||
@@ -3275,6 +3275,20 @@ class MlsGroup private constructor(
|
||||
* [com.vitorpamplona.quartz.marmot.groups.MarmotCapabilities.currentProfileRequired].
|
||||
*/
|
||||
requiredCapabilities: Extension? = policy.defaultRequiredCapabilities,
|
||||
/**
|
||||
* The RFC 9420 `group_id` for the new group. Null generates a
|
||||
* random 32-byte one, which is the right default: §8.1 only
|
||||
* requires it to be unique, and a random id tells a receiver
|
||||
* nothing.
|
||||
*
|
||||
* A binding may need to choose it. cordn's reference client sets
|
||||
* `group_id = utf8(gid)` so that a joiner can recover the delivery
|
||||
* id from a Welcome, which otherwise carries no way to learn it —
|
||||
* see `CordnGroupManager`. Anything a binding puts here is visible
|
||||
* to whoever handles the ciphertext, so it must carry nothing the
|
||||
* group would not publish.
|
||||
*/
|
||||
groupId: ByteArray? = null,
|
||||
): MlsGroup {
|
||||
val sigKp =
|
||||
signingKey?.let { key ->
|
||||
@@ -3283,7 +3297,7 @@ class MlsGroup private constructor(
|
||||
} ?: Ed25519.generateKeyPair()
|
||||
|
||||
val encKp = X25519.generateKeyPair()
|
||||
val groupId = MlsCryptoProvider.randomBytes(32)
|
||||
val newGroupId = groupId ?: MlsCryptoProvider.randomBytes(32)
|
||||
|
||||
val leafNode =
|
||||
buildLeafNode(
|
||||
@@ -3307,7 +3321,7 @@ class MlsGroup private constructor(
|
||||
val baseExtensions = listOfNotNull(requiredCapabilities)
|
||||
val groupContext =
|
||||
GroupContext(
|
||||
groupId = groupId,
|
||||
groupId = newGroupId,
|
||||
epoch = 0,
|
||||
treeHash = treeHash,
|
||||
confirmedTranscriptHash = ByteArray(0),
|
||||
|
||||
Reference in New Issue
Block a user