fix: make the sync benchmark opt-in; stop advertising the QUIC preview path

Two unrelated things that both amount to not paying for something nobody
asked for.

**`MirrorSyncThroughputTest` is a benchmark, so it now opts in.** It
preloaded a million events and pulled them over a real WebSocket on every
ordinary test run: 4,584 s of `:geode:test`'s 4,636 s — 98.9% of the
module's test time for one test that asserts nothing about correctness and
reported `skipped` at the end anyway. Every other benchmark in the module
is already gated this way (`perf.LoadBenchmark`). It now bails before
building anything, and enables on `-DrunLoadBenchmark=true` OR on any of
its own sizing properties, so every invocation its kdoc documents still
runs it — naming a size is itself the opt-in. Measured after: the test
takes 5 ms, the module takes 64.8 s, and `-DsyncN=2000` still prints a
throughput number.

**The agent text stream QUIC path is kept but no longer advertised, and
nothing starts it.** Nothing in the deployed network publishes those
previews. So:

- `SUPPORTED_COMPONENTS` drops `0x8006` and the leaf capabilities drop
  `0xF2D1`/`0xF2D2`/`0xF2D4`. A capability is a standing promise to every
  peer that reads our KeyPackage, and one for a path nobody exercises
  costs something and buys nothing. The captured reference KeyPackage in
  our own conformance vector does not advertise `0x8006` either.
- The Android chat screen no longer builds a stream watcher and dials the
  brokers a kind:1200 advertises. That was a UDP connection attempt to a
  third-party endpoint on every feed change, on behalf of a feature with
  nothing to show — a service we start, not a capability we hold.

The implementation stays and stays tested: `:marmotQuic`, the codecs,
`amy marmot stream`, the direct path, the certificate pinning and the
interop tests are all untouched. The module README records the posture and
the exact way back.

Three tests asserted the old advertisement and were reworked rather than
deleted. The role-enforcement gate is still covered — the tests now build
leaves that explicitly carry the roles, which is the better shape anyway,
since a test that exercised the gate through OUR default was really
asserting the default and stopped testing the gate the moment it changed.
A new test pins the new default: our KeyPackage carries no role and is
therefore refused by a group requiring one. That refusal is the deliberate
cost, so it is asserted rather than discovered.

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 18:29:36 +00:00
parent 9e1120fc07
commit 1794ed5249
8 changed files with 165 additions and 69 deletions
@@ -44,10 +44,8 @@ import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.Color import androidx.compose.ui.graphics.Color
import androidx.compose.ui.platform.LocalContext import androidx.compose.ui.platform.LocalContext
import androidx.compose.ui.unit.dp import androidx.compose.ui.unit.dp
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import androidx.lifecycle.viewmodel.compose.viewModel import androidx.lifecycle.viewmodel.compose.viewModel
import com.vitorpamplona.amethyst.R import com.vitorpamplona.amethyst.R
import com.vitorpamplona.amethyst.commons.marmot.MarmotAgentStreamWatcher
import com.vitorpamplona.amethyst.commons.resources.Res import com.vitorpamplona.amethyst.commons.resources.Res
import com.vitorpamplona.amethyst.commons.resources.marmot_group_default_name import com.vitorpamplona.amethyst.commons.resources.marmot_group_default_name
import com.vitorpamplona.amethyst.ui.actions.MentionPreservingInputTransformation import com.vitorpamplona.amethyst.ui.actions.MentionPreservingInputTransformation
@@ -78,7 +76,6 @@ import com.vitorpamplona.quartz.nip01Core.core.HexKey
import kotlinx.collections.immutable.ImmutableList import kotlinx.collections.immutable.ImmutableList
import kotlinx.collections.immutable.persistentListOf import kotlinx.collections.immutable.persistentListOf
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
@Composable @Composable
@@ -123,31 +120,19 @@ fun MarmotGroupChatView(
} }
} }
// The live agent-preview watcher. It follows the newest kind:1200 in the // The live agent-preview watcher is NOT started here.
// group and folds the QUIC records behind it; a group with no stream, no //
// broker candidate or no reachable broker simply never shows a preview, // Opening the chat used to build a [MarmotAgentStreamWatcher] and call
// and the durable kind:9 still arrives as ordinary chat either way. // `watchLatest` on every feed change, which dials the QUIC brokers a
val marmot = accountViewModel.account.marmotManager // kind:1200 advertises. Nothing in the deployed network publishes those
val streamScope = rememberCoroutineScope() // streams, so that was a UDP connection attempt to a third-party endpoint
val streamWatcher = // on behalf of a feature no one is using — a service we start, not a
remember(nostrGroupId, marmot) { // capability we hold.
marmot?.let { //
MarmotAgentStreamWatcher(it, accountViewModel.account.marmotStreamTransport, streamScope) // The watcher, the transport and the banner all still exist and are still
} // tested; `amy marmot stream watch` drives the same code on demand. Wiring
} // it back is re-adding the watcher, the LaunchedEffect and the banner
val streamPreview by (streamWatcher?.preview ?: remember { MutableStateFlow(null) }).collectAsStateWithLifecycle() // below, once there is something to watch.
// Re-check on every feed change: a kind:1200 arrives as an ordinary group
// message, so "the feed moved" is exactly when a new stream may have been
// anchored. watchLatest is idempotent for a stream already being followed.
val feedState by feedViewModel.feedState.feedContent.collectAsStateWithLifecycle()
LaunchedEffect(feedState, streamWatcher) {
streamWatcher?.watchLatest(nostrGroupId)
}
DisposableEffect(streamWatcher) {
onDispose { streamWatcher?.stop() }
}
Column(Modifier.fillMaxHeight()) { Column(Modifier.fillMaxHeight()) {
Column( Column(
@@ -166,8 +151,6 @@ fun MarmotGroupChatView(
) )
} }
AgentStreamPreviewBanner(streamPreview)
Spacer(modifier = DoubleVertSpacer) Spacer(modifier = DoubleVertSpacer)
MarmotGroupMessageComposer( MarmotGroupMessageComposer(
@@ -81,6 +81,15 @@ import kotlin.test.Test
* *
* Size with `-DsyncN` (default 1,000,000). Timed from first byte to the * Size with `-DsyncN` (default 1,000,000). Timed from first byte to the
* downstream reaching the target (or plateauing). * downstream reaching the target (or plateauing).
*
* **Opt-in.** This is a benchmark, not a regression test: it preloads a
* million events and pulls them over a real WebSocket, which took 4,584 s of
* `:geode:test`'s 4,636 s total — 98.9% of the module's test time for one test
* that asserts nothing about correctness. Every other benchmark in this module
* is already gated the same way (see `perf.LoadBenchmark`), and every
* invocation documented above passes `-DsyncN` or `-DsyncSourceUrl`, so those
* still run it. A plain `./gradlew test` — which is what the pre-push hook
* runs — now skips it in milliseconds.
*/ */
class MirrorSyncThroughputTest { class MirrorSyncThroughputTest {
// geode's real relay config (deferred FTS, live negentropy index) — not the // geode's real relay config (deferred FTS, live negentropy index) — not the
@@ -148,9 +157,30 @@ class MirrorSyncThroughputTest {
return String(out) return String(out)
} }
/**
* True when someone actually asked for a throughput number: either the
* module-wide benchmark switch, or any of this test's own sizing/source
* properties. Naming a size IS the opt-in — a run that says `-DsyncN=…`
* plainly wants the measurement and should not need a second flag.
*/
private val enabled =
System.getProperty("runLoadBenchmark") == "true" ||
System.getProperty("syncN") != null ||
System.getProperty("syncSourceUrl") != null
@Test @Test
fun mirrorSyncThroughput() = fun mirrorSyncThroughput() =
runBlocking { runBlocking {
// Bail before building anything. The old code decided nothing up
// front and spent over an hour preloading and syncing a million
// events on every ordinary test run.
if (!enabled) {
println(
"[skip] mirrorSyncThroughput — benchmark. Enable with -DrunLoadBenchmark=true, " +
"or size it directly with -DsyncN=… / -DsyncSourceUrl=…",
)
return@runBlocking
}
val n = System.getProperty("syncN")?.toInt() ?: 1_000_000 val n = System.getProperty("syncN")?.toInt() ?: 1_000_000
val externalUrl = System.getProperty("syncSourceUrl") val externalUrl = System.getProperty("syncSourceUrl")
val expect = System.getProperty("syncExpect")?.toInt() ?: n val expect = System.getProperty("syncExpect")?.toInt() ?: n
+18
View File
@@ -107,6 +107,24 @@ amy marmot stream watch GID --stream-id …
amy marmot stream finish GID --stream-id … --transcript-hash … --chunk-count N "hello" amy marmot stream finish GID --stream-id … --transcript-hash … --chunk-count N "hello"
``` ```
## Not wired into the app
The implementation is complete and tested, and nothing in the app starts it.
Nothing in the deployed network publishes agent text stream previews, so the
Android chat screen no longer builds a watcher and dials the brokers a kind:1200
advertises, and our published KeyPackage no longer advertises component `0x8006`
or the `receive`/`send`/`fanout` role capabilities. A capability is a standing
promise to every peer that reads the KeyPackage; making one for a path nobody
exercises costs something and buys nothing.
What that leaves: the codecs, this module, the CLI (`amy marmot stream …`) and
the interop tests all still work and still run. Turning the feature back on is
re-adding `AppComponentIds.AGENT_TEXT_STREAM_QUIC_V1` to
`CurrentProfileGroupFactory.SUPPORTED_COMPONENTS`, the three roles to
`MlsGroup.currentProfileLeafCapabilities()`, and the watcher to
`MarmotGroupChatView`.
## Not done ## Not done
- The Android GUI renders previews but does not originate a stream — that is - The Android GUI renders previews but does not originate a stream — that is
@@ -64,10 +64,13 @@ object CurrentProfileGroupFactory {
* later as a group we cannot actually participate in. Add an id here only * later as a group we cannot actually participate in. Add an id here only
* when the component is implemented. * when the component is implemented.
* *
* `0x8006` (agent-text-stream over QUIC) is listed for every role * `0x8006` (agent-text-stream over QUIC) is deliberately absent even
* [MlsGroup.currentProfileLeafCapabilities] advertises — receive, send and * though it is implemented. Nothing in the deployed network uses the QUIC
* fanout. Publishing needs durable per-stream sequence state so a restart * preview path, and this list is a promise rather than a description: a
* cannot reuse an AEAD nonce, and that store exists. * group may require any id we advertise, and we would then owe every peer
* behaviour for a feature no one exercises. The code stays (`:marmotQuic`,
* the codecs, `amy marmot stream`), and the id goes back on the list the
* day the feature is actually used.
*/ */
val SUPPORTED_COMPONENTS: List<Int> = val SUPPORTED_COMPONENTS: List<Int> =
listOf( listOf(
@@ -78,7 +81,6 @@ object CurrentProfileGroupFactory {
AppComponentIds.ADMIN_POLICY_V1, AppComponentIds.ADMIN_POLICY_V1,
AppComponentIds.NOSTR_ROUTING_V1, AppComponentIds.NOSTR_ROUTING_V1,
AppComponentIds.MESSAGE_RETENTION_V1, AppComponentIds.MESSAGE_RETENTION_V1,
AppComponentIds.AGENT_TEXT_STREAM_QUIC_V1,
AppComponentIds.ACCOUNT_IDENTITY_PROOF_V2, AppComponentIds.ACCOUNT_IDENTITY_PROOF_V2,
AppComponentIds.GROUP_ENCRYPTED_MEDIA_V2, AppComponentIds.GROUP_ENCRYPTED_MEDIA_V2,
AppComponentIds.GROUP_LIFECYCLE_V1, AppComponentIds.GROUP_LIFECYCLE_V1,
@@ -210,6 +210,12 @@ object AgentTextStreamRoles {
const val SEND_CAPABILITY = 0xF2D2 const val SEND_CAPABILITY = 0xF2D2
const val FANOUT_CAPABILITY = 0xF2D4 const val FANOUT_CAPABILITY = 0xF2D4
/**
* Every role capability, for callers that need to ask "does this leaf
* advertise any of them" without enumerating the three by hand.
*/
val ALL_CAPABILITIES = listOf(RECEIVE_CAPABILITY, SEND_CAPABILITY, FANOUT_CAPABILITY)
fun capabilityFor(role: Int): Int = fun capabilityFor(role: Int): Int =
when (role) { when (role) {
RECEIVE -> RECEIVE_CAPABILITY RECEIVE -> RECEIVE_CAPABILITY
@@ -25,7 +25,6 @@ import com.vitorpamplona.quartz.marmot.appComponents.AppComponentIds
import com.vitorpamplona.quartz.marmot.appComponents.MarmotGroupState import com.vitorpamplona.quartz.marmot.appComponents.MarmotGroupState
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamCrypto import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamCrypto
import com.vitorpamplona.quartz.marmot.appComponents.agentTextStream.AgentTextStreamQuicPolicyV1 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.mip01Groups.MarmotGroupData
import com.vitorpamplona.quartz.marmot.mls.codec.TlsReader import com.vitorpamplona.quartz.marmot.mls.codec.TlsReader
import com.vitorpamplona.quartz.marmot.mls.codec.TlsWriter import com.vitorpamplona.quartz.marmot.mls.codec.TlsWriter
@@ -3442,21 +3441,21 @@ class MlsGroup private constructor(
listOf( listOf(
AppDataDictionary.EXTENSION_TYPE, AppDataDictionary.EXTENSION_TYPE,
MarmotGroupData.EXTENSION_ID_INT, MarmotGroupData.EXTENSION_ID_INT,
// All three agent-stream roles, matching what MDK puts // The agent-stream roles (`0xF2D1` receive, `0xF2D2`
// on every KeyPackage it publishes. `receive` is the // send, `0xF2D4` fanout) are deliberately NOT here.
// baseline compatibility role; `send` says we can
// originate preview records, which we can now that the
// publisher, the raw-QUIC binding and the app wiring
// exist; `fanout` says records may be forwarded on our
// behalf, which is what using a broker at all means.
// //
// A capability is only a claim about what we support, // The implementation exists and stays — see
// not a duty to stream: a group that requires `send` // [AgentTextStreamRoles] and the `:marmotQuic` module —
// wants members that COULD originate, and a member that // but nothing in the deployed network uses the QUIC
// never does is a quiet member, not a broken one. // preview path, and an advertised capability is a
AgentTextStreamRoles.RECEIVE_CAPABILITY, // standing promise to every peer that reads our
AgentTextStreamRoles.SEND_CAPABILITY, // KeyPackage. Advertising a role no one exercises buys
AgentTextStreamRoles.FANOUT_CAPABILITY, // nothing and commits us to answering for it; the
// reference KeyPackage in our own conformance vector
// does not advertise it either.
//
// Re-adding them is a one-line change once the feature
// is actually in use.
), ),
proposals = listOf(APP_DATA_UPDATE_PROPOSAL_TYPE, SELF_REMOVE_PROPOSAL_TYPE), proposals = listOf(APP_DATA_UPDATE_PROPOSAL_TYPE, SELF_REMOVE_PROPOSAL_TYPE),
) )
@@ -22,7 +22,6 @@ package com.vitorpamplona.quartz.marmot.appComponents
import com.vitorpamplona.quartz.TestResourceLoader import com.vitorpamplona.quartz.TestResourceLoader
import com.vitorpamplona.quartz.marmot.appComponents.accountIdentityProof.AccountIdentityProofV2 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.MarmotGroupData
import com.vitorpamplona.quartz.marmot.mip01Groups.MlsCiphersuite import com.vitorpamplona.quartz.marmot.mip01Groups.MlsCiphersuite
import com.vitorpamplona.quartz.marmot.mls.codec.TlsReader import com.vitorpamplona.quartz.marmot.mls.codec.TlsReader
@@ -111,21 +110,22 @@ class CurrentProfileGroupFactoryTest {
// agent-text-stream roles — the same set MDK puts on every // agent-text-stream roles — the same set MDK puts on every
// KeyPackage it publishes. // KeyPackage it publishes.
// //
// The extra entries are deliberate and are NOT drift from the MDK // `0xF2EE` is deliberate and is NOT drift from the MDK reference:
// reference. A capability says "this client can handle it", and a // a legacy group REQUIRES it, and a group refuses to add a leaf
// group that REQUIRES 0xF2EE (legacy) or a role (any group MDK // that does not advertise what it requires, so without it a
// creates) refuses to add a leaf that does not advertise it — so // current-profile KeyPackage would be un-addable to every legacy
// without these a current-profile KeyPackage would be un-addable // group that already exists.
// to every legacy group that already exists and to every group MDK //
// makes. Advertising more than a group requires is always // The agent-stream roles are deliberately absent. Advertising more
// acceptable; advertising less is what gets a leaf rejected. // than a group requires is harmless to that group but is not free:
// it is a standing claim to every peer that reads this KeyPackage,
// and nothing in the deployed network uses the QUIC preview path.
// The reference KeyPackage in `mls/marmot-current-profile.json`
// does not advertise `0x8006` either.
assertEquals( assertEquals(
listOf( listOf(
AppDataDictionary.EXTENSION_TYPE, AppDataDictionary.EXTENSION_TYPE,
MarmotGroupData.EXTENSION_ID_INT, MarmotGroupData.EXTENSION_ID_INT,
AgentTextStreamRoles.RECEIVE_CAPABILITY,
AgentTextStreamRoles.SEND_CAPABILITY,
AgentTextStreamRoles.FANOUT_CAPABILITY,
), ),
kp.leafNode.capabilities.extensions, kp.leafNode.capabilities.extensions,
) )
@@ -135,7 +135,7 @@ class CurrentProfileWelcomeTest {
* admits us. * admits us.
*/ */
@Test @Test
fun aJoinerFillsEveryRoleTheProfileDefines() = fun aLeafAdvertisingEveryRoleFillsAGroupThatRequiresThem() =
runBlocking<Unit> { runBlocking<Unit> {
val group = val group =
aGroup( aGroup(
@@ -147,7 +147,7 @@ class CurrentProfileWelcomeTest {
paddingBucketBytes = 0, paddingBucketBytes = 0,
), ),
) )
val invitee = CurrentProfileGroupFactory.createKeyPackage(signer(0x66)) val invitee = allRolesKeyPackage(signer(0x66))
group.proposeAdd(invitee.keyPackage.toTlsBytes()) group.proposeAdd(invitee.keyPackage.toTlsBytes())
val welcome = assertNotNull(group.commit().welcomeBytes) val welcome = assertNotNull(group.commit().welcomeBytes)
@@ -155,13 +155,71 @@ class CurrentProfileWelcomeTest {
assertEquals(nostrGroupId.toHexKey(), joined.currentNostrGroupId()) assertEquals(nostrGroupId.toHexKey(), joined.currentNostrGroupId())
} }
/** A current-profile leaf with the `send` and `fanout` roles stripped. */ /**
private suspend fun receiveOnlyKeyPackage(signer: NostrSignerInternal): KeyPackageBundle { * Our published KeyPackage advertises NO agent-stream role, and is
* therefore refused by a group that requires one.
*
* That refusal is the deliberate cost of not advertising, so it is asserted
* rather than discovered: the implementation is still here and still
* tested, but a capability is a standing promise to every peer that reads
* the KeyPackage, and we do not make one for a path nothing uses. If this
* test starts failing because the default advertises a role again, that is
* a decision to take on purpose, not a drift to absorb.
*/
@Test
fun ourDefaultLeafAdvertisesNoStreamRoleAndIsRefusedByAGroupThatNeedsOne() =
runBlocking<Unit> {
val group = aGroup(AgentTextStreamQuicPolicyV1.userToAgentDefault())
val invitee = CurrentProfileGroupFactory.createKeyPackage(signer(0x77))
assertTrue(
invitee.keyPackage.leafNode.capabilities.extensions
.none { it in AgentTextStreamRoles.ALL_CAPABILITIES },
"the default leaf must carry no agent-stream role, got ${invitee.keyPackage.leafNode.capabilities.extensions}",
)
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}",
)
}
/** A current-profile leaf carrying the `receive` role and nothing beyond it. */
private suspend fun receiveOnlyKeyPackage(signer: NostrSignerInternal): KeyPackageBundle = keyPackageAdvertising(signer, listOf(AgentTextStreamRoles.RECEIVE_CAPABILITY))
/** A current-profile leaf carrying every role the profile defines. */
private suspend fun allRolesKeyPackage(signer: NostrSignerInternal): KeyPackageBundle =
keyPackageAdvertising(
signer,
listOf(
AgentTextStreamRoles.RECEIVE_CAPABILITY,
AgentTextStreamRoles.SEND_CAPABILITY,
AgentTextStreamRoles.FANOUT_CAPABILITY,
),
)
/**
* A current-profile KeyPackage whose leaf advertises exactly [roles] on top
* of the default capability set.
*
* The default set no longer carries any agent-stream role, so these tests
* build the leaf they need instead of relying on it. That is the right
* shape regardless: a test that asserted the gate through OUR default was
* really asserting the default, and stopped testing the gate the moment the
* default changed — which is exactly what happened.
*/
private suspend fun keyPackageAdvertising(
signer: NostrSignerInternal,
roles: List<Int>,
): KeyPackageBundle {
val full = CurrentProfileGroupFactory.createKeyPackage(signer) val full = CurrentProfileGroupFactory.createKeyPackage(signer)
val reduced = val reduced =
MlsGroup.currentProfileLeafCapabilities().let { MlsGroup.currentProfileLeafCapabilities().let {
Capabilities( Capabilities(
extensions = it.extensions.filterNot { ext -> ext == AgentTextStreamRoles.SEND_CAPABILITY || ext == AgentTextStreamRoles.FANOUT_CAPABILITY }, extensions = it.extensions + roles,
proposals = it.proposals, proposals = it.proposals,
) )
} }
@@ -188,10 +246,10 @@ class CurrentProfileWelcomeTest {
} }
@Test @Test
fun aJoinerAcceptsAGroupRequiringOnlyTheReceiveRole() = fun aLeafAdvertisingReceiveJoinsAGroupRequiringOnlyThatRole() =
runBlocking<Unit> { runBlocking<Unit> {
val group = aGroup(AgentTextStreamQuicPolicyV1.userToAgentDefault()) val group = aGroup(AgentTextStreamQuicPolicyV1.userToAgentDefault())
val invitee = CurrentProfileGroupFactory.createKeyPackage(signer(0x55)) val invitee = receiveOnlyKeyPackage(signer(0x55))
group.proposeAdd(invitee.keyPackage.toTlsBytes()) group.proposeAdd(invitee.keyPackage.toTlsBytes())
val welcome = assertNotNull(group.commit().welcomeBytes) val welcome = assertNotNull(group.commit().welcomeBytes)