feat(marmot): implement agent-text-stream over QUIC (0x8006), receive role

MDK requires component `0x8006` in every group it creates, and its
`required_member_roles` mask names MLS leaf capabilities each member must
advertise. A client that carries neither is refused at the Add — which is
why `wn groups create "Interop-03" <us>` failed outright, taking five
interop scenarios with it. The `0xf2d1`/`0xf2d2`/`0xf2d4` extensions MDK
advertises are exactly those role capabilities.

Implements the component and the record layer it gates:

  - `AgentTextStreamQuicPolicyV1` — the 12-byte component state, decoded
    strictly (a short, long, or out-of-range payload is rejected, never
    defaulted: these bytes sit in signed group state and a guessed role
    mask admits a member the group refuses).
  - `AgentTextStreamRecordV1` — the wire record, with the QUIC varint
    length prefixes the Marmot binary profile uses. Unknown record types
    decode fine on purpose; a newer advisory record must not tear down an
    otherwise valid preview stream.
  - `AgentTextStreamCrypto` — HKDF-Expand-only key and nonce derivation
    over the full key context, `nonce_base XOR uint96_be(seq)`, and the
    record AAD. `seq` is in both the nonce and the AAD, so a replayed or
    reordered record fails to open rather than being noticed afterwards.
  - `AgentTextStreamTranscriptV1` — the rolling hash the final kind-9
    chat publishes, so a receiver can tell that it saw exactly the stream
    the publisher sent.
  - `AgentTextStreamStart` / `AgentTextStreamFinal` — the kind-1200 anchor
    tags and the kind-9 closing tags.

We advertise the RECEIVE role only, and the group state validator refuses
a Welcome whose policy requires a role we do not advertise. Publishing
would need durable per-stream sequence state to avoid reusing an AEAD
nonce across a restart, and we have none — advertising `send` without it
would be a claim we cannot keep.

Also fixes three places that read only the legacy `0xF2EE` extension and
therefore did nothing at all on a current-profile group:

  - The Welcome's `nostr_group_id`, which is the `h` tag every kind-445
    event carries. Without it a joiner cannot subscribe, so MDK's welcome
    decrypted and was then discarded — "GroupContext is missing the
    NostrGroupData extension" — for a routing id that was present the
    whole time in the `0x8004` component.
  - The admin gate on GroupContextExtensions changes, which read
    `adminsConfigured` as false for every current-profile group and so
    skipped the check instead of failing closed.
  - Disappearing-message expiration, which silently never applied.

