mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix(cordn): a cancelled subscription really does save, and stop() waits
Two halves of one bug, the second caused by fixing the first. `CordnGroupManager.subscribe` persisted in a `finally` and its comment claimed that covered cancellation. It did not: `persistAll` suspends, and a suspend call in a cancelled coroutine throws before it writes anything. Cancellation is also the normal way this ends — CordnSyncLoop cancels the subscription whenever the group set changes — so the case the `finally` was written for was the one case it never served. `withContext(NonCancellable)` fixes it. That alone broke two CordnRuntime tests, and rightly. `CordnSyncLoop.stop()` was `job?.cancel()` with no join, so it returned while the loop's `finally` blocks were still running. `importArchive` stops the loops, deletes the account's directory and restores from the archive — and a persist that landed a moment late wrote a group back onto disk after the delete, which the restore then read in again. That is exactly the merge importArchive exists to prevent, and `adoptMigration` had it too. stop() is suspend now and does `cancelAndJoin`; CordnRuntime.stop() already suspended, so nothing else had to change. Nobody had noticed the missing join because the persist it races was silently failing. Making the write survive cancellation is what made the race reachable. The test needed two goes to be worth anything. The first version passed against the mutant: `send()` already persists, so the cursor it asserted on was in the store before the subscription ran, and the fake returned normally so no cancellation ever landed mid-subscription. The second still passed, because InMemoryCordnGroupStore's methods are `suspend` but never reach a suspension point — and cancellation is only observed at one. It takes a store whose writes actually suspend, like the file-backed one, and a coordinator that hangs after delivering. With both, removing NonCancellable fails the test. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
+9
-1
@@ -46,9 +46,11 @@ import com.vitorpamplona.quartz.mls.group.MlsGroupState
|
||||
import com.vitorpamplona.quartz.mls.messages.KeyPackageBundle
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||
import kotlinx.coroutines.NonCancellable
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlin.io.encoding.Base64
|
||||
import kotlin.io.encoding.ExperimentalEncodingApi
|
||||
|
||||
@@ -612,7 +614,13 @@ class CordnGroupManager(
|
||||
} finally {
|
||||
// finally, not after: a subscription normally ends by timing out or
|
||||
// being cancelled, and those are the cases with progress to keep.
|
||||
persistAll()
|
||||
//
|
||||
// NonCancellable because the cancelled case is the one this exists
|
||||
// for and the one it could not serve: persistAll suspends, and a
|
||||
// suspend call in a cancelled coroutine throws before it writes
|
||||
// anything. Without this the `finally` looked like it saved on
|
||||
// cancellation and silently did not.
|
||||
withContext(NonCancellable) { persistAll() }
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+18
-3
@@ -25,6 +25,7 @@ import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.TimeoutCancellationException
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.cancelAndJoin
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
@@ -153,9 +154,23 @@ class CordnSyncLoop(
|
||||
job = scope.launch { run() }
|
||||
}
|
||||
|
||||
/** Stops the loop. Safe to call when it is not running. */
|
||||
fun stop() {
|
||||
job?.cancel()
|
||||
/**
|
||||
* Stops the loop and waits for it to actually be stopped. Safe to call
|
||||
* when it is not running.
|
||||
*
|
||||
* `cancel()` alone returns while the loop's `finally` blocks are still
|
||||
* running, and one of those persists the group state a subscription
|
||||
* ingested. `CordnRuntime.importArchive` stops the loops, deletes the
|
||||
* account's directory and restores from the archive — so a persist that
|
||||
* landed a moment late wrote a group back onto disk after the delete, and
|
||||
* the restore then read it in again. That is the merge importArchive
|
||||
* exists to prevent.
|
||||
*
|
||||
* Nothing was noticed while the persist silently failed on cancellation;
|
||||
* making it survive cancellation is what made the missing join matter.
|
||||
*/
|
||||
suspend fun stop() {
|
||||
job?.cancelAndJoin()
|
||||
job = null
|
||||
_state.value = State.Stopped
|
||||
}
|
||||
|
||||
+97
@@ -22,7 +22,10 @@ package com.vitorpamplona.amethyst.commons.cordn
|
||||
|
||||
import com.vitorpamplona.quartz.cordn.groups.CordnCredential
|
||||
import com.vitorpamplona.quartz.cordn.groups.CordnGroupPolicy
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.GroupMessage
|
||||
import com.vitorpamplona.quartz.cordn.spec00Coordinator.ICoordinator
|
||||
import com.vitorpamplona.quartz.cordn.spec01GroupMetadata.CordnGroupMetadata
|
||||
import com.vitorpamplona.quartz.cordn.sync.GroupCursor
|
||||
import com.vitorpamplona.quartz.mls.group.MlsGroup
|
||||
import com.vitorpamplona.quartz.mls.messages.KeyPackageBundle
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
@@ -31,13 +34,19 @@ import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.EventTemplate
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerSync
|
||||
import kotlinx.coroutines.awaitCancellation
|
||||
import kotlinx.coroutines.cancelAndJoin
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.advanceUntilIdle
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import kotlinx.coroutines.yield
|
||||
import kotlin.io.encoding.Base64
|
||||
import kotlin.io.encoding.ExperimentalEncodingApi
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertNotEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
@@ -601,6 +610,94 @@ class CordnGroupManagerTest {
|
||||
assertTrue(aliceManager.group(gid)!!.memberCount == 1)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a cancelled subscription still saves what it ingested`() =
|
||||
runTest {
|
||||
// The case the `finally` exists for, and the one it could not
|
||||
// serve: persistAll suspends, and a suspend call inside a
|
||||
// cancelled coroutine throws before writing anything. Cancelling
|
||||
// is also the normal way this ends — CordnSyncLoop cancels the
|
||||
// subscription whenever the group set changes.
|
||||
val aliceCoordinator = FakeCoordinator(callerPubKey = alice)
|
||||
val aliceManager = manager(alice, aliceCoordinator)
|
||||
val (bundle, stored) = bobsPublication()
|
||||
aliceCoordinator.keyPackages[stored.keyPackageRef] = stored
|
||||
aliceManager.createGroup(gid, CordnGroupMetadata(name = "Cancelled"))
|
||||
aliceManager.invite(gid, bob, stored.keyPackageRef)
|
||||
aliceManager.send(gid, "a message bob has not seen")
|
||||
|
||||
// Bob joins with an empty store, so any cursor in it afterwards
|
||||
// can only have come from the subscription.
|
||||
val bobStore = SuspendingStore()
|
||||
val bobManager =
|
||||
CordnGroupManager(
|
||||
accountPubKey = bob,
|
||||
config = config,
|
||||
coordinator = HangingCoordinator(bobsCoordinatorOver(aliceCoordinator)),
|
||||
store = bobStore,
|
||||
clock = { 1_757_000_000L },
|
||||
)
|
||||
bobManager.pendingWelcomes({ bundle }).pending.forEach { bobManager.accept(it) }
|
||||
val before = bobStore.loadCursor(gid)
|
||||
|
||||
// Delivers, then hangs — so the cancellation lands while the
|
||||
// subscription is open, which is the whole point. A fake that
|
||||
// returns normally leaves the `finally` running in a live
|
||||
// coroutine and proves nothing.
|
||||
val job = launch { bobManager.subscribe(timeoutMs = 60_000) { } }
|
||||
advanceUntilIdle()
|
||||
job.cancelAndJoin()
|
||||
|
||||
assertNotEquals(
|
||||
before,
|
||||
bobStore.loadCursor(gid),
|
||||
"a cursor the subscription advanced must reach the store even when cancelled",
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* A store whose writes actually suspend, like the file-backed one.
|
||||
*
|
||||
* [InMemoryCordnGroupStore]'s methods are `suspend` but never reach a
|
||||
* suspension point, and cancellation is only observed at one — so against
|
||||
* it a cancelled `finally` writes happily and proves nothing. The real
|
||||
* store goes through `Dispatchers.IO`, where the same code throws before
|
||||
* it writes. One `yield()` is the difference.
|
||||
*/
|
||||
private class SuspendingStore(
|
||||
private val inner: CordnGroupStore = InMemoryCordnGroupStore(),
|
||||
) : CordnGroupStore by inner {
|
||||
override suspend fun saveCursor(
|
||||
gid: String,
|
||||
cursor: GroupCursor,
|
||||
) {
|
||||
yield()
|
||||
inner.saveCursor(gid, cursor)
|
||||
}
|
||||
|
||||
override suspend fun saveGroup(
|
||||
gid: String,
|
||||
state: ByteArray,
|
||||
) {
|
||||
yield()
|
||||
inner.saveGroup(gid, state)
|
||||
}
|
||||
}
|
||||
|
||||
/** Delivers whatever it is wrapping, then never returns. */
|
||||
private class HangingCoordinator(
|
||||
private val inner: FakeCoordinator,
|
||||
) : ICoordinator by inner {
|
||||
override suspend fun subscribeMessages(
|
||||
cursors: Map<String, Long?>,
|
||||
timeoutMs: Long,
|
||||
onMessage: (GroupMessage) -> Unit,
|
||||
) {
|
||||
inner.subscribeMessages(cursors, timeoutMs, onMessage)
|
||||
awaitCancellation()
|
||||
}
|
||||
}
|
||||
|
||||
/** Bob's own view of the coordinator, carrying whatever Alice's has stored. */
|
||||
private fun bobsCoordinatorOver(alices: FakeCoordinator): FakeCoordinator {
|
||||
val bobCoordinator = FakeCoordinator(callerPubKey = bob)
|
||||
|
||||
Reference in New Issue
Block a user