commit 8f5156a06cf6c84ea3c85c95e096bbfc6800be61
parent a6c7eb3daee5d7ca86386a1dba0ba57434085de1
Author: triesap <tyson@radroots.org>
Date: Mon, 10 Aug 2026 17:22:19 +0000
runtime: make snapshot delivery gap-aware
- conflate native callback bursts while checking every send result
- terminate unexpected delivery failures with typed application problems
- resnapshot on predecessor gaps and reject insufficient refreshes
- cover burst convergence, stale changes, gaps, failures, and cleanup
Diffstat:
4 files changed, 198 insertions(+), 9 deletions(-)
diff --git a/app/desktop/src/main/kotlin/org/harvestcircle/application/NativeHarvestCircleRuntime.kt b/app/desktop/src/main/kotlin/org/harvestcircle/application/NativeHarvestCircleRuntime.kt
@@ -1,9 +1,11 @@
package org.harvestcircle.application
import kotlinx.coroutines.CancellationException
+import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.NonCancellable
-import kotlinx.coroutines.awaitCancellation
+import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.flow.Flow
+import kotlinx.coroutines.flow.buffer
import kotlinx.coroutines.flow.channelFlow
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
@@ -22,6 +24,7 @@ import org.harvestcircle.ffi.ShutdownReceiptDto
import org.harvestcircle.ffi.SnapshotChangeDto
import org.harvestcircle.ffi.compatibilityDescriptor
import org.harvestcircle.ffi.generateOperationIdV7
+import java.util.concurrent.atomic.AtomicBoolean
import org.harvestcircle.ffi.buildInfo as nativeBuildInfo
class NativeHarvestCircleRuntime internal constructor(
@@ -50,18 +53,31 @@ class NativeHarvestCircleRuntime internal constructor(
override fun changes(): Flow<ApplicationChange> =
channelFlow {
+ val acceptingCallbacks = AtomicBoolean(true)
+ val deliveryFailure = CompletableDeferred<ApplicationFailure>()
val subscription =
callNative {
native.subscribe { change ->
- trySend(change.toApplicationChange())
+ if (!acceptingCallbacks.get()) return@subscribe
+ val mapped =
+ try {
+ change.toApplicationChange()
+ } catch (_: Exception) {
+ deliveryFailure.complete(observerDeliveryFailure())
+ return@subscribe
+ }
+ if (trySend(mapped).isFailure && acceptingCallbacks.get()) {
+ deliveryFailure.complete(observerDeliveryFailure())
+ }
}
}
try {
- awaitCancellation()
+ throw deliveryFailure.await()
} finally {
+ acceptingCallbacks.set(false)
withContext(NonCancellable) { subscription.unsubscribe() }
}
- }
+ }.buffer(Channel.CONFLATED)
override suspend fun execute(command: ApplicationCommand): ApplicationCommandResult =
when (command) {
@@ -236,6 +252,18 @@ class NativeHarvestCircleRuntime internal constructor(
),
)
+ private fun observerDeliveryFailure(): ApplicationFailure =
+ ApplicationFailure(
+ ApplicationProblem(
+ code = ApplicationErrorCode.ObserverRegistrationFailed,
+ category = ApplicationErrorCategory.Lifecycle,
+ retryable = false,
+ recoveryAction = RecoveryAction.RestartApplication,
+ operationId = null,
+ safeMessage = "Application updates could not be delivered.",
+ ),
+ )
+
private suspend fun <T> callNative(
fallbackOperationId: OperationId? = null,
operation: suspend () -> T,
diff --git a/app/desktop/src/test/kotlin/org/harvestcircle/application/NativeRuntimeMappingsTest.kt b/app/desktop/src/test/kotlin/org/harvestcircle/application/NativeRuntimeMappingsTest.kt
@@ -1,7 +1,14 @@
+@file:OptIn(kotlinx.coroutines.ExperimentalCoroutinesApi::class)
+
package org.harvestcircle.application
+import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.async
+import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.flow.first
+import kotlinx.coroutines.flow.onEach
+import kotlinx.coroutines.flow.take
+import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
import org.harvestcircle.ffi.ActiveIdentityDto
@@ -253,6 +260,60 @@ class NativeHarvestCircleRuntimeTest {
assertEquals(SnapshotRevision(2UL), change.snapshot.revision)
runCurrent()
assertTrue(port.subscriptionClosed)
+ assertEquals(1, port.subscriptionCloseCalls)
+ }
+
+ @Test
+ fun observerChangesConflateBurstsToTheLatestSnapshot() =
+ runTest {
+ val port = FakeNativeCorePort()
+ val runtime = NativeHarvestCircleRuntime(port)
+ val releaseCollector = CompletableDeferred<Unit>()
+ val pending =
+ async {
+ runtime
+ .changes()
+ .onEach { change ->
+ if (change.snapshot.revision == SnapshotRevision(2UL)) releaseCollector.await()
+ }.take(2)
+ .toList()
+ }
+ runCurrent()
+ port.emit(SnapshotChangeDto(populatedSnapshot(2UL), 1UL))
+ runCurrent()
+
+ (3UL..100UL).forEach { revision ->
+ port.emit(SnapshotChangeDto(populatedSnapshot(revision), revision - 1UL))
+ }
+ releaseCollector.complete(Unit)
+ val changes = pending.await()
+
+ assertEquals(listOf(2UL, 100UL), changes.map { it.snapshot.revision.value })
+ assertEquals(1, port.subscriptionCloseCalls)
+ }
+
+ @Test
+ fun invalidObserverDeliveryFailsTypedAndUnsubscribesOnce() =
+ runTest {
+ val port = FakeNativeCorePort()
+ val runtime = NativeHarvestCircleRuntime(port)
+ val pending =
+ async {
+ try {
+ runtime.changes().collect { }
+ null
+ } catch (failure: ApplicationFailure) {
+ failure
+ }
+ }
+ runCurrent()
+
+ port.emit(SnapshotChangeDto(populatedSnapshot(2UL), 2UL))
+ val failure = assertIs<ApplicationFailure>(pending.await())
+
+ assertEquals(ApplicationErrorCode.ObserverRegistrationFailed, failure.problem.code)
+ assertEquals(RecoveryAction.RestartApplication, failure.problem.recoveryAction)
+ assertEquals(1, port.subscriptionCloseCalls)
}
}
@@ -263,6 +324,7 @@ private class FakeNativeCorePort : NativeCorePort {
var shutdownCalls = 0
var closed = false
var subscriptionClosed = false
+ var subscriptionCloseCalls = 0
private var observer: ((SnapshotChangeDto) -> Unit)? = null
private val snapshot = populatedSnapshot(2UL)
@@ -272,7 +334,10 @@ private class FakeNativeCorePort : NativeCorePort {
override suspend fun subscribe(onChange: (SnapshotChangeDto) -> Unit): NativeSubscriptionHandle {
observer = onChange
- return NativeSubscriptionHandle { subscriptionClosed = true }
+ return NativeSubscriptionHandle {
+ subscriptionCloseCalls += 1
+ subscriptionClosed = true
+ }
}
fun emit(change: SnapshotChangeDto) {
diff --git a/app/shared/src/commonMain/kotlin/org/harvestcircle/application/HarvestCirclePresenter.kt b/app/shared/src/commonMain/kotlin/org/harvestcircle/application/HarvestCirclePresenter.kt
@@ -40,7 +40,7 @@ class HarvestCirclePresenter(
subscriptionJob =
scope.launch {
try {
- runtime.changes().collect { change -> acceptSnapshot(change.snapshot) }
+ runtime.changes().collect(::acceptChange)
} catch (error: CancellationException) {
throw error
} catch (error: Exception) {
@@ -367,6 +367,29 @@ class HarvestCirclePresenter(
}
}
+ private fun acceptChange(change: ApplicationChange) {
+ val acceptedRevision = state.value.snapshot.revision
+ if (change.snapshot.revision.value <= acceptedRevision.value) return
+ if (change.previousRevision == acceptedRevision) {
+ acceptSnapshot(change.snapshot)
+ return
+ }
+ val refreshed = runtime.currentSnapshot()
+ if (refreshed.revision.value < change.snapshot.revision.value) {
+ throw ApplicationFailure(
+ ApplicationProblem(
+ code = ApplicationErrorCode.ObserverRegistrationFailed,
+ category = ApplicationErrorCategory.Lifecycle,
+ retryable = false,
+ recoveryAction = RecoveryAction.RestartApplication,
+ operationId = null,
+ safeMessage = "Application updates could not be synchronized.",
+ ),
+ )
+ }
+ acceptSnapshot(refreshed)
+ }
+
private fun acceptFailure(
error: Throwable,
operationId: OperationId?,
diff --git a/app/shared/src/commonTest/kotlin/org/harvestcircle/application/HarvestCirclePresenterTest.kt b/app/shared/src/commonTest/kotlin/org/harvestcircle/application/HarvestCirclePresenterTest.kt
@@ -28,7 +28,7 @@ class HarvestCirclePresenterTest {
assertEquals(1UL, presenter.state.value.snapshot.revision.value)
runtime.emit(snapshot(3UL))
- runtime.emit(snapshot(2UL))
+ runtime.emitChange(snapshot(2UL), SnapshotRevision(1UL))
runCurrent()
assertEquals(3UL, presenter.state.value.snapshot.revision.value)
@@ -37,6 +37,63 @@ class HarvestCirclePresenterTest {
}
@Test
+ fun revisionGapResnapshotsAndAcceptsOnlyTheAuthoritativeLatestState() =
+ runTest {
+ val runtime = FakePresenterRuntime()
+ val presenter = presenter(runtime)
+ runCurrent()
+ val callsBeforeGap = runtime.currentSnapshotCalls
+ runtime.setCurrent(snapshot(5UL))
+
+ runtime.emitChange(snapshot(4UL), SnapshotRevision(2UL))
+ runCurrent()
+
+ assertEquals(5UL, presenter.state.value.snapshot.revision.value)
+ assertEquals(callsBeforeGap + 1, runtime.currentSnapshotCalls)
+ presenter.close()
+ }
+
+ @Test
+ fun duplicateAndStaleChangesAreIgnoredWithoutResnapshotting() =
+ runTest {
+ val runtime = FakePresenterRuntime()
+ val presenter = presenter(runtime)
+ runCurrent()
+ runtime.emit(snapshot(2UL))
+ runCurrent()
+ val callsBeforeStale = runtime.currentSnapshotCalls
+
+ runtime.emitChange(snapshot(2UL), SnapshotRevision(1UL))
+ runtime.emitChange(snapshot(1UL), null)
+ runCurrent()
+
+ assertEquals(2UL, presenter.state.value.snapshot.revision.value)
+ assertEquals(callsBeforeStale, runtime.currentSnapshotCalls)
+ presenter.close()
+ }
+
+ @Test
+ fun insufficientGapResnapshotSurfacesATypedTerminalProblem() =
+ runTest {
+ val runtime = FakePresenterRuntime()
+ val presenter = presenter(runtime)
+ runCurrent()
+ runtime.setCurrent(snapshot(3UL))
+
+ runtime.emitChange(snapshot(5UL), SnapshotRevision(3UL))
+ runCurrent()
+
+ assertEquals(1UL, presenter.state.value.snapshot.revision.value)
+ assertEquals(
+ ApplicationErrorCode.ObserverRegistrationFailed,
+ presenter.state.value.lastProblem
+ ?.code,
+ )
+ assertEquals(CommandStatus.FAILED_TERMINAL, presenter.state.value.commandStatus)
+ presenter.close()
+ }
+
+ @Test
fun busyAdmissionRejectsOverlappingCommands() =
runTest {
val bootstrapGate = CompletableDeferred<Unit>()
@@ -204,6 +261,7 @@ private class FakePresenterRuntime(
private val changes = MutableSharedFlow<ApplicationChange>(extraBufferCapacity = 8)
private var current = snapshot(0UL)
+ var currentSnapshotCalls = 0
var nextFailure: ApplicationProblem? = null
var executeCalls = 0
var executeCancelled = false
@@ -218,7 +276,10 @@ private class FakePresenterRuntime(
return snapshot(1UL).also { current = it }
}
- override fun currentSnapshot(): ApplicationSnapshot = current
+ override fun currentSnapshot(): ApplicationSnapshot {
+ currentSnapshotCalls += 1
+ return current
+ }
override fun changes(): Flow<ApplicationChange> = changes
@@ -271,8 +332,20 @@ private class FakePresenterRuntime(
}
fun emit(snapshot: ApplicationSnapshot) {
+ val previousRevision = current.revision
+ current = snapshot
+ emitChange(snapshot, previousRevision)
+ }
+
+ fun emitChange(
+ snapshot: ApplicationSnapshot,
+ previousRevision: SnapshotRevision?,
+ ) {
+ check(changes.tryEmit(ApplicationChange(snapshot, previousRevision)))
+ }
+
+ fun setCurrent(snapshot: ApplicationSnapshot) {
current = snapshot
- changes.tryEmit(ApplicationChange(snapshot, null))
}
}