diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt index 5e3a2e6300..32fed3bc1b 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt @@ -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() } } } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnSyncLoop.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnSyncLoop.kt index ee53f0e717..29c7313912 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnSyncLoop.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnSyncLoop.kt @@ -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 } diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt index fe15fdf48c..58fc3bfbd9 100644 --- a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManagerTest.kt @@ -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, + 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)