commit 3e2b0e55a26ca3e3f6d71619087770cd05b5297c
parent 3d2cd67529e8c074febac78e2620dd792cdc0a3a
Author: triesap <tyson@radroots.org>
Date: Fri, 2 Oct 2026 06:12:36 +0000
runtime: carry overflow recovery across FFI and Kotlin
- Preserve initial delivery before the coalesced latest Kotlin snapshot.
- Verify final-tail recovery through native callbacks and generated JNI.
- Preserve predecessor metadata and typed safe observer failures.
- Retain deterministic regressions with governed Rust and desktop checks.
Diffstat:
4 files changed, 327 insertions(+), 8 deletions(-)
diff --git a/app/desktop/src/integrationTest/kotlin/org/harvestcircle/integration/NativeRuntimeIntegrationTest.kt b/app/desktop/src/integrationTest/kotlin/org/harvestcircle/integration/NativeRuntimeIntegrationTest.kt
@@ -26,6 +26,64 @@ import kotlin.test.fail
class NativeRuntimeIntegrationTest {
@Test
+ fun generatedActorObserverRecoversFinalTailAfterSaturation() {
+ val dataRoot = Files.createTempDirectory("harvestcircle-generated-observer-")
+ try {
+ val bridge = HarvestCircleTestBridge.open(dataRoot.toString())
+ try {
+ bridge.bootstrap()
+ val publicKeys =
+ listOf(
+ "00000000-0000-7000-8000-000000000031",
+ "00000000-0000-7000-8000-000000000032",
+ ).map { requestId ->
+ val request = bridge.beginGeneratedIdentity()
+ try {
+ val publicKey = request.identity().publicKeyHex
+ bridge.acknowledgeGeneratedIdentity(requestId, bridge.snapshot().revision, 2_000UL, request)
+ publicKey
+ } finally {
+ request.close()
+ }
+ }
+ assertEquals(2, publicKeys.distinct().size)
+ val initial = bridge.selectIdentity(publicKeys[1])
+ bridge.startObserver()
+ assertEquals(initial, assertNotNull(bridge.nextObservedSnapshot(2_000UL)))
+
+ val queuedCapacity = 16
+ val publications = queuedCapacity + 2
+ repeat(publications) { offset ->
+ val publicKey = publicKeys[offset % publicKeys.size]
+ val snapshot = bridge.selectIdentity(publicKey)
+ assertEquals(initial.revision + offset.toULong() + 1UL, snapshot.revision)
+ assertEquals(publicKey, snapshot.selectedPublicKeyHex)
+ }
+ val finalRevision = initial.revision + publications.toULong()
+ assertEquals(finalRevision, bridge.snapshot().revision)
+
+ repeat(queuedCapacity) { offset ->
+ val queued = assertNotNull(bridge.nextObservedSnapshot(2_000UL))
+ assertEquals(initial.revision + offset.toULong() + 1UL, queued.revision)
+ assertEquals(publicKeys[offset % publicKeys.size], queued.selectedPublicKeyHex)
+ }
+ val finalTail = assertNotNull(bridge.nextObservedSnapshot(2_000UL))
+ assertEquals(finalRevision, finalTail.revision)
+ assertEquals(publicKeys[(publications - 1) % publicKeys.size], finalTail.selectedPublicKeyHex)
+ assertEquals(2, finalTail.identities.size)
+ } finally {
+ try {
+ bridge.shutdown()
+ } finally {
+ bridge.close()
+ }
+ }
+ } finally {
+ deleteTree(dataRoot)
+ }
+ }
+
+ @Test
fun generatedFfiClassifiesCanonicalReferencesAndRedactsPrivateKeys() {
val eventId = "d94a3f4dd87b9a3b0bed183b32e916fa29c8020107845d1752d72697fe5309a5"
val event = classifyNostrReference(eventId)
diff --git a/app/desktop/src/main/kotlin/org/harvestcircle/application/NativeHarvestCircleRuntime.kt b/app/desktop/src/main/kotlin/org/harvestcircle/application/NativeHarvestCircleRuntime.kt
@@ -5,8 +5,8 @@ import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.flow.Flow
-import kotlinx.coroutines.flow.buffer
-import kotlinx.coroutines.flow.channelFlow
+import kotlinx.coroutines.flow.flow
+import kotlinx.coroutines.selects.select
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
@@ -55,9 +55,11 @@ class NativeHarvestCircleRuntime internal constructor(
}
override fun changes(): Flow<ApplicationChange> =
- channelFlow {
+ flow {
val acceptingCallbacks = AtomicBoolean(true)
val deliveryFailure = CompletableDeferred<ApplicationFailure>()
+ val initialChange = CompletableDeferred<ApplicationChange>()
+ val latestChanges = Channel<ApplicationChange>(Channel.CONFLATED)
val subscription =
callNative {
native.subscribe { change ->
@@ -69,18 +71,33 @@ class NativeHarvestCircleRuntime internal constructor(
deliveryFailure.complete(observerDeliveryFailure())
return@subscribe
}
- if (trySend(mapped).isFailure && acceptingCallbacks.get()) {
+ if (initialChange.complete(mapped)) return@subscribe
+ if (latestChanges.trySend(mapped).isFailure && acceptingCallbacks.get()) {
deliveryFailure.complete(observerDeliveryFailure())
}
}
}
try {
- throw deliveryFailure.await()
+ val initial =
+ select<ApplicationChange> {
+ deliveryFailure.onAwait { throw it }
+ initialChange.onAwait { it }
+ }
+ emit(initial)
+ while (true) {
+ val latest =
+ select<ApplicationChange> {
+ deliveryFailure.onAwait { throw it }
+ latestChanges.onReceive { it }
+ }
+ emit(latest)
+ }
} finally {
acceptingCallbacks.set(false)
+ latestChanges.close()
withContext(NonCancellable) { subscription.unsubscribe() }
}
- }.buffer(Channel.CONFLATED)
+ }
override suspend fun execute(command: ApplicationCommand): ApplicationCommandResult =
when (command) {
diff --git a/app/desktop/src/test/kotlin/org/harvestcircle/application/NativeRuntimeMappingsTest.kt b/app/desktop/src/test/kotlin/org/harvestcircle/application/NativeRuntimeMappingsTest.kt
@@ -9,8 +9,10 @@ import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.take
import kotlinx.coroutines.flow.toList
+import kotlinx.coroutines.test.UnconfinedTestDispatcher
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
+import kotlinx.coroutines.withTimeout
import org.harvestcircle.ffi.ActiveIdentityDto
import org.harvestcircle.ffi.AppLifecycleDto
import org.harvestcircle.ffi.AppSnapshotDto
@@ -214,6 +216,26 @@ class NativeRuntimeMappingsTest {
}
@Test
+ fun snapshotChangeMappingsPreserveInitialAndCoalescedPredecessorMetadata() {
+ val initial = SnapshotChangeDto(emptySnapshot(), null).toApplicationChange()
+ assertEquals(SnapshotRevision(0UL), initial.snapshot.revision)
+ assertNull(initial.previousRevision)
+
+ val latest = SnapshotChangeDto(populatedSnapshot(100UL), 99UL).toApplicationChange()
+ assertEquals(SnapshotRevision(100UL), latest.snapshot.revision)
+ assertEquals(SnapshotRevision(99UL), latest.previousRevision)
+ }
+
+ @Test
+ fun snapshotChangeMappingsRejectEqualAndFuturePredecessors() {
+ listOf(2UL, 3UL).forEach { previousRevision ->
+ assertFailsWith<IllegalArgumentException> {
+ SnapshotChangeDto(populatedSnapshot(2UL), previousRevision).toApplicationChange()
+ }
+ }
+ }
+
+ @Test
fun nativeAndUnknownErrorsBecomeStructuredAndSecretSafe() {
val native =
HarvestCircleException
@@ -349,6 +371,80 @@ class NativeHarvestCircleRuntimeTest {
}
@Test
+ fun synchronousRegistrationBurstPreservesInitialBeforeLatestSnapshot() =
+ runTest(UnconfinedTestDispatcher()) {
+ val burst =
+ listOf(SnapshotChangeDto(emptySnapshot(), null)) +
+ (1UL..100UL).map { revision ->
+ SnapshotChangeDto(populatedSnapshot(revision), revision - 1UL)
+ }
+ val port = FakeNativeCorePort(registrationChanges = burst)
+ val runtime = NativeHarvestCircleRuntime(port)
+ val changes = withTimeout(1_000) { runtime.changes().take(2).toList() }
+
+ assertEquals(listOf(0UL, 100UL), changes.map { it.snapshot.revision.value })
+ assertNull(changes[0].previousRevision)
+ assertEquals(SnapshotRevision(99UL), changes[1].previousRevision)
+ assertEquals(1, port.subscriptionCalls)
+ assertEquals(1, port.subscriptionCloseCalls)
+ }
+
+ @Test
+ fun observerRegistrationFailurePreservesTypedSafeNativeProblem() =
+ runTest {
+ val port =
+ FakeNativeCorePort(
+ registrationFailure =
+ HarvestCircleException.Failure(
+ code = WireErrorCode.OBSERVER_REGISTRATION_FAILED,
+ category = WireErrorCategory.LIFECYCLE,
+ retryable = true,
+ recoveryAction = WireRecoveryAction.RETRY,
+ correlationId = null,
+ safeMessage = "The change observer could not be registered.",
+ ),
+ )
+ val runtime = NativeHarvestCircleRuntime(port)
+ val failure =
+ assertFailsWith<ApplicationFailure> {
+ withTimeout(1_000) { runtime.changes().first() }
+ }
+
+ assertEquals(ApplicationErrorCode.ObserverRegistrationFailed, failure.problem.code)
+ assertEquals(ApplicationErrorCategory.Lifecycle, failure.problem.category)
+ assertTrue(failure.problem.retryable)
+ assertEquals(RecoveryAction.Retry, failure.problem.recoveryAction)
+ assertEquals("The change observer could not be registered.", failure.problem.safeMessage)
+ assertEquals(1, port.subscriptionCalls)
+ assertEquals(0, port.subscriptionCloseCalls)
+ }
+
+ @Test
+ fun invalidObserverSnapshotFailsTypedWithoutExposingPayload() =
+ runTest {
+ val invalidPublicKey = "isolated-invalid-observer-public-key"
+ val port = FakeNativeCorePort()
+ val runtime = NativeHarvestCircleRuntime(port)
+ val pending =
+ async {
+ assertFailsWith<ApplicationFailure> {
+ withTimeout(1_000) { runtime.changes().first() }
+ }
+ }
+ runCurrent()
+
+ port.emit(SnapshotChangeDto(populatedSnapshot(2UL).copy(selectedPublicKeyHex = invalidPublicKey), 1UL))
+ val failure = pending.await()
+
+ assertEquals(ApplicationErrorCode.ObserverRegistrationFailed, failure.problem.code)
+ assertEquals(ApplicationErrorCategory.Lifecycle, failure.problem.category)
+ assertFalse(failure.problem.retryable)
+ assertEquals(RecoveryAction.RestartApplication, failure.problem.recoveryAction)
+ assertFalse(failure.problem.safeMessage.contains(invalidPublicKey))
+ assertEquals(1, port.subscriptionCloseCalls)
+ }
+
+ @Test
fun invalidObserverDeliveryFailsTypedAndUnsubscribesOnce() =
runTest {
val port = FakeNativeCorePort()
@@ -373,7 +469,10 @@ class NativeHarvestCircleRuntimeTest {
}
}
-private class FakeNativeCorePort : NativeCorePort {
+private class FakeNativeCorePort(
+ private val registrationChanges: List<SnapshotChangeDto> = emptyList(),
+ private val registrationFailure: Exception? = null,
+) : NativeCorePort {
val generated = FakeGeneratedRecoveryHandle()
val removal = FakeRemovalHandle()
var importedSecret: ByteArray? = null
@@ -381,6 +480,7 @@ private class FakeNativeCorePort : NativeCorePort {
var closed = false
var subscriptionClosed = false
var subscriptionCloseCalls = 0
+ var subscriptionCalls = 0
private var observer: ((SnapshotChangeDto) -> Unit)? = null
private val snapshot = populatedSnapshot(2UL)
@@ -389,7 +489,10 @@ private class FakeNativeCorePort : NativeCorePort {
override suspend fun bootstrap(): AppSnapshotDto = snapshot
override suspend fun subscribe(onChange: (SnapshotChangeDto) -> Unit): NativeSubscriptionHandle {
+ subscriptionCalls += 1
+ registrationFailure?.let { throw it }
observer = onChange
+ registrationChanges.forEach(onChange)
return NativeSubscriptionHandle {
subscriptionCloseCalls += 1
subscriptionClosed = true
diff --git a/core/crates/harvestcircle_ffi/src/observer.rs b/core/crates/harvestcircle_ffi/src/observer.rs
@@ -245,8 +245,9 @@ mod tests {
use std::time::Duration;
use harvestcircle_application::{
- RelayAccess, RelayConfiguration, RelayEndpoint, RelayUrlPolicy,
+ DurableRequestId, RelayAccess, RelayConfiguration, RelayEndpoint, RelayUrlPolicy,
};
+ use harvestcircle_domain::SecretKeyInput;
use nostr::{EventBuilder, Keys, Metadata};
use nostr_relay_builder::MockRelay;
use nostr_sdk::Client;
@@ -267,6 +268,28 @@ mod tests {
struct PanickingObserver;
+ struct GatedObserver {
+ changes: Mutex<Vec<SnapshotChangeDto>>,
+ entered: Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
+ release: Mutex<std::sync::mpsc::Receiver<()>>,
+ delivered: tokio::sync::Notify,
+ }
+
+ impl HarvestCircleChangeObserver for Arc<GatedObserver> {
+ fn on_change(&self, change: SnapshotChangeDto) {
+ self.changes.lock().expect("changes").push(change);
+ if let Some(entered) = self.entered.lock().expect("initial callback gate").take() {
+ entered.send(()).expect("initial callback entered");
+ self.release
+ .lock()
+ .expect("callback release gate")
+ .recv_timeout(OBSERVER_DELIVERY_TIMEOUT)
+ .expect("release initial callback");
+ }
+ self.delivered.notify_one();
+ }
+ }
+
impl HarvestCircleChangeObserver for PanickingObserver {
fn on_change(&self, _change: SnapshotChangeDto) {
panic!("injected host callback failure");
@@ -362,6 +385,124 @@ mod tests {
}
#[test]
+ fn slow_callback_recovers_final_tail_after_actor_queue_saturation() {
+ let runtime = tokio::runtime::Builder::new_multi_thread()
+ .worker_threads(2)
+ .enable_all()
+ .build()
+ .expect("two-worker callback test runtime");
+ runtime.block_on(async {
+ let core = core().await;
+ let mut public_keys = Vec::with_capacity(2);
+ for request_id in [
+ "01890f3e-7b1c-7000-8000-000000000050",
+ "01890f3e-7b1c-7000-8000-000000000051",
+ ] {
+ let imported = core
+ .inner
+ .actor
+ .import_secret_key(
+ DurableRequestId::parse(request_id).expect("test import request"),
+ core.inner.actor.snapshot().revision(),
+ SecretKeyInput::parse(Keys::generate().secret_key().to_secret_hex())
+ .expect("ephemeral identity input"),
+ OBSERVER_DELIVERY_TIMEOUT,
+ )
+ .await
+ .expect("import into isolated memory secret store");
+ public_keys.push(imported.identity().public_key());
+ }
+ assert_ne!(public_keys[0], public_keys[1]);
+ core.inner
+ .actor
+ .select_identity(public_keys[1])
+ .await
+ .expect("select second identity before observation");
+ let initial_revision = core.snapshot().revision;
+ let (entered_sender, entered_receiver) = tokio::sync::oneshot::channel();
+ let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1);
+ let observer = Arc::new(GatedObserver {
+ changes: Mutex::new(Vec::new()),
+ entered: Mutex::new(Some(entered_sender)),
+ release: Mutex::new(release_receiver),
+ delivered: tokio::sync::Notify::new(),
+ });
+ let subscription = core
+ .subscribe_changes_v2(Box::new(Arc::clone(&observer)))
+ .await
+ .expect("subscribe slow callback");
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered_receiver)
+ .await
+ .expect("initial callback deadline")
+ .expect("initial callback entered");
+
+ let queued_changes = super::OBSERVER_CHANGE_CAPACITY.get();
+ let publications = queued_changes + 2;
+ let final_revision = tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async {
+ let mut final_revision = initial_revision;
+ for offset in 0..publications {
+ final_revision = core
+ .inner
+ .actor
+ .select_identity(public_keys[offset % public_keys.len()])
+ .await
+ .expect("alternate selected identity")
+ .revision()
+ .value();
+ }
+ final_revision
+ })
+ .await
+ .expect("bounded actor publications");
+ assert_eq!(
+ final_revision,
+ initial_revision + u64::try_from(publications).expect("publication count")
+ );
+ assert_eq!(observer.changes.lock().expect("changes").len(), 1);
+ release_sender.send(()).expect("release slow callback");
+
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async {
+ loop {
+ let delivered = observer.delivered.notified();
+ if observer
+ .changes
+ .lock()
+ .expect("changes")
+ .last()
+ .is_some_and(|change| change.snapshot.revision == final_revision)
+ {
+ break;
+ }
+ delivered.await;
+ }
+ })
+ .await
+ .expect("final callback without a later publication");
+
+ let changes = observer.changes.lock().expect("changes").clone();
+ let mut expected_revisions = vec![initial_revision];
+ expected_revisions.extend(
+ (1..=queued_changes)
+ .map(|offset| initial_revision + u64::try_from(offset).expect("queue offset")),
+ );
+ expected_revisions.push(final_revision);
+ assert_eq!(
+ changes
+ .iter()
+ .map(|change| change.snapshot.revision)
+ .collect::<Vec<_>>(),
+ expected_revisions
+ );
+ assert_eq!(changes[0].previous_revision, None);
+ for change in &changes[1..] {
+ assert_eq!(change.previous_revision, Some(change.snapshot.revision - 1));
+ }
+ subscription.unsubscribe().await;
+ core.shutdown_v2().await.expect("shutdown");
+ });
+ }
+
+ #[test]
fn core_close_deregisters_all_observers_and_rejects_new_subscriptions() {
test_runtime().block_on(async {
let core = core().await;