And `syncMetadataTo`, which left every current-profile group with a blank
name, no admins, no relays and no avatar in the UI.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
Claude
2026-09-09 04:14:43 +00:00
parent cfd44c2108
commit 4a89d29f6c
18 changed files with 1629 additions and 47 deletions
@@ -798,10 +798,12 @@ class Context(
val kp = keyPackageRelaysOf(pubKey)
val nip65 = relaysOf(pubKey)
if (dm == null && kp == null && nip65 == null) return null
val dmInbox = dm?.relays().orEmpty()
return RecipientRelayFetcher.Lists(
dmInbox = dm?.relays().orEmpty(),
dmInbox = dmInbox,
keyPackage = kp?.relays().orEmpty(),
nip65 = nip65,
dmInboxWithheld = dmInbox.isEmpty() && dm?.allRelays().orEmpty().isNotEmpty(),
)
}
@@ -127,9 +127,17 @@ object GroupAddMemberCommand {
// bootstrapped Amethyst accounts listen on these)
// Our own outbox is added as belt-and-braces so we
// can re-ingest the welcome ourselves too.
//
// The default set is a bootstrap for someone who has
// advertised NOTHING, not a fallback for someone whose
// advertised inbox we declined to use. When they named
// a local-network relay we refuse to reach, sending
// their invite to a public default set instead is a
// different action than they asked for — so we keep it
// on the group's own relays and say so.
buildSet {
addAll(recipient.dmInboxOrFallback())
if (isEmpty()) {
if (isEmpty() && !recipient.dmInboxWithheld) {
addAll(DefaultDMRelayList)
}
addAll(ctx.outboxRelays())
@@ -154,6 +162,7 @@ object GroupAddMemberCommand {
"commit_accepted_by" to commitAck.filterValues { it.accepted }.keys.map { it.url },
"welcome_accepted_by" to welcomeAck.filterValues { it.accepted }.keys.map { it.url },
"welcome_targets" to welcomeTargets.map { it.url },
"welcome_inbox_withheld" to recipient.dmInboxWithheld,
"key_package_relays" to kpRelays.map { it.url },
),
)
@@ -32,6 +32,7 @@ import com.vitorpamplona.quartz.marmot.WelcomeDelivery
import com.vitorpamplona.quartz.marmot.WelcomeResult
import com.vitorpamplona.quartz.marmot.appComponents.CurrentProfileGroupFactory
import com.vitorpamplona.quartz.marmot.appComponents.GroupProfileV1
import com.vitorpamplona.quartz.marmot.appComponents.MarmotGroupState
import com.vitorpamplona.quartz.marmot.appComponents.MessageRetentionV1
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageBundleStore
import com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageEvent
@@ -1011,6 +1012,13 @@ class MarmotManager(
return MarmotGroupData.fromExtensions(group.extensions)
}
/**
* The current profile's GroupContext components for a group, or null when
* the group is unknown. A legacy group returns a state whose component
* fields are all null — read [groupMetadata] for those.
*/
fun groupState(nostrGroupId: HexKey): MarmotGroupState? = groupManager.getGroup(nostrGroupId)?.currentGroupState()
/**
* Sync MIP-01 metadata and member info from the MLS group into a [MarmotGroupChatroom].
* Call after joining a group, processing a commit, or restoring from storage.
@@ -1019,27 +1027,44 @@ class MarmotManager(
nostrGroupId: HexKey,
chatroom: MarmotGroupChatroom,
) {
val metadata = groupMetadata(nostrGroupId)
if (metadata != null) {
if (metadata.name.isNotEmpty()) {
chatroom.displayName.value = metadata.name
}
if (metadata.description.isNotEmpty()) {
chatroom.description.value = metadata.description
}
chatroom.adminPubkeys.value = metadata.adminPubkeys
chatroom.relays.value = metadata.relays
chatroom.image.value =
if (metadata.hasImage()) {
// Read the current profile's components first, then the legacy
// 0xF2EE extension. Reading only the legacy one left every
// current-profile group with a blank name, no admins, no relays and no
// avatar in the UI — the group worked, it just looked empty.
val state = groupState(nostrGroupId)
val legacy = groupMetadata(nostrGroupId)
val name = state?.profile?.name?.takeIf { it.isNotEmpty() } ?: legacy?.name
if (!name.isNullOrEmpty()) chatroom.displayName.value = name
val description = state?.profile?.description?.takeIf { it.isNotEmpty() } ?: legacy?.description
if (!description.isNullOrEmpty()) chatroom.description.value = description
val admins = state?.adminPolicy?.adminHexKeys ?: legacy?.adminPubkeys
if (admins != null) chatroom.adminPubkeys.value = admins
val relays = state?.routing?.relays ?: legacy?.relays
if (relays != null) chatroom.relays.value = relays
val image = state?.image
chatroom.image.value =
when {
image?.imageHash != null ->
MarmotGroupImage(
hash = metadata.imageHash!!,
key = metadata.imageKey!!,
nonce = metadata.imageNonce!!,
hash = image.imageHash!!.toHexKey(),
key = image.imageKey!!,
nonce = image.imageNonce!!,
)
} else {
null
}
}
legacy?.hasImage() == true ->
MarmotGroupImage(
hash = legacy.imageHash!!,
key = legacy.imageKey!!,
nonce = legacy.imageNonce!!,
)
else -> null
}
val previousCount = chatroom.members.value.size
val members = memberPubkeys(nostrGroupId)
chatroom.members.value = members
@@ -210,9 +210,21 @@ class MarmotOutboundProcessor(
nostrGroupId: HexKey,
createdAt: Long,
): Long? {
val extensions = groupManager.getGroup(nostrGroupId)?.extensions ?: return null
val marmotData = MarmotGroupData.fromExtensions(extensions) ?: return null
val secs = marmotData.disappearingMessageSecs ?: return null
return createdAt + secs.toLong()
val group = groupManager.getGroup(nostrGroupId) ?: return null
// The current profile carries retention in the 0x8005 component; the
// legacy profile in the monolithic 0xF2EE extension. Reading only the
// legacy one silently dropped the expiration tag on every
// current-profile group, so disappearing messages simply did not
// disappear.
val secs =
group
.currentGroupState()
.retention
?.takeIf { it.isEnabled }
?.disappearingMessageSecs
?.toLong()
?: MarmotGroupData.fromExtensions(group.extensions)?.disappearingMessageSecs?.toLong()
?: return null
return createdAt + secs
}
}
@@ -53,6 +53,17 @@ object RecipientRelayFetcher {
val keyPackage: List<NormalizedRelayUrl>,
/** Latest kind:10002 the user published (null if none seen). */
val nip65: AdvertisedRelayListEvent?,
/**
* True when the user DID advertise a kind:10050 inbox but every entry
* was dropped by the local-network filter, leaving [dmInbox] empty.
*
* "Advertised nothing" and "advertised only relays we refuse to reach"
* are different facts, and the caller has to be able to tell them
* apart. Treating the second as the first sends a gift wrap addressed
* to this user to a default relay set they never chose — the opposite
* of what the filter is for.
*/
val dmInboxWithheld: Boolean = false,
) {
/** Read-marker relays from kind:10002. Mirrors `User.inboxRelays()`. */
fun nip65Read(): List<NormalizedRelayUrl> = nip65?.readRelaysNorm().orEmpty()
@@ -125,10 +136,12 @@ object RecipientRelayFetcher {
}
}
val dmInbox = dm?.relays().orEmpty()
return Lists(
dmInbox = dm?.relays().orEmpty(),
dmInbox = dmInbox,
keyPackage = kp?.relays().orEmpty(),
nip65 = nip65,
dmInboxWithheld = dmInbox.isEmpty() && dm?.allRelays().orEmpty().isNotEmpty(),
)
}
}
@@ -21,6 +21,7 @@
package com.vitorpamplona.quartz.marmot.appComponents
import com.vitorpamplona.quartz.marmot.appComponents.accountIdentityProof.AccountIdentityProofV2
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamQuicPolicyV1
import com.vitorpamplona.quartz.marmot.mip01Groups.MlsCiphersuite
import com.vitorpamplona.quartz.marmot.mls.components.AppDataDictionary
import com.vitorpamplona.quartz.marmot.mls.components.ComponentData
@@ -62,6 +63,14 @@ object CurrentProfileGroupFactory {
* invisible, and one we list but do not implement is a lie that surfaces
* later as a group we cannot actually participate in. Add an id here only
* when the component is implemented.
*
* `0x8006` (agent-text-stream over QUIC) is listed for the RECEIVE role
* only, which is what [MlsGroup.currentProfileLeafCapabilities] advertises:
* we decode the group's policy, derive the per-stream record keys, open
* records and fold the transcript. We do not advertise the `send`
* (`0xF2D2`) or `fanout` (`0xF2D4`) capabilities, because publishing needs
* durable per-stream sequence state to avoid reusing an AEAD nonce across
* a restart, and we have none.
*/
val SUPPORTED_COMPONENTS: List<Int> =
listOf(
@@ -71,6 +80,7 @@ object CurrentProfileGroupFactory {
AppComponentIds.ADMIN_POLICY_V1,
AppComponentIds.NOSTR_ROUTING_V1,
AppComponentIds.MESSAGE_RETENTION_V1,
AppComponentIds.AGENT_TEXT_STREAM_QUIC_V1,
AppComponentIds.ACCOUNT_IDENTITY_PROOF_V2,
AppComponentIds.GROUP_ENCRYPTED_MEDIA_V2,
AppComponentIds.GROUP_LIFECYCLE_V1,
@@ -166,6 +176,7 @@ object CurrentProfileGroupFactory {
profile: GroupProfileV1? = null,
additionalAdmins: List<ByteArray> = emptyList(),
retention: MessageRetentionV1? = null,
agentTextStream: AgentTextStreamQuicPolicyV1? = null,
ciphersuite: MlsCiphersuite = MlsCiphersuite.DEFAULT,
): MlsGroup {
val identity = signer.pubKey.hexToByteArray()
@@ -178,6 +189,7 @@ object CurrentProfileGroupFactory {
profile = profile,
retention = retention,
lifecycle = GroupLifecycleV1.ACTIVE,
agentTextStream = agentTextStream,
)
return MlsGroup.create(
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.quartz.marmot.appComponents
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamQuicPolicyV1
import com.vitorpamplona.quartz.marmot.mls.components.AppDataDictionary
import com.vitorpamplona.quartz.marmot.mls.components.ComponentsList
import com.vitorpamplona.quartz.marmot.mls.tree.Extension
@@ -57,6 +58,7 @@ data class MarmotGroupState(
val image: GroupBlossomImageV1?,
val retention: MessageRetentionV1?,
val lifecycle: GroupLifecycleV1?,
val agentTextStream: AgentTextStreamQuicPolicyV1?,
) {
/** True once a disband Commit has been applied. Absorbing and terminal. */
val isDisbanded: Boolean get() = lifecycle == GroupLifecycleV1.DISBANDED
@@ -99,6 +101,10 @@ data class MarmotGroupState(
image = dictionary[GroupBlossomImageV1.COMPONENT_ID]?.let { GroupBlossomImageV1.decode(it) },
retention = dictionary[MessageRetentionV1.COMPONENT_ID]?.let { MessageRetentionV1.decode(it) },
lifecycle = dictionary[GroupLifecycleV1.COMPONENT_ID]?.let { GroupLifecycleV1.decode(it) },
agentTextStream =
dictionary[AgentTextStreamQuicPolicyV1.COMPONENT_ID]?.let {
AgentTextStreamQuicPolicyV1.decode(it)
},
)
fun fromExtensions(extensions: List<Extension>): MarmotGroupState = fromDictionary(AppDataDictionary.fromExtensionsOrEmpty(extensions))
@@ -117,6 +123,7 @@ data class MarmotGroupState(
image: GroupBlossomImageV1? = null,
retention: MessageRetentionV1? = null,
lifecycle: GroupLifecycleV1? = GroupLifecycleV1.ACTIVE,
agentTextStream: AgentTextStreamQuicPolicyV1? = null,
extraRequiredComponents: Collection<Int> = emptyList(),
): AppDataDictionary {
var dictionary = AppDataDictionary.EMPTY
@@ -145,6 +152,10 @@ data class MarmotGroupState(
required.add(GroupLifecycleV1.COMPONENT_ID)
dictionary = dictionary.with(GroupLifecycleV1.COMPONENT_ID, it.encode())
}
agentTextStream?.let {
required.add(AgentTextStreamQuicPolicyV1.COMPONENT_ID)
dictionary = dictionary.with(AgentTextStreamQuicPolicyV1.COMPONENT_ID, it.encode())
}
return dictionary.with(ComponentsList.APP_COMPONENTS_ID, ComponentsList.encode(required))
}
@@ -0,0 +1,180 @@
/*
* 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.marmot.appComponents.agentTextStream
import com.vitorpamplona.quartz.marmot.mls.crypto.MlsCryptoProvider
import com.vitorpamplona.quartz.nip44Encryption.crypto.ChaCha20Poly1305
/**
* The context every agent-text-stream key derivation is bound to.
*
* ```
* len("v1") || "v1" || len(group_id) || group_id || len(stream_id) || stream_id
* || mls_epoch (uint64) || len(sender_id) || sender_id
* || len(start_event_id) || start_event_id
* ```
*
* The whole tuple goes into the HKDF info, which is what separates two streams
* that share one group exporter secret: change the stream, the epoch, the
* sender or the anchoring kind-1200 event and the record key changes with it.
*/
class AgentTextStreamKeyContextV1(
val groupId: ByteArray,
val streamId: ByteArray,
val mlsEpoch: Long,
val senderId: ByteArray,
val startEventId: ByteArray,
) {
fun encode(): ByteArray {
val out = ArrayList<Byte>()
out.addLengthPrefixed(VERSION)
out.addLengthPrefixed(groupId)
out.addLengthPrefixed(streamId)
for (i in 7 downTo 0) out.add((mlsEpoch shr (8 * i)).toByte())
out.addLengthPrefixed(senderId)
out.addLengthPrefixed(startEventId)
return out.toByteArray()
}
companion object {
val VERSION = "v1".encodeToByteArray()
}
}
/**
* Per-stream record AEAD.
*
* The stream secret is the group's `MLS-Exporter("marmot",
* "agent-text-stream-quic", 32)`, so every member of the epoch can derive it;
* per-stream and per-record separation comes entirely from the key context and
* the sequence number:
*
* - key = `HKDF-Expand(secret, len("record key") || "record key" || context, 32)`
* - nonce = `HKDF-Expand(secret, len("record nonce") || "record nonce" || context, 12)`
* XOR the 96-bit big-endian sequence number
* - aad = `version || SHA-256(group_id) || len(stream_id) || stream_id ||
* mls_epoch || len(sender_id) || sender_id || seq || record_type || flags`
*
* HKDF-**Expand** with no Extract: the exporter secret is already a uniformly
* random PRK, so extracting again would only discard entropy.
*/
class AgentTextStreamCrypto(
val streamSecret: ByteArray,
val context: AgentTextStreamKeyContextV1,
) {
init {
require(streamSecret.size == SECRET_LENGTH) { "agent text stream secret must be $SECRET_LENGTH bytes" }
require(context.streamId.size == AgentTextStreamRecordV1.PROFILE_STREAM_ID_LEN) {
"agent text stream id must be ${AgentTextStreamRecordV1.PROFILE_STREAM_ID_LEN} bytes"
}
require(context.startEventId.size == AgentTextStreamRecordV1.START_EVENT_ID_LEN) {
"agent text stream start event id must be ${AgentTextStreamRecordV1.START_EVENT_ID_LEN} bytes"
}
}
fun recordKey(): ByteArray = derive(KEY_LABEL, KEY_LENGTH)
/** `nonce_base` XOR uint96_be(seq); `seq = 0` yields `nonce_base` itself. */
fun recordNonce(seq: Long): ByteArray {
val nonce = derive(NONCE_LABEL, NONCE_LENGTH)
// uint96_be(seq): the top 4 of the 12 bytes stay zero for any uint64.
for (i in 0 until 8) {
nonce[NONCE_LENGTH - 1 - i] = (nonce[NONCE_LENGTH - 1 - i].toInt() xor (seq shr (8 * i)).toInt()).toByte()
}
return nonce
}
fun recordAad(record: AgentTextStreamRecordV1): ByteArray {
val out = ArrayList<Byte>()
out.add(AgentTextStreamRecordV1.VERSION.toByte())
for (b in MlsCryptoProvider.hash(context.groupId)) out.add(b)
out.addLengthPrefixed(record.streamId)
for (i in 7 downTo 0) out.add((context.mlsEpoch shr (8 * i)).toByte())
out.addLengthPrefixed(context.senderId)
for (i in 7 downTo 0) out.add((record.seq shr (8 * i)).toByte())
out.add(record.recordType.toByte())
out.add(record.flags.toByte())
return out.toByteArray()
}
/** Encrypt [record]'s plaintext frame in place, returning the wire record. */
fun seal(
record: AgentTextStreamRecordV1,
maxPlaintextFrameLen: Long = AgentTextStreamQuicPolicyV1.MAX_PLAINTEXT_FRAME_LEN,
): AgentTextStreamRecordV1 {
requireSameStream(record)
require(record.frame.size <= maxPlaintextFrameLen) {
"agent text stream plaintext frame is larger than the group's limit"
}
return record.copyWithFrame(
ChaCha20Poly1305.encrypt(record.frame, recordAad(record), recordNonce(record.seq), recordKey()),
)
}
/** Decrypt a wire record, returning it with its plaintext frame. */
fun open(record: AgentTextStreamRecordV1): AgentTextStreamRecordV1 {
requireSameStream(record)
require(record.frame.size >= AgentTextStreamRecordV1.AEAD_TAG_LEN) {
"encrypted agent text stream frame is shorter than the AEAD tag"
}
return record.copyWithFrame(
ChaCha20Poly1305.decrypt(record.frame, recordAad(record), recordNonce(record.seq), recordKey()),
)
}
fun openOrNull(record: AgentTextStreamRecordV1): AgentTextStreamRecordV1? =
try {
open(record)
} catch (_: IllegalArgumentException) {
null
} catch (_: IllegalStateException) {
null
}
private fun requireSameStream(record: AgentTextStreamRecordV1) {
require(record.streamId.contentEquals(context.streamId)) {
"agent text stream record belongs to a different stream than its key context"
}
}
private fun derive(
label: ByteArray,
length: Int,
): ByteArray {
val info = ArrayList<Byte>()
info.addLengthPrefixed(label)
for (b in context.encode()) info.add(b)
return MlsCryptoProvider.hkdfExpand(streamSecret, info.toByteArray(), length)
}
companion object {
const val SECRET_LENGTH = 32
const val KEY_LENGTH = 32
const val NONCE_LENGTH = 12
/** `MLS-Exporter("marmot", "agent-text-stream-quic", 32)`. */
const val EXPORTER_LABEL = "marmot"
val EXPORTER_CONTEXT = "agent-text-stream-quic".encodeToByteArray()
private val KEY_LABEL = "record key".encodeToByteArray()
private val NONCE_LABEL = "record nonce".encodeToByteArray()
}
}
@@ -0,0 +1,232 @@
/*
* 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.marmot.appComponents.agentTextStream
import com.vitorpamplona.quartz.marmot.appComponents.AppComponentIds
import com.vitorpamplona.quartz.marmot.mls.components.ComponentData
/**
* `marmot.group.agent-text-stream.quic.v1` (component `0x8006`).
*
* A group carrying this component streams an agent's text as it is produced,
* over QUIC, out of band from the kind-445 group timeline; a kind-1200 event
* anchors each stream and a kind-9 chat carries its final transcript. The
* component itself is only the group's POLICY for those streams: which member
* roles are required and allowed, and the record bounds every participant
* enforces.
*
* Twelve bytes, fixed:
*
* ```
* required_member_roles uint8
* allowed_member_roles uint8
* max_plaintext_frame_len uint32
* replay_ttl_secs uint32
* padding_bucket_bytes uint16
* ```
*
* The role masks are the reason this component cannot be treated as opaque
* bytes. `required_member_roles` names MLS leaf capabilities every member MUST
* advertise ([AgentTextStreamRoles.capabilityFor]), so a group requiring
* `receive` refuses to add a leaf that does not advertise `0xF2D1` — which is
* exactly how an implementation that ignores this component gets locked out of
* every group that carries it.
*/
class AgentTextStreamQuicPolicyV1(
val requiredMemberRoles: Int,
val allowedMemberRoles: Int,
val maxPlaintextFrameLen: Long,
val replayTtlSecs: Long,
val paddingBucketBytes: Int,
) {
init {
require(requiredMemberRoles != 0) { "required agent text stream roles cannot be empty" }
require(requiredMemberRoles and AgentTextStreamRoles.MASK.inv() == 0) {
"required agent text stream role mask contains unknown bits"
}
require(allowedMemberRoles and AgentTextStreamRoles.MASK.inv() == 0) {
"allowed agent text stream role mask contains unknown bits"
}
require(requiredMemberRoles and allowedMemberRoles.inv() == 0) {
"required agent text stream roles must be a subset of allowed roles"
}
require(maxPlaintextFrameLen > 0) { "agent text stream plaintext frame limit cannot be zero" }
require(maxPlaintextFrameLen <= MAX_PLAINTEXT_FRAME_LEN) {
"agent text stream plaintext frame limit exceeds the app profile max"
}
require(replayTtlSecs in 0..MAX_REPLAY_TTL_SECS) {
"agent text stream replay ttl exceeds the app profile max"
}
require(paddingBucketBytes in 0..MAX_PADDING_BUCKET_BYTES) {
"agent text stream padding bucket exceeds the app profile max"
}
}
fun requires(role: Int) = requiredMemberRoles and role != 0
fun allows(role: Int) = allowedMemberRoles and role != 0
/** The MLS leaf capabilities this policy demands of every member. */
fun requiredRoleCapabilities(): List<Int> = AgentTextStreamRoles.capabilitiesFor(requiredMemberRoles)
fun encode(): ByteArray {
val out = ByteArray(STATE_LENGTH)
out[0] = requiredMemberRoles.toByte()
out[1] = allowedMemberRoles.toByte()
out[2] = (maxPlaintextFrameLen shr 24).toByte()
out[3] = (maxPlaintextFrameLen shr 16).toByte()
out[4] = (maxPlaintextFrameLen shr 8).toByte()
out[5] = maxPlaintextFrameLen.toByte()
out[6] = (replayTtlSecs shr 24).toByte()
out[7] = (replayTtlSecs shr 16).toByte()
out[8] = (replayTtlSecs shr 8).toByte()
out[9] = replayTtlSecs.toByte()
out[10] = (paddingBucketBytes shr 8).toByte()
out[11] = paddingBucketBytes.toByte()
return out
}
fun toComponentData() = ComponentData(COMPONENT_ID, encode())
override fun equals(other: Any?): Boolean {
if (this === other) return true
if (other !is AgentTextStreamQuicPolicyV1) return false
return requiredMemberRoles == other.requiredMemberRoles &&
allowedMemberRoles == other.allowedMemberRoles &&
maxPlaintextFrameLen == other.maxPlaintextFrameLen &&
replayTtlSecs == other.replayTtlSecs &&
paddingBucketBytes == other.paddingBucketBytes
}
override fun hashCode(): Int {
var result = requiredMemberRoles
result = 31 * result + allowedMemberRoles
result = 31 * result + maxPlaintextFrameLen.hashCode()
result = 31 * result + replayTtlSecs.hashCode()
result = 31 * result + paddingBucketBytes
return result
}
companion object {
const val COMPONENT_ID = AppComponentIds.AGENT_TEXT_STREAM_QUIC_V1
/** Fixed encoding: 1 + 1 + 4 + 4 + 2. */
const val STATE_LENGTH = 12
/**
* App-profile cap on one frame's plaintext. 65519 keeps the ciphertext
* (plaintext + the 16-byte AEAD tag) inside the record's
* `ciphertext<0..2^16-1>` wire field.
*/
const val MAX_PLAINTEXT_FRAME_LEN = 65519L
const val MAX_REPLAY_TTL_SECS = 5L * 60L
const val MAX_PADDING_BUCKET_BYTES = 4096
/**
* The policy a user-to-agent group installs: every member must be able
* to receive, and a member may additionally send.
*/
fun userToAgentDefault() =
AgentTextStreamQuicPolicyV1(
requiredMemberRoles = AgentTextStreamRoles.RECEIVE,
allowedMemberRoles = AgentTextStreamRoles.RECEIVE or AgentTextStreamRoles.SEND,
maxPlaintextFrameLen = 4096,
replayTtlSecs = 0,
paddingBucketBytes = 0,
)
/**
* Decode exactly [STATE_LENGTH] bytes. A short, long or out-of-range
* payload is rejected rather than defaulted: these bytes live in signed
* group state, and a receiver that guessed a role mask would silently
* admit a member the group refuses.
*/
fun decode(bytes: ByteArray): AgentTextStreamQuicPolicyV1 {
require(bytes.size == STATE_LENGTH) {
"agent text stream component state must be $STATE_LENGTH bytes, got ${bytes.size}"
}
fun u32(at: Int): Long =
((bytes[at].toLong() and 0xff) shl 24) or
((bytes[at + 1].toLong() and 0xff) shl 16) or
((bytes[at + 2].toLong() and 0xff) shl 8) or
(bytes[at + 3].toLong() and 0xff)
return AgentTextStreamQuicPolicyV1(
requiredMemberRoles = bytes[0].toInt() and 0xff,
allowedMemberRoles = bytes[1].toInt() and 0xff,
maxPlaintextFrameLen = u32(2),
replayTtlSecs = u32(6),
paddingBucketBytes = ((bytes[10].toInt() and 0xff) shl 8) or (bytes[11].toInt() and 0xff),
)
}
fun decodeOrNull(bytes: ByteArray): AgentTextStreamQuicPolicyV1? =
try {
decode(bytes)
} catch (_: IllegalArgumentException) {
null
}
}
}
/**
* The three agent-text-stream roles, and the MLS leaf capability that backs
* each one.
*
* Each role is a distinct private-use MLS extension type rather than one flag
* on the component, so a member can advertise `receive` without claiming
* `send` or `fanout`. That is what makes `required_member_roles` enforceable
* per role: MLS itself refuses a leaf that does not advertise a required
* extension, so the group's policy is checked by the CGKA rather than by the
* application.
*/
object AgentTextStreamRoles {
const val RECEIVE = 0x01
const val SEND = 0x02
const val FANOUT = 0x04
const val MASK = RECEIVE or SEND or FANOUT
/** MLS leaf capability (extension type) for each role. */
const val RECEIVE_CAPABILITY = 0xF2D1
const val SEND_CAPABILITY = 0xF2D2
const val FANOUT_CAPABILITY = 0xF2D4
fun capabilityFor(role: Int): Int =
when (role) {
RECEIVE -> RECEIVE_CAPABILITY
SEND -> SEND_CAPABILITY
FANOUT -> FANOUT_CAPABILITY
else -> throw IllegalArgumentException("unknown agent text stream role $role")
}
/**
* Capabilities a role mask demands, in role order. Unknown bits are
* ignored here — the component decoder rejects them, so a mask that
* reaches this function has already been validated.
*/
fun capabilitiesFor(mask: Int): List<Int> =
buildList {
if (mask and RECEIVE != 0) add(RECEIVE_CAPABILITY)
if (mask and SEND != 0) add(SEND_CAPABILITY)
if (mask and FANOUT != 0) add(FANOUT_CAPABILITY)
}
}
@@ -0,0 +1,171 @@
/*
* 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.marmot.appComponents.agentTextStream
/**
* One `AgentTextStreamRecordV1` frame, the unit that travels over the QUIC
* stream.
*
* ```
* version uint8 (= 1)
* stream_id opaque<V>
* seq uint64
* record_type uint8
* flags uint8
* frame opaque<V>
* ```
*
* [frame] holds plaintext before [AgentTextStreamCrypto.seal] and ciphertext
* after it; the wire always carries the sealed form. `seq` is authenticated
* through the AEAD nonce and the AAD, so a reordered or replayed record fails
* to open rather than being detected afterwards.
*/
class AgentTextStreamRecordV1(
val streamId: ByteArray,
val seq: Long,
val recordType: Int,
val flags: Int = 0,
val frame: ByteArray,
) {
init {
require(streamId.isNotEmpty()) { "agent text stream id cannot be empty" }
require(streamId.size <= MAX_STREAM_ID_LEN) { "agent text stream id is too long: ${streamId.size}" }
require(seq >= 0) { "agent text stream record seq cannot be negative" }
require(recordType in 0..0xff) { "agent text stream record type is not a uint8" }
require(flags in 0..0xff) { "agent text stream record flags are not a uint8" }
require(frame.size <= MAX_CIPHERTEXT_LEN) { "agent text stream frame is too large: ${frame.size}" }
}
fun copyWithFrame(newFrame: ByteArray) =
AgentTextStreamRecordV1(
streamId = streamId,
seq = seq,
recordType = recordType,
flags = flags,
frame = newFrame,
)
fun encode(): ByteArray {
val out = ArrayList<Byte>(encodedLength())
out.add(VERSION.toByte())
out.addLengthPrefixed(streamId)
for (i in 7 downTo 0) out.add((seq shr (8 * i)).toByte())
out.add(recordType.toByte())
out.add(flags.toByte())
out.addLengthPrefixed(frame)
return out.toByteArray()
}
fun encodedLength(): Int =
1 +
QuicVarInt.encodedLength(streamId.size.toLong()) + streamId.size +
8 + 1 + 1 +
QuicVarInt.encodedLength(frame.size.toLong()) + frame.size
override fun equals(other: Any?): Boolean {
if (this === other) return true
if (other !is AgentTextStreamRecordV1) return false
return seq == other.seq &&
recordType == other.recordType &&
flags == other.flags &&
streamId.contentEquals(other.streamId) &&
frame.contentEquals(other.frame)
}
override fun hashCode(): Int {
var result = streamId.contentHashCode()
result = 31 * result + seq.hashCode()
result = 31 * result + recordType
result = 31 * result + flags
result = 31 * result + frame.contentHashCode()
return result
}
companion object {
const val VERSION = 1
const val TYPE_TEXT_DELTA = 0x01
const val TYPE_PROGRESS_DELTA = 0x02
const val TYPE_STATUS = 0x03
const val TYPE_CHECKPOINT = 0x04
const val TYPE_ABORT = 0x05
const val TYPE_FINAL_NOTICE = 0x06
const val MAX_STREAM_ID_LEN = 64
/** The profile pins the stream id to 32 bytes; the framing allows less. */
const val PROFILE_STREAM_ID_LEN = 32
const val START_EVENT_ID_LEN = 32
/** Wire bound of the record's `ciphertext<0..2^16-1>` field. */
const val MAX_CIPHERTEXT_LEN = 0xffff
const val AEAD_TAG_LEN = 16
fun textDelta(
streamId: ByteArray,
seq: Long,
frame: ByteArray,
) = AgentTextStreamRecordV1(streamId, seq, TYPE_TEXT_DELTA, 0, frame)
/**
* Decode one record. Trailing bytes are an error: a record is framed
* by the QUIC stream, so anything after the frame field means the
* sender and receiver disagree about where this record ends.
*
* An UNKNOWN [recordType] is deliberately not an error. A newer
* advisory record type must not tear down an otherwise valid preview
* stream, so receivers ignore semantics they do not understand.
*/
fun decode(bytes: ByteArray): AgentTextStreamRecordV1 {
require(bytes.isNotEmpty()) { "agent text stream record is truncated while reading version" }
require(bytes[0].toInt() and 0xff == VERSION) {
"unsupported agent text stream record version: ${bytes[0].toInt() and 0xff}"
}
var at = 1
val streamIdLen = QuicVarInt.decode(bytes, at)
at += streamIdLen.length
require(at + streamIdLen.value <= bytes.size) { "agent text stream record is truncated while reading stream_id" }
val streamId = bytes.copyOfRange(at, at + streamIdLen.value.toInt())
at += streamIdLen.value.toInt()
require(at + 8 <= bytes.size) { "agent text stream record is truncated while reading seq" }
var seq = 0L
for (i in 0 until 8) seq = (seq shl 8) or (bytes[at + i].toLong() and 0xff)
at += 8
require(at + 2 <= bytes.size) { "agent text stream record is truncated while reading record_type" }
val recordType = bytes[at].toInt() and 0xff
val flags = bytes[at + 1].toInt() and 0xff
at += 2
val frameLen = QuicVarInt.decode(bytes, at)
at += frameLen.length
require(at + frameLen.value <= bytes.size) { "agent text stream record is truncated while reading frame" }
val frame = bytes.copyOfRange(at, at + frameLen.value.toInt())
at += frameLen.value.toInt()
require(at == bytes.size) { "agent text stream record contains trailing bytes: ${bytes.size - at}" }
return AgentTextStreamRecordV1(streamId, seq, recordType, flags, frame)
}
}
}
@@ -0,0 +1,128 @@
/*
* 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.marmot.appComponents.agentTextStream
import com.vitorpamplona.quartz.nip01Core.core.HexKey
/**
* The kind-1200 event that anchors one agent text stream, read off an inner
* app payload's tags.
*
* A stream never carries its own key material: the anchoring event's id is
* part of [AgentTextStreamKeyContextV1], so a receiver can only derive record
* keys for a stream it has already seen announced inside the group. That is
* also why the anchor is an ordinary in-group app payload rather than
* something the broker hands out — the broker relays ciphertext and learns
* nothing.
*
* ```
* ["stream", <32-byte stream id, hex>]
* ["route", "quic"] (optional; "quic" when absent)
* ["broker", <candidate URL>] (repeatable)
* ```
*
* The matching END of a stream is an ordinary kind-9 chat carrying
* [STREAM_TAG], [STREAM_HASH_TAG] and [STREAM_CHUNKS_TAG]; see
* [AgentTextStreamFinal].
*/
class AgentTextStreamStart(
val streamId: HexKey,
val route: String,
val brokerCandidates: List<String>,
) {
val isQuicRoute: Boolean get() = route == ROUTE_QUIC
companion object {
const val KIND = 1200
const val STREAM_TAG = "stream"
const val ROUTE_TAG = "route"
const val BROKER_TAG = "broker"
const val ROUTE_QUIC = "quic"
fun tags(
streamId: HexKey,
brokerCandidates: List<String>,
route: String = ROUTE_QUIC,
): Array<Array<String>> =
buildList {
add(arrayOf(STREAM_TAG, streamId))
add(arrayOf(ROUTE_TAG, route))
brokerCandidates.forEach { add(arrayOf(BROKER_TAG, it)) }
}.toTypedArray()
/**
* Read the start view from a kind-1200 payload's tags, or null when
* this is not a stream start. A missing `route` means `quic`; a
* missing `stream` tag means the payload is not usable as an anchor at
* all, so it is rejected rather than defaulted.
*/
fun fromTags(
kind: Int,
tags: Array<Array<String>>,
): AgentTextStreamStart? {
if (kind != KIND) return null
val streamId = tags.firstOrNull { it.size >= 2 && it[0] == STREAM_TAG }?.get(1) ?: return null
val route = tags.firstOrNull { it.size >= 2 && it[0] == ROUTE_TAG }?.get(1) ?: ROUTE_QUIC
val brokers = tags.filter { it.size >= 2 && it[0] == BROKER_TAG }.map { it[1] }
return AgentTextStreamStart(streamId, route, brokers)
}
}
}
/**
* The `stream` / `stream-hash` / `stream-chunks` tags a kind-9 chat carries to
* close a stream out.
*
* The final chat is the durable message; the stream was a live preview of the
* text it now contains. A receiver that folded every record compares its own
* [AgentTextStreamTranscriptV1] against these two values: agreement means it
* saw exactly the stream the publisher sent, and disagreement means records
* were dropped, reordered, or injected — even though each one opened.
*/
class AgentTextStreamFinal(
val streamId: HexKey,
val transcriptHash: HexKey,
val chunkCount: Long,
) {
companion object {
const val STREAM_HASH_TAG = "stream-hash"
const val STREAM_CHUNKS_TAG = "stream-chunks"
fun tags(
streamId: HexKey,
transcriptHash: HexKey,
chunkCount: Long,
): Array<Array<String>> =
arrayOf(
arrayOf(AgentTextStreamStart.STREAM_TAG, streamId),
arrayOf(STREAM_HASH_TAG, transcriptHash),
arrayOf(STREAM_CHUNKS_TAG, chunkCount.toString()),
)
fun fromTags(tags: Array<Array<String>>): AgentTextStreamFinal? {
val streamId = tags.firstOrNull { it.size >= 2 && it[0] == AgentTextStreamStart.STREAM_TAG }?.get(1) ?: return null
val hash = tags.firstOrNull { it.size >= 2 && it[0] == STREAM_HASH_TAG }?.get(1) ?: return null
val chunks = tags.firstOrNull { it.size >= 2 && it[0] == STREAM_CHUNKS_TAG }?.get(1)?.toLongOrNull() ?: return null
return AgentTextStreamFinal(streamId, hash, chunks)
}
}
}
@@ -0,0 +1,101 @@
/*
* 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.marmot.appComponents.agentTextStream
import com.vitorpamplona.quartz.marmot.mls.crypto.MlsCryptoProvider
/**
* The rolling hash a stream's final kind-9 chat publishes as `stream-hash`,
* alongside the `stream-chunks` count.
*
* ```
* h0 = SHA-256("marmot agent text stream transcript v1"
* || len(stream_id) || stream_id
* || len(start_event_id) || start_event_id)
* h_n = SHA-256(h_{n-1} || seq || record_type || plaintext_frame)
* ```
*
* Every receiver folds the same records in the same order, so a receiver whose
* hash disagrees with the publisher's saw a different stream — a dropped,
* reordered or injected record — even though each individual record opened.
*/
class AgentTextStreamTranscriptV1 private constructor(
val streamId: ByteArray,
val startEventId: ByteArray,
private var rollingHash: ByteArray,
private var chunks: Long,
) {
val hash: ByteArray get() = rollingHash.copyOf()
val chunkCount: Long get() = chunks
fun append(
seq: Long,
recordType: Int,
plaintextFrame: ByteArray,
) {
val input = ArrayList<Byte>(rollingHash.size + 9 + plaintextFrame.size)
for (b in rollingHash) input.add(b)
for (i in 7 downTo 0) input.add((seq shr (8 * i)).toByte())
input.add(recordType.toByte())
for (b in plaintextFrame) input.add(b)
rollingHash = MlsCryptoProvider.hash(input.toByteArray())
chunks += 1
}
fun append(record: AgentTextStreamRecordV1) = append(record.seq, record.recordType, record.frame)
companion object {
val CONTEXT = "marmot agent text stream transcript v1".encodeToByteArray()
fun start(
streamId: ByteArray,
startEventId: ByteArray,
): AgentTextStreamTranscriptV1 {
val input = ArrayList<Byte>()
for (b in CONTEXT) input.add(b)
input.addLengthPrefixed(streamId)
input.addLengthPrefixed(startEventId)
return AgentTextStreamTranscriptV1(
streamId = streamId,
startEventId = startEventId,
rollingHash = MlsCryptoProvider.hash(input.toByteArray()),
chunks = 0,
)
}
/**
* Resume from durable state. The caller must have bound [hash] and
* [chunkCount] to this same stream and start event — nothing here can
* check that, and resuming against a different stream silently forks
* the transcript.
*/
fun resume(
streamId: ByteArray,
startEventId: ByteArray,
hash: ByteArray,
chunkCount: Long,
): AgentTextStreamTranscriptV1 {
require(hash.size == 32) { "agent text stream transcript hash must be 32 bytes" }
require(chunkCount >= 0) { "agent text stream chunk count cannot be negative" }
return AgentTextStreamTranscriptV1(streamId, startEventId, hash.copyOf(), chunkCount)
}
}
}
@@ -0,0 +1,99 @@
/*
* 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.marmot.appComponents.agentTextStream
/**
* The Marmot binary profile's QUIC variable-length integer (RFC 9000 §16),
* used for every length prefix inside the agent-text-stream wire formats.
*
* `TlsWriter.putOpaqueVarInt` writes the same encoding but always couples it to
* the bytes it measures. The stream formats prefix values that are not the
* immediately following field (the record's `plaintext_frame` length is read
* back before the frame is consumed, and the key context hashes prefixes on
* their own), so the codec is exposed here as a standalone pair.
*/
object QuicVarInt {
const val MAX_VALUE: Long = (1L shl 62) - 1
fun encodedLength(value: Long): Int =
when {
value < 0 -> throw IllegalArgumentException("QUIC varint cannot encode a negative value")
value < 64 -> 1
value < 16_384 -> 2
value < 1_073_741_824 -> 4
value <= MAX_VALUE -> 8
else -> throw IllegalArgumentException("QUIC varint cannot encode $value")
}
fun encode(value: Long): ByteArray {
val out = ByteArray(encodedLength(value))
when (out.size) {
1 -> out[0] = value.toByte()
2 -> {
out[0] = (((value shr 8) and 0x3f) or 0x40).toByte()
out[1] = value.toByte()
}
4 -> {
out[0] = (((value shr 24) and 0x3f) or 0x80).toByte()
out[1] = (value shr 16).toByte()
out[2] = (value shr 8).toByte()
out[3] = value.toByte()
}
else -> {
out[0] = (((value shr 56) and 0x3f) or 0xc0).toByte()
for (i in 1 until 8) out[i] = (value shr (8 * (7 - i))).toByte()
}
}
return out
}
/** The decoded value plus the number of bytes it consumed. */
class Decoded(
val value: Long,
val length: Int,
)
/**
* Decode at [offset]. Rejects a truncated prefix rather than returning a
* short read: these bytes carry field boundaries, so a silent zero would
* turn a truncated record into a differently-shaped valid one.
*/
fun decode(
bytes: ByteArray,
offset: Int = 0,
): Decoded {
require(offset < bytes.size) { "QUIC varint is truncated" }
val first = bytes[offset].toInt() and 0xff
val length = 1 shl (first shr 6)
require(offset + length <= bytes.size) { "QUIC varint is truncated" }
var value = (first and 0x3f).toLong()
for (i in 1 until length) {
value = (value shl 8) or (bytes[offset + i].toLong() and 0xff)
}
return Decoded(value, length)
}
}
/** Append `varint(bytes.size) || bytes`. */
fun MutableList<Byte>.addLengthPrefixed(bytes: ByteArray) {
for (b in QuicVarInt.encode(bytes.size.toLong())) add(b)
for (b in bytes) add(b)
}
@@ -21,7 +21,11 @@
package com.vitorpamplona.quartz.marmot.mls.group
import com.vitorpamplona.quartz.marmot.appComponents.AdminPolicyV1
import com.vitorpamplona.quartz.marmot.appComponents.AppComponentIds
import com.vitorpamplona.quartz.marmot.appComponents.MarmotGroupState
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamCrypto
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamQuicPolicyV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRoles
import com.vitorpamplona.quartz.marmot.mip01Groups.MarmotGroupData
import com.vitorpamplona.quartz.marmot.mls.codec.TlsReader
import com.vitorpamplona.quartz.marmot.mls.codec.TlsWriter
@@ -65,6 +69,7 @@ import com.vitorpamplona.quartz.marmot.mls.tree.LeafNodeSource
import com.vitorpamplona.quartz.marmot.mls.tree.Lifetime
import com.vitorpamplona.quartz.marmot.mls.tree.RatchetTree
import com.vitorpamplona.quartz.marmot.mls.tree.UpdatePathNode
import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import com.vitorpamplona.quartz.utils.TimeUtils
import com.vitorpamplona.quartz.utils.mac.MacInstance
@@ -193,6 +198,20 @@ class MlsGroup private constructor(
/** The current profile's component view of this GroupContext. */
fun currentGroupState(): MarmotGroupState = MarmotGroupState.fromExtensions(groupContext.extensions)
/**
* The `nostr_group_id` this group routes kind-445 traffic under, from
* whichever profile the group is actually using.
*
* A current-profile group carries it in the `marmot.transport.nostr.routing.v1`
* component (`0x8004`); a legacy group carries it inside the monolithic
* `0xF2EE` extension. Reading only the legacy one leaves us unable to join
* any group a current-profile client created — the routing id is required
* to subscribe at all, so the failure is total rather than partial.
*/
fun currentNostrGroupId(): HexKey? =
currentGroupState().routing?.nostrGroupIdHex
?: currentMarmotData()?.nostrGroupId
/**
* The group's configured admin account identities, as lowercase hex.
*
@@ -1905,6 +1924,19 @@ class MlsGroup private constructor(
length: Int,
): ByteArray = KeySchedule.mlsExporter(epochSecrets.exporterSecret, label, context, length)
/**
* `MLS-Exporter("marmot", "agent-text-stream-quic", 32)` — the secret every
* member of this epoch derives per-stream record keys from. Per-stream and
* per-record separation is entirely in the HKDF key context, so this one
* secret covers every stream in the epoch.
*/
fun agentTextStreamSecret(): ByteArray =
exporterSecret(
AgentTextStreamCrypto.EXPORTER_LABEL,
AgentTextStreamCrypto.EXPORTER_CONTEXT,
AgentTextStreamCrypto.SECRET_LENGTH,
)
// --- External Join Support (RFC 9420 Section 8.3, 12.4.3.2) ---
/**
@@ -3258,6 +3290,37 @@ class MlsGroup private constructor(
proposals = listOf(SELF_REMOVE_PROPOSAL_TYPE),
)
/**
* Enforce the `0x8006` component's `required_member_roles` mask over
* the joining tree.
*
* A group carrying the agent-text-stream component requires each named
* role as an MLS leaf capability (`0xF2D1` receive, `0xF2D2` send,
* `0xF2D4` fanout). Advertising the component id alone is not enough —
* that only says "understands the component"; the role capability says
* "can actually do this".
*/
private fun requireAgentTextStreamRoles(
extensions: List<Extension>,
tree: RatchetTree,
myLeafIndex: Int,
) {
val policy =
AppDataDictionary
.fromExtensionsOrEmpty(extensions)[AgentTextStreamQuicPolicyV1.COMPONENT_ID]
?.let { AgentTextStreamQuicPolicyV1.decode(it) } ?: return
val required = policy.requiredRoleCapabilities()
if (required.isEmpty()) return
val myLeaf = tree.getLeaf(myLeafIndex)
requireNotNull(myLeaf) { "Joiner's leaf is blank after tree reconstruction" }
val missing = required.filterNot { myLeaf.capabilities.extensions.contains(it) }
require(missing.isEmpty()) {
"Joiner does not advertise agent text stream roles this group requires: " +
missing.joinToString { AppComponentIds.toHex(it) }
}
}
/**
* Leaf capabilities for the current profile.
*
@@ -3272,10 +3335,20 @@ class MlsGroup private constructor(
* this line a current-profile KeyPackage would be un-addable to every
* legacy group that already exists — the exact mirror of the interop
* failure the current profile was adopted to fix.
*
* `0xF2D1` is the agent-text-stream RECEIVE role, for the same reason:
* a group carrying component `0x8006` with `required_member_roles`
* naming `receive` refuses a leaf that does not advertise it. We stop
* at receive — see `CurrentProfileGroupFactory.SUPPORTED_COMPONENTS`.
*/
fun currentProfileLeafCapabilities(): Capabilities =
Capabilities(
extensions = listOf(AppDataDictionary.EXTENSION_TYPE, MarmotGroupData.EXTENSION_ID_INT),
extensions =
listOf(
AppDataDictionary.EXTENSION_TYPE,
MarmotGroupData.EXTENSION_ID_INT,
AgentTextStreamRoles.RECEIVE_CAPABILITY,
),
proposals = listOf(APP_DATA_UPDATE_PROPOSAL_TYPE, SELF_REMOVE_PROPOSAL_TYPE),
)
@@ -3526,6 +3599,14 @@ class MlsGroup private constructor(
}
}
// The agent-text-stream component (0x8006) states its own
// per-member requirement OUTSIDE MLS `required_capabilities`:
// `required_member_roles` names role capabilities every member
// must advertise. MLS cannot enforce it, so a joiner that skipped
// this check would join a group it can never satisfy and have
// every one of its commits refused by peers that do check.
requireAgentTextStreamRoles(groupContext.extensions, tree, myLeafIndex)
// Derive epoch secrets directly from memberSecret (RFC 9420 Section 8.3)
// For Welcome, epoch_secret = ExpandWithLabel(member_secret, "epoch", GroupContext, Nh)
val epochSecret =
@@ -287,9 +287,10 @@ class MlsGroupManager(
val group = MlsGroup.processWelcome(welcomeBytes, bundle)
val derivedId =
group.currentMarmotData()?.nostrGroupId
group.currentNostrGroupId()
?: throw IllegalArgumentException(
"Welcome GroupContext is missing the NostrGroupData extension — cannot derive nostrGroupId",
"Welcome GroupContext carries no nostr routing: neither the current profile's " +
"0x8004 component nor the legacy 0xF2EE extension — cannot derive nostrGroupId",
)
if (hintNostrGroupId != null && hintNostrGroupId != derivedId) {
@@ -404,17 +405,34 @@ class MlsGroupManager(
targetLeafIndex: Int,
): StagedCommit = stage(nostrGroupId) { it.removeMember(targetLeafIndex) }
/**
* Refuse a GroupContextExtensions change from a non-admin.
*
* The admin set is read profile-agnostically: a current-profile group
* keeps it in the `0x8003` admin-policy component, a legacy group inside
* the `0xF2EE` extension. Reading only the legacy one made
* `adminsConfigured` false for every current-profile group, which skipped
* the gate entirely rather than failing closed — the group would then
* refuse the commit on arrival at every peer, so the only thing the
* missing check bought was a locally-diverged copy.
*
* A group with NO admin set at all is still open: that is the MIP-01
* bootstrap state, before any admin policy has been installed.
*/
private fun requireAdminForExtensionChange(group: MlsGroup) {
val admins = group.currentAdminIdentities()
check(admins.isEmpty() || group.isLocalAdmin()) {
"MIP-01: only admins may update group extensions"
}
}
/** Stage a GroupContextExtensions change. See [StagedCommit]. */
suspend fun stageUpdateGroupExtensions(
nostrGroupId: HexKey,
extensions: List<Extension>,
): StagedCommit {
val live = requireGroup(nostrGroupId)
val currentMarmot = live.currentMarmotData()
val adminsConfigured = currentMarmot != null && currentMarmot.adminPubkeys.isNotEmpty()
check(!adminsConfigured || live.isLocalAdmin()) {
"MIP-01: only admins may update group extensions"
}
requireAdminForExtensionChange(live)
return stage(nostrGroupId) { clone ->
clone.proposeGroupContextExtensions(extensions)
clone.commit()
@@ -636,11 +654,7 @@ class MlsGroupManager(
): CommitResult =
mutex.withLock {
val group = requireGroup(nostrGroupId)
val currentMarmot = group.currentMarmotData()
val adminsConfigured = currentMarmot != null && currentMarmot.adminPubkeys.isNotEmpty()
check(!adminsConfigured || group.isLocalAdmin()) {
"MIP-01: only admins may update group extensions"
}
requireAdminForExtensionChange(group)
val retainedBefore = group.retainedSecrets()
group.proposeGroupContextExtensions(extensions)
val result = group.commit()
@@ -22,6 +22,7 @@ package com.vitorpamplona.quartz.marmot.appComponents
import com.vitorpamplona.quartz.TestResourceLoader
import com.vitorpamplona.quartz.marmot.appComponents.accountIdentityProof.AccountIdentityProofV2
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamRoles
import com.vitorpamplona.quartz.marmot.mip01Groups.MarmotGroupData
import com.vitorpamplona.quartz.marmot.mip01Groups.MlsCiphersuite
import com.vitorpamplona.quartz.marmot.mls.codec.TlsReader
@@ -106,17 +107,23 @@ class CurrentProfileGroupFactoryTest {
assertContentEquals(ByteArray(0), kpDictionary[AppComponentIds.LAST_RESORT_KEY_PACKAGE])
// Capabilities advertise the draft extension the current profile
// needs, plus the legacy 0xF2EE group-data extension.
// needs, plus the legacy 0xF2EE group-data extension and the
// agent-text-stream RECEIVE role.
//
// The extra entry is deliberate and is NOT drift from the MDK
// The extra entries are deliberate and are NOT drift from the MDK
// reference. A capability says "this client can handle it", and a
// group that REQUIRES 0xF2EE refuses to add a leaf that does not
// advertise it — so without this a current-profile KeyPackage
// would be un-addable to every legacy group that already exists.
// Advertising more than a group requires is always acceptable;
// advertising less is what gets a leaf rejected.
// group that REQUIRES 0xF2EE (legacy) or 0xF2D1 (any group MDK
// creates) refuses to add a leaf that does not advertise it — so
// without these a current-profile KeyPackage would be un-addable
// to every legacy group that already exists and to every group MDK
// makes. Advertising more than a group requires is always
// acceptable; advertising less is what gets a leaf rejected.
assertEquals(
listOf(AppDataDictionary.EXTENSION_TYPE, MarmotGroupData.EXTENSION_ID_INT),
listOf(
AppDataDictionary.EXTENSION_TYPE,
MarmotGroupData.EXTENSION_ID_INT,
AgentTextStreamRoles.RECEIVE_CAPABILITY,
),
kp.leafNode.capabilities.extensions,
)
assertTrue(
@@ -0,0 +1,346 @@
/*
* 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.marmot.appComponents.agentTextStream
import com.vitorpamplona.quartz.nip01Core.core.hexToByteArray
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import kotlin.test.Test
import kotlin.test.assertContentEquals
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
* `marmot.group.agent-text-stream.quic.v1` (0x8006), checked against the bytes
* MDK 0.9.20 actually installs. The fixture is the component state read off a
* group `wn groups create` made on the interop harness.
*
* This component is the one that decides whether MDK will let us into a group
* it created: its `required_member_roles` mask names MLS leaf capabilities, so
* a client that neither advertises `0x8006` nor the `0xF2D1` receive
* capability is refused at the Add, before any of our own code runs.
*/
class AgentTextStreamQuicPolicyV1Test {
/** `wn groups create` default: require receive, allow receive+send, 4096-byte frames. */
private val mdkDefaultState = "010300001000000000000000".hexToByteArray()
@Test
fun decodesTheComponentStateMdkInstalls() {
val policy = AgentTextStreamQuicPolicyV1.decode(mdkDefaultState)
assertEquals(AgentTextStreamRoles.RECEIVE, policy.requiredMemberRoles)
assertEquals(AgentTextStreamRoles.RECEIVE or AgentTextStreamRoles.SEND, policy.allowedMemberRoles)
assertEquals(4096L, policy.maxPlaintextFrameLen)
assertEquals(0L, policy.replayTtlSecs)
assertEquals(0, policy.paddingBucketBytes)
assertTrue(policy.requires(AgentTextStreamRoles.RECEIVE))
assertTrue(policy.allows(AgentTextStreamRoles.SEND))
}
@Test
fun ourDefaultEncodesToTheSameBytes() {
assertContentEquals(mdkDefaultState, AgentTextStreamQuicPolicyV1.userToAgentDefault().encode())
}
@Test
fun requiredRolesMapToLeafCapabilities() {
val policy = AgentTextStreamQuicPolicyV1.decode(mdkDefaultState)
assertEquals(listOf(AgentTextStreamRoles.RECEIVE_CAPABILITY), policy.requiredRoleCapabilities())
assertEquals(0xF2D1, AgentTextStreamRoles.RECEIVE_CAPABILITY)
assertEquals(0xF2D2, AgentTextStreamRoles.SEND_CAPABILITY)
assertEquals(0xF2D4, AgentTextStreamRoles.FANOUT_CAPABILITY)
}
@Test
fun roundTripsEveryFieldAtItsBound() {
val policy =
AgentTextStreamQuicPolicyV1(
requiredMemberRoles = AgentTextStreamRoles.MASK,
allowedMemberRoles = AgentTextStreamRoles.MASK,
maxPlaintextFrameLen = AgentTextStreamQuicPolicyV1.MAX_PLAINTEXT_FRAME_LEN,
replayTtlSecs = AgentTextStreamQuicPolicyV1.MAX_REPLAY_TTL_SECS,
paddingBucketBytes = AgentTextStreamQuicPolicyV1.MAX_PADDING_BUCKET_BYTES,
)
assertEquals(policy, AgentTextStreamQuicPolicyV1.decode(policy.encode()))
}
/**
* These bytes sit in signed group state. A decoder that repaired them
* would admit a member the group's own policy refuses, so every one of
* these is a hard failure rather than a default.
*/
@Test
fun rejectsStateItCannotActOn() {
// Wrong length.
assertFailsWith<IllegalArgumentException> { AgentTextStreamQuicPolicyV1.decode(ByteArray(11)) }
assertFailsWith<IllegalArgumentException> { AgentTextStreamQuicPolicyV1.decode(ByteArray(13)) }
// No required role at all.
assertFailsWith<IllegalArgumentException> {
AgentTextStreamQuicPolicyV1.decode("000300001000000000000000".hexToByteArray())
}
// Unknown role bit — a mask we cannot map to a capability.
assertFailsWith<IllegalArgumentException> {
AgentTextStreamQuicPolicyV1.decode("080b00001000000000000000".hexToByteArray())
}
// Required roles outside the allowed set.
assertFailsWith<IllegalArgumentException> {
AgentTextStreamQuicPolicyV1.decode("020100001000000000000000".hexToByteArray())
}
// Zero frame limit, and one past the app-profile cap (65519 -> 65520).
assertFailsWith<IllegalArgumentException> {
AgentTextStreamQuicPolicyV1.decode("010300000000000000000000".hexToByteArray())
}
assertFailsWith<IllegalArgumentException> {
AgentTextStreamQuicPolicyV1.decode("01030000fff0000000000000".hexToByteArray())
}
assertNull(AgentTextStreamQuicPolicyV1.decodeOrNull(ByteArray(0)))
}
}
/**
* The QUIC variable-length integer is the length prefix for every field in the
* stream wire formats, and it is also hashed into the transcript and the key
* context — so a wrong prefix is not a parse error but a different key.
*/
class QuicVarIntTest {
@Test
fun encodesEachWidthAtItsBoundary() {
assertContentEquals(byteArrayOf(0x00), QuicVarInt.encode(0))
assertContentEquals(byteArrayOf(0x20), QuicVarInt.encode(32))
assertContentEquals(byteArrayOf(0x3f), QuicVarInt.encode(63))
assertContentEquals(byteArrayOf(0x40, 0x40), QuicVarInt.encode(64))
assertContentEquals(byteArrayOf(0x7f, 0xff.toByte()), QuicVarInt.encode(16_383))
assertContentEquals(byteArrayOf(0x80.toByte(), 0x00, 0x40, 0x00), QuicVarInt.encode(16_384))
}
@Test
fun roundTripsAcrossWidths() {
for (value in listOf(0L, 1L, 63L, 64L, 16_383L, 16_384L, 1_073_741_823L, 1_073_741_824L)) {
val encoded = QuicVarInt.encode(value)
val decoded = QuicVarInt.decode(encoded)
assertEquals(value, decoded.value)
assertEquals(encoded.size, decoded.length)
}
}
@Test
fun rejectsATruncatedPrefix() {
assertFailsWith<IllegalArgumentException> { QuicVarInt.decode(ByteArray(0)) }
assertFailsWith<IllegalArgumentException> { QuicVarInt.decode(byteArrayOf(0x40)) }
}
}
class AgentTextStreamRecordV1Test {
private val streamId = ByteArray(32) { 0x5a }
@Test
fun roundTripsAFrame() {
val record = AgentTextStreamRecordV1.textDelta(streamId, 7, "hello".encodeToByteArray())
val decoded = AgentTextStreamRecordV1.decode(record.encode())
assertEquals(record, decoded)
assertEquals(record.encodedLength(), record.encode().size)
}
/**
* A newer advisory record type must not tear down an otherwise valid
* preview stream, so framing accepts types it has no semantics for.
*/
@Test
fun acceptsAnUnknownRecordType() {
val record = AgentTextStreamRecordV1(streamId, 1, recordType = 0x7f, frame = ByteArray(4))
assertEquals(0x7f, AgentTextStreamRecordV1.decode(record.encode()).recordType)
}
@Test
fun rejectsFramingItCannotTrust() {
val encoded = AgentTextStreamRecordV1.textDelta(streamId, 1, ByteArray(4)).encode()
// Trailing bytes: the sender and receiver disagree where the record ends.
assertFailsWith<IllegalArgumentException> { AgentTextStreamRecordV1.decode(encoded + 0x00) }
// Truncated.
assertFailsWith<IllegalArgumentException> { AgentTextStreamRecordV1.decode(encoded.copyOf(encoded.size - 1)) }
// A version we do not speak.
val wrongVersion = encoded.copyOf()
wrongVersion[0] = 2
assertFailsWith<IllegalArgumentException> { AgentTextStreamRecordV1.decode(wrongVersion) }
// An empty stream id names no stream.
assertFailsWith<IllegalArgumentException> {
AgentTextStreamRecordV1(ByteArray(0), 1, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, frame = ByteArray(0))
}
}
}
class AgentTextStreamCryptoTest {
private val secret = ByteArray(32) { it.toByte() }
private val streamId = ByteArray(32) { 0x11 }
private val startEventId = ByteArray(32) { 0x22 }
private fun crypto(
epoch: Long = 3,
stream: ByteArray = streamId,
) = AgentTextStreamCrypto(
secret,
AgentTextStreamKeyContextV1(
groupId = ByteArray(16) { 0x33 },
streamId = stream,
mlsEpoch = epoch,
senderId = ByteArray(32) { 0x44 },
startEventId = startEventId,
),
)
@Test
fun sealAndOpenRoundTrip() {
val c = crypto()
val record = AgentTextStreamRecordV1.textDelta(streamId, 5, "streamed text".encodeToByteArray())
val sealed = c.seal(record)
assertTrue(sealed.frame.size == record.frame.size + AgentTextStreamRecordV1.AEAD_TAG_LEN)
assertContentEquals(record.frame, c.open(sealed).frame)
}
/**
* `seq` is in both the nonce and the AAD, so a replayed or reordered
* record does not merely look wrong — it fails to open.
*/
@Test
fun aRecordDoesNotOpenAtADifferentSequence() {
val c = crypto()
val sealed = c.seal(AgentTextStreamRecordV1.textDelta(streamId, 5, "abc".encodeToByteArray()))
val replayed =
AgentTextStreamRecordV1(
streamId = sealed.streamId,
seq = 6,
recordType = sealed.recordType,
flags = sealed.flags,
frame = sealed.frame,
)
assertNull(c.openOrNull(replayed))
}
/** Nonce 0 is the base itself, and each seq flips only the low 64 bits. */
@Test
fun theNonceIsTheBaseXorTheSequence() {
val c = crypto()
val base = c.recordNonce(0)
val one = c.recordNonce(1)
assertEquals(AgentTextStreamCrypto.NONCE_LENGTH, base.size)
assertContentEquals(base.copyOf(11), one.copyOf(11))
assertEquals((base[11].toInt() xor 1).toByte(), one[11])
}
/**
* The key context is the only thing separating two streams that share one
* group exporter secret, so changing any part of it must change the key.
*/
@Test
fun everyKeyContextFieldSeparatesTheKey() {
val base = crypto().recordKey().toHexKey()
assertTrue(crypto(epoch = 4).recordKey().toHexKey() != base)
assertTrue(crypto(stream = ByteArray(32) { 0x12 }).recordKey().toHexKey() != base)
}
@Test
fun refusesARecordFromAnotherStream() {
val c = crypto()
val foreign = AgentTextStreamRecordV1.textDelta(ByteArray(32) { 0x77 }, 1, ByteArray(4))
assertFailsWith<IllegalArgumentException> { c.seal(foreign) }
}
@Test
fun refusesAFrameOverTheGroupLimit() {
val c = crypto()
val record = AgentTextStreamRecordV1.textDelta(streamId, 1, ByteArray(5000))
assertFailsWith<IllegalArgumentException> { c.seal(record, maxPlaintextFrameLen = 4096) }
}
}
class AgentTextStreamTranscriptV1Test {
private val streamId = ByteArray(32) { 0x11 }
private val startEventId = ByteArray(32) { 0x22 }
@Test
fun foldsRecordsInOrder() {
val a = AgentTextStreamTranscriptV1.start(streamId, startEventId)
val b = AgentTextStreamTranscriptV1.start(streamId, startEventId)
a.append(0, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "one".encodeToByteArray())
a.append(1, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "two".encodeToByteArray())
b.append(1, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "two".encodeToByteArray())
b.append(0, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "one".encodeToByteArray())
assertEquals(2, a.chunkCount)
assertTrue(!a.hash.contentEquals(b.hash), "a reordered stream must not hash the same")
}
@Test
fun resumesFromDurableState() {
val original = AgentTextStreamTranscriptV1.start(streamId, startEventId)
original.append(0, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "one".encodeToByteArray())
val resumed = AgentTextStreamTranscriptV1.resume(streamId, startEventId, original.hash, original.chunkCount)
resumed.append(1, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "two".encodeToByteArray())
original.append(1, AgentTextStreamRecordV1.TYPE_TEXT_DELTA, "two".encodeToByteArray())
assertContentEquals(original.hash, resumed.hash)
assertEquals(original.chunkCount, resumed.chunkCount)
}
}
class AgentTextStreamStartTest {
@Test
fun readsTheAnchorTags() {
val tags =
AgentTextStreamStart.tags(
streamId = ByteArray(32) { 0x5a }.toHexKey(),
brokerCandidates = listOf("https://broker.example:4443", "https://alt.example:4443"),
)
val start = assertNotNull(AgentTextStreamStart.fromTags(AgentTextStreamStart.KIND, tags))
assertTrue(start.isQuicRoute)
assertEquals(2, start.brokerCandidates.size)
assertEquals(ByteArray(32) { 0x5a }.toHexKey(), start.streamId)
}
@Test
fun defaultsAMissingRouteToQuicButNeverAMissingStreamId() {
assertTrue(
assertNotNull(
AgentTextStreamStart.fromTags(AgentTextStreamStart.KIND, arrayOf(arrayOf("stream", "ab"))),
).isQuicRoute,
)
assertNull(AgentTextStreamStart.fromTags(AgentTextStreamStart.KIND, arrayOf(arrayOf("route", "quic"))))
assertNull(AgentTextStreamStart.fromTags(9, arrayOf(arrayOf("stream", "ab"))))
}
@Test
fun readsTheFinalTranscriptTags() {
val tags = AgentTextStreamFinal.tags("aa", "bb", 12)
val final = assertNotNull(AgentTextStreamFinal.fromTags(tags))
assertEquals("aa", final.streamId)
assertEquals("bb", final.transcriptHash)
assertEquals(12, final.chunkCount)
assertNull(AgentTextStreamFinal.fromTags(arrayOf(arrayOf("stream", "aa"))))
}
}
@@ -0,0 +1,139 @@
/*
* 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.marmot.mls.group
import com.vitorpamplona.quartz.marmot.appComponents.CurrentProfileGroupFactory
import com.vitorpamplona.quartz.marmot.appComponents.GroupProfileV1
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamQuicPolicyV1
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import kotlinx.coroutines.runBlocking
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
* A joiner has to learn the group's `nostr_group_id` from the Welcome itself —
* it is the `h` tag every kind-445 event in the group carries, so without it
* the joiner cannot even subscribe.
*
* Reading it only from the legacy `0xF2EE` extension made us unable to join
* ANY group a current-profile client created: MDK's welcome arrived, decrypted,
* and was then thrown away with "GroupContext is missing the NostrGroupData
* extension". The routing id had been there the whole time, in the
* `marmot.transport.nostr.routing.v1` component.
*/
class CurrentProfileWelcomeTest {
private fun signer(seed: Byte) = NostrSignerInternal(KeyPair(ByteArray(32) { seed }))
private val nostrGroupId = ByteArray(32) { 0x4d }
private fun aGroup(agentTextStream: AgentTextStreamQuicPolicyV1? = null) =
runBlocking {
CurrentProfileGroupFactory.createGroup(
signer = signer(0x11),
nostrGroupId = nostrGroupId,
relays = listOf("wss://relay.example"),
profile = GroupProfileV1("Routing", ""),
agentTextStream = agentTextStream,
)
}
@Test
fun theRoutingIdComesFromTheCurrentProfileComponent() {
val group = aGroup()
assertEquals(nostrGroupId.toHexKey(), group.currentNostrGroupId())
// The legacy extension is genuinely absent — this is not a group that
// happens to carry both.
assertNull(group.currentMarmotData())
}
@Test
fun aJoinerReadsTheSameRoutingIdOutOfTheWelcome() =
runBlocking<Unit> {
val group = aGroup()
val inviteeSigner = signer(0x33)
val invitee = CurrentProfileGroupFactory.createKeyPackage(inviteeSigner)
group.proposeAdd(invitee.keyPackage.toTlsBytes())
val commit = group.commit()
val welcome = assertNotNull(commit.welcomeBytes, "adding a member must produce a Welcome")
val joined = MlsGroup.processWelcome(welcome, invitee)
assertEquals(nostrGroupId.toHexKey(), joined.currentNostrGroupId())
assertEquals(group.currentGroupState().profile?.name, joined.currentGroupState().profile?.name)
}
/**
* The `0x8006` policy names MLS leaf capabilities every member must
* advertise. Our current-profile leaf advertises `receive` only, so a
* group that also requires `send` must be refused at join rather than
* joined into a state where every commit we make is rejected by peers.
*/
@Test
fun aJoinerRefusesAGroupWhoseStreamRolesItCannotFill() =
runBlocking<Unit> {
val group =
aGroup(
AgentTextStreamQuicPolicyV1(
requiredMemberRoles = AgentTextStreamRolesFixture.RECEIVE_AND_SEND,
allowedMemberRoles = AgentTextStreamRolesFixture.RECEIVE_AND_SEND,
maxPlaintextFrameLen = 4096,
replayTtlSecs = 0,
paddingBucketBytes = 0,
),
)
val invitee = CurrentProfileGroupFactory.createKeyPackage(signer(0x44))
group.proposeAdd(invitee.keyPackage.toTlsBytes())
val welcome = assertNotNull(group.commit().welcomeBytes)
val failure = assertFailsWith<IllegalArgumentException> { MlsGroup.processWelcome(welcome, invitee) }
assertTrue(
failure.message.orEmpty().contains("agent text stream roles"),
"expected a role-capability refusal, got: ${failure.message}",
)
}
@Test
fun aJoinerAcceptsAGroupRequiringOnlyTheReceiveRole() =
runBlocking<Unit> {
val group = aGroup(AgentTextStreamQuicPolicyV1.userToAgentDefault())
val invitee = CurrentProfileGroupFactory.createKeyPackage(signer(0x55))
group.proposeAdd(invitee.keyPackage.toTlsBytes())
val welcome = assertNotNull(group.commit().welcomeBytes)
val joined = MlsGroup.processWelcome(welcome, invitee)
assertEquals(
AgentTextStreamQuicPolicyV1.userToAgentDefault(),
joined.currentGroupState().agentTextStream,
)
// Both sides derive the same per-epoch stream secret.
assertEquals(group.agentTextStreamSecret().toHexKey(), joined.agentTextStreamSecret().toHexKey())
}
}
private object AgentTextStreamRolesFixture {
const val RECEIVE_AND_SEND = 0x03
}