mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix(concord): show a joined community's channels without a restart
Joining fetched the Control Plane before the community's session existed. Those wraps reached LocalCache as unclaimed notes, the live subscription's copies of the same wraps were deduplicated away, and the community sat on "No channels yet" until a restart emptied the cache. The join now waits for the session and hands it the wraps it fetched (creation hands over its genesis the same way). The live plane subscription also applied each relay's since cursor, set by the other communities on it, to a community whose Control Plane had not folded yet, skipping its older editions. Unfolded communities are now requested in full. The Community List load a join must finish before writing (a ~30s drain of every stock and outbox relay) now starts when the join starts and is single flight, so it overlaps the bundle and plane fetches. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
99513a92b6
commit
89a8e658dc
+12
-6
@@ -258,16 +258,22 @@ object ConcordSubscriptionPlanner {
|
||||
accountPubKey: HexKey? = null,
|
||||
stateOf: (ConcordCommunityListEntry) -> ConcordCommunityState?,
|
||||
): List<RelayBasedFilter> {
|
||||
val controlSubs = controlPlaneSubs(entries)
|
||||
// A community whose Control Plane hasn't folded yet (just joined, or its editions never
|
||||
// arrived) is fetched without the relay's `since`: that cursor was set by the OTHER
|
||||
// communities on the relay, and applying it here skips every edition older than it — the
|
||||
// genesis included — so the community showed no channels until a restart dropped the cursor.
|
||||
val (folded, unfolded) = entries.partition { stateOf(it) != null }
|
||||
|
||||
val otherSubs = ArrayList<ConcordPlaneSub>()
|
||||
otherSubs += auxiliaryPlaneSubs(entries)
|
||||
for (entry in entries) {
|
||||
val state = stateOf(entry) ?: continue
|
||||
otherSubs += channelPlaneSubs(entry, state)
|
||||
otherSubs += auxiliaryPlaneSubs(folded)
|
||||
for (entry in folded) {
|
||||
otherSubs += channelPlaneSubs(entry, stateOf(entry) ?: continue)
|
||||
}
|
||||
|
||||
return relayBasedFilters(controlSubs, since, accountPubKey).orEmpty() + relayBasedFilters(otherSubs, since, accountPubKey).orEmpty()
|
||||
return relayBasedFilters(controlPlaneSubs(folded), since, accountPubKey).orEmpty() +
|
||||
relayBasedFilters(otherSubs, since, accountPubKey).orEmpty() +
|
||||
relayBasedFilters(controlPlaneSubs(unfolded), null, accountPubKey).orEmpty() +
|
||||
relayBasedFilters(auxiliaryPlaneSubs(unfolded), null, accountPubKey).orEmpty()
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+40
-5
@@ -92,10 +92,18 @@ import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
|
||||
/** Name of the default Concord community Admin role minted by "Make admin". */
|
||||
private const val CONCORD_ADMIN_ROLE = "Admin"
|
||||
|
||||
/** How long a join waits for the new community's session before handing it the wraps it fetched. */
|
||||
private const val SESSION_WAIT_MS = 10_000L
|
||||
|
||||
/**
|
||||
* How often a joined Concord community's stored invite link is re-resolved to check whether
|
||||
* we were left out of a Refounding (see `recoverStrandedConcordCommunities`). Stranding is
|
||||
@@ -137,27 +145,48 @@ class AccountConcordActions(
|
||||
entry: ConcordCommunityListEntry,
|
||||
inviteCreator: HexKey? = null,
|
||||
inviteLabel: String? = null,
|
||||
fetchedWraps: List<Event> = emptyList(),
|
||||
) {
|
||||
if (!persistConcordEntry(entry)) return
|
||||
// The session is built asynchronously from the Community List flow. Every wrap that reaches
|
||||
// the cache before it exists is kept as an unclaimed note, and the live subscription's copy
|
||||
// of the same wrap is then deduplicated away, so the community showed "No channels yet"
|
||||
// until a restart emptied the cache. Wait for the session, then hand it what we fetched.
|
||||
awaitConcordSession(entry.id)
|
||||
fetchedWraps.forEach { account.concordSessions.ingest(it) }
|
||||
announceConcordGuestbookJoin(entry, inviteCreator, inviteLabel)
|
||||
}
|
||||
|
||||
private suspend fun awaitConcordSession(communityId: HexKey) {
|
||||
withTimeoutOrNull(SESSION_WAIT_MS) {
|
||||
account.concordSessions.revision.first { account.concordSessions.sessionFor(communityId) != null }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Makes the Community List fetched before any write can depend on it: an empty fragment set
|
||||
* only means "no List" once the relays have been asked (CORD-02 §8 — a write built on an
|
||||
* unloaded List replaces fragments another device published).
|
||||
*/
|
||||
private suspend fun ensureConcordListLoaded() {
|
||||
if (!account.concordChannelList.relaysConfirmed) importConcordCommunities()
|
||||
suspend fun preloadConcordList() {
|
||||
if (account.concordChannelList.relaysConfirmed) return
|
||||
// Single-flight: the import is a ~30s drain of every stock/outbox relay. A join starts it
|
||||
// as soon as it begins, and the join's own List write then waits for that same fetch
|
||||
// instead of starting a second one.
|
||||
concordListImport.withLock {
|
||||
if (!account.concordChannelList.relaysConfirmed) importConcordCommunities()
|
||||
}
|
||||
}
|
||||
|
||||
private val concordListImport = Mutex()
|
||||
|
||||
/**
|
||||
* Read-modify-writes the Community List through [change] and publishes the fragments it
|
||||
* produced. Returns false — logged, never thrown into a UI coroutine — when the List can't be
|
||||
* written safely yet (fragments unreadable or not loaded) or a fragment would pass the ceiling.
|
||||
*/
|
||||
private suspend fun writeConcordList(change: suspend (ConcordChannelListState) -> List<Event>): Boolean {
|
||||
ensureConcordListLoaded()
|
||||
preloadConcordList()
|
||||
return try {
|
||||
account.sendMyPublicAndPrivateOutbox(change(account.concordChannelList))
|
||||
true
|
||||
@@ -227,6 +256,7 @@ class AccountConcordActions(
|
||||
name = name,
|
||||
addedAt = TimeUtils.nowMillis(),
|
||||
),
|
||||
fetchedWraps = community.genesisWraps,
|
||||
)
|
||||
return community.communityIdHex
|
||||
}
|
||||
@@ -540,6 +570,11 @@ class AccountConcordActions(
|
||||
if (!account.isWriteable()) return ConcordInviteResult.InvalidLink
|
||||
val parsed = ConcordActions.parseInviteLink(url) ?: return ConcordInviteResult.InvalidLink
|
||||
|
||||
// Joining ends in a Community List write, which first needs the List loaded (a slow drain of
|
||||
// every stock and outbox relay). Start it now so it overlaps the bundle and plane fetches
|
||||
// below instead of running after them.
|
||||
account.scope.launch { preloadConcordList() }
|
||||
|
||||
val relays =
|
||||
(parsed.fragment.relays.mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) } + account.outboxRelays.flow.value).toSet()
|
||||
if (relays.isEmpty()) return ConcordInviteResult.NotReachable
|
||||
@@ -628,7 +663,7 @@ class AccountConcordActions(
|
||||
if (rejoined != null) {
|
||||
if (!adoptedConcordRotations.add("${rejoined.id}:${rejoined.rootEpoch}")) return ConcordInviteResult.Joined(bundle.communityId)
|
||||
Log.i("Concord") { "Stranded rejoin by explicit invite: ${rejoined.id} -> epoch ${rejoined.rootEpoch}" }
|
||||
joinConcordCommunity(rejoined, inviteCreator, inviteLabel)
|
||||
joinConcordCommunity(rejoined, inviteCreator, inviteLabel, planeWraps)
|
||||
_strandedConcordCommunities.value -= rejoined.id
|
||||
return ConcordInviteResult.Joined(bundle.communityId)
|
||||
}
|
||||
@@ -654,7 +689,7 @@ class AccountConcordActions(
|
||||
// recoverStrandedConcordCommunities().
|
||||
inviteRef = ConcordActions.bareInviteRef(url),
|
||||
)
|
||||
joinConcordCommunity(entry, inviteCreator, inviteLabel)
|
||||
joinConcordCommunity(entry, inviteCreator, inviteLabel, planeWraps)
|
||||
return ConcordInviteResult.Joined(bundle.communityId)
|
||||
}
|
||||
|
||||
|
||||
+42
@@ -22,6 +22,7 @@ package com.vitorpamplona.amethyst.commons.actions
|
||||
|
||||
import com.vitorpamplona.amethyst.commons.relays.MutableTime
|
||||
import com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityFactory
|
||||
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.relay.normalizer.RelayUrlNormalizer
|
||||
@@ -243,4 +244,45 @@ class ConcordSubscriptionPlannerTest {
|
||||
assertTrue(generalPk in otherFilter.authors.orEmpty(), "channel plane missing from the non-control filter")
|
||||
assertTrue(controlPk !in otherFilter.authors.orEmpty(), "Control Plane leaked into the Guestbook/channel filter")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aJustJoinedCommunityIsFetchedWithoutTheRelaysSinceCursor() =
|
||||
runTest {
|
||||
// Regression: joining a community on a relay that already carries another one asked for its
|
||||
// Control Plane `since` that relay's EOSE cursor, so the genesis editions (older than the
|
||||
// cursor) never arrived and the community showed "No channels yet" until an app restart.
|
||||
fun entryOf(c: NewConcordCommunity) =
|
||||
com.vitorpamplona.quartz.concord.cord02Community.ConcordCommunityListEntry(
|
||||
id = c.communityIdHex,
|
||||
owner = c.ownerPubKey,
|
||||
ownerSalt = c.ownerSalt.toHexKey(),
|
||||
root = c.communityRoot.toHexKey(),
|
||||
rootEpoch = c.rootEpoch,
|
||||
controlPk = c.controlPkHex,
|
||||
relays = listOf("wss://r.example"),
|
||||
name = "x",
|
||||
)
|
||||
val folded = ConcordCommunityFactory.create(owner, "Folded", createdAt = 1L, relays = listOf("wss://r.example"))
|
||||
val joined = ConcordCommunityFactory.create(owner, "Joined", createdAt = 1L, relays = listOf("wss://r.example"))
|
||||
val foldedEntry = entryOf(folded)
|
||||
val joinedEntry = entryOf(joined)
|
||||
val foldedState = ConcordActions.foldCommunity(folded.genesisWraps, folded.controlPlane, folded.communityId, folded.ownerPubKey)
|
||||
val relay = RelayUrlNormalizer.normalizeOrNull("wss://r.example")!!
|
||||
|
||||
val filters =
|
||||
ConcordSubscriptionPlanner.controlIsolatedFilters(
|
||||
listOf(foldedEntry, joinedEntry),
|
||||
since = mutableMapOf(relay to MutableTime(1234L)),
|
||||
stateOf = { if (it.id == foldedEntry.id) foldedState else null },
|
||||
)
|
||||
|
||||
val joinedControl = filters.map { it.filter }.single { joined.controlPlane.address in it.authors.orEmpty() }
|
||||
assertNull(joinedControl.since, "an unfolded community's Control Plane must be fetched in full")
|
||||
val joinedGuestbook = ConcordActions.guestbookPlane(joined.communityRoot, joined.communityId, joined.rootEpoch).publicKeyHex
|
||||
assertNull(filters.map { it.filter }.single { joinedGuestbook in it.authors.orEmpty() }.since)
|
||||
|
||||
// The already-folded community keeps its incremental cursor.
|
||||
val foldedControl = filters.map { it.filter }.single { folded.controlPlane.address in it.authors.orEmpty() }
|
||||
assertEquals(1234L, foldedControl.since)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user