From faa9aa732e439006d7dfb5fa14414f9276d37099 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 10 Jul 2026 19:30:55 +0000 Subject: [PATCH] feat(concord): add ConcordSessionRegistry account-wide coordinator Holds one live ConcordCommunitySession per joined community, fans inbound kind-1059 wraps out to the owning session, and exposes the union of control/channel plane addresses to subscribe. sync() reconciles sessions against the joined list while preserving already-folded state. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01CzJ2Cwo8tg4oZq43oRa3ig --- .../model/concord/ConcordSessionRegistry.kt | 107 +++++++++++++++++ .../concord/ConcordSessionRegistryTest.kt | 110 ++++++++++++++++++ 2 files changed, 217 insertions(+) create mode 100644 commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistry.kt create mode 100644 commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistryTest.kt diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistry.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistry.kt new file mode 100644 index 0000000000..f4d6f7ab54 --- /dev/null +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistry.kt @@ -0,0 +1,107 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.amethyst.commons.model.concord + +import com.vitorpamplona.amethyst.commons.util.KmpLock +import com.vitorpamplona.amethyst.commons.util.withLock +import com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityListEntry +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey + +/** + * The account-wide coordinator for every joined Concord community: it holds one + * live [ConcordCommunitySession] per community id and fans inbound stream wraps + * out to whichever session owns them. This is the read-path analog of + * [ConcordChannelListState] on the write path — the list yields the joined + * [ConcordCommunityListEntry] set, this expands each into a folding read-model. + * + * The app layer drives it from two directions: + * - [sync] whenever the joined list changes (from `liveCommunities`), which + * creates sessions for new communities and drops sessions for departed ones + * while **preserving** the already-folded state of the ones that remain. + * - [ingest] for every inbound kind-1059 wrap, which routes it to the matching + * session (control plane → re-fold; channel plane → re-project messages). + * + * [subscribeAddresses] returns the union of every session's control- and + * channel-plane addresses — exactly the `authors` set a subscription must watch + * for kind-1059 wraps. Thread-safe: the ingest path and UI share one instance. + */ +class ConcordSessionRegistry { + private val lock = KmpLock() + + // communityId -> live folding session. Insertion-ordered for stable iteration. + private val sessions = LinkedHashMap() + + /** + * Reconcile the held sessions with the current joined [entries]. Sessions for + * communities still present are kept as-is (their folded state survives); + * sessions for communities no longer joined are dropped; new communities get a + * fresh session. Returns the set of community ids whose sessions were created. + */ + fun sync( + entries: List, + myPubKey: HexKey, + ): Set = + lock.withLock { + val wanted = entries.associateBy { it.id } + // Drop sessions for communities we've left. + sessions.keys.retainAll(wanted.keys) + // Add sessions for newly-joined communities. + val created = mutableSetOf() + for ((id, entry) in wanted) { + if (id !in sessions) { + sessions[id] = ConcordCommunitySession(entry, myPubKey) + created += id + } + } + created + } + + fun sessionFor(communityId: HexKey): ConcordCommunitySession? = lock.withLock { sessions[communityId] } + + fun sessions(): List = lock.withLock { sessions.values.toList() } + + /** The union of control- and channel-plane addresses across all sessions to subscribe to. */ + fun subscribeAddresses(): Set = + lock.withLock { + val out = HashSet() + for (session in sessions.values) { + out += session.controlPlaneAddress + out += session.channelAddresses() + } + out + } + + /** + * Routes an inbound stream [wrap] to whichever session recognizes it. Returns + * true if some session applied it. A wrap belongs to at most one plane, so the + * first accepting session wins. + */ + fun ingest(wrap: Event): Boolean { + val snapshot = lock.withLock { sessions.values.toList() } + for (session in snapshot) { + if (session.ingest(wrap)) return true + } + return false + } + + fun clear() = lock.withLock { sessions.clear() } +} diff --git a/commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistryTest.kt b/commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistryTest.kt new file mode 100644 index 0000000000..29b600f620 --- /dev/null +++ b/commons/src/commonTest/kotlin/com/vitorpamplona/amethyst/commons/model/concord/ConcordSessionRegistryTest.kt @@ -0,0 +1,110 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.amethyst.commons.model.concord + +import com.vitorpamplona.amethyst.commons.actions.ConcordActions +import com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityFactory +import com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityListEntry +import com.vitorpamplona.quartz.concord.cord02Community.NewConcordCommunity +import com.vitorpamplona.quartz.nip01Core.core.toHexKey +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal +import kotlinx.coroutines.test.runTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue + +class ConcordSessionRegistryTest { + private val owner = NostrSignerInternal(KeyPair()) + + private fun entryFor( + community: NewConcordCommunity, + name: String, + ) = ConcordCommunityListEntry( + id = community.communityIdHex, + owner = community.ownerPubKey, + ownerSalt = community.ownerSalt.toHexKey(), + root = community.communityRoot.toHexKey(), + rootEpoch = community.rootEpoch, + relays = listOf("wss://r.example"), + name = name, + ) + + @Test + fun syncsSessionsRoutesWrapsAndDropsDepartedCommunities() = + runTest { + val alpha = ConcordCommunityFactory.create(owner, "Alpha", createdAt = 1L, relays = listOf("wss://r.example")) + val beta = ConcordCommunityFactory.create(owner, "Beta", createdAt = 1L, relays = listOf("wss://r.example")) + + val registry = ConcordSessionRegistry() + + // First sync creates a session for each joined community. + val created = registry.sync(listOf(entryFor(alpha, "Alpha"), entryFor(beta, "Beta")), owner.pubKey) + assertEquals(setOf(alpha.communityIdHex, beta.communityIdHex), created) + assertNotNull(registry.sessionFor(alpha.communityIdHex)) + assertNotNull(registry.sessionFor(beta.communityIdHex)) + + // Both control-plane addresses are in the subscribe set from the entries alone. + assertTrue(registry.subscribeAddresses().contains(alpha.controlPlane.publicKeyHex)) + assertTrue(registry.subscribeAddresses().contains(beta.controlPlane.publicKeyHex)) + + // A genesis control wrap routes to Alpha's session and folds it. + alpha.genesisWraps.forEach { assertTrue(registry.ingest(it)) } + val alphaState = registry.sessionFor(alpha.communityIdHex)!!.state.value + assertEquals("Alpha", alphaState?.metadata?.name) + + // After the fold, Alpha's #general channel plane joins the subscribe set. + val general = ConcordActions.publicChannel(alpha.communityRoot, alpha.generalChannelId, alpha.rootEpoch) + assertTrue(registry.subscribeAddresses().contains(general.publicKeyHex)) + + // A channel message routes to Alpha's #general flow, not Beta. + val msg = ConcordActions.buildChannelMessage(owner, general, alpha.generalChannelIdHex, alpha.rootEpoch, "gm", 2L) + assertTrue(registry.ingest(msg)) + assertEquals( + 1, + registry + .sessionFor(alpha.communityIdHex)!! + .messagesFlow(alpha.generalChannelIdHex) + .value.size, + ) + + // A re-sync that keeps Alpha but drops Beta preserves Alpha's folded state and removes Beta. + val createdAgain = registry.sync(listOf(entryFor(alpha, "Alpha")), owner.pubKey) + assertTrue(createdAgain.isEmpty()) + assertNotNull(registry.sessionFor(alpha.communityIdHex)) + assertNull(registry.sessionFor(beta.communityIdHex)) + assertEquals( + "Alpha", + registry + .sessionFor(alpha.communityIdHex)!! + .state.value + ?.metadata + ?.name, + ) + + // A wrap from an unknown community is routed nowhere. + val gamma = ConcordCommunityFactory.create(owner, "Gamma", createdAt = 1L, relays = listOf("wss://r.example")) + assertFalse(registry.ingest(gamma.genesisWraps.first())) + } +}