commit ad175ba7ce1b615c748ef3e867b71996622d5946
parent 3e2b0e55a26ca3e3f6d71619087770cd05b5297c
Author: triesap <tyson@radroots.org>
Date: Fri, 2 Oct 2026 09:00:52 +0000
runtime: qualify observation cancellation and close
- Retain exact observer handles and admission through canceled join waits
- Share fixed observer admission across pending and installed registrations
- Reclaim abandoned subscriptions and protect native collector acquisition
- Verify seventeen lifecycle regressions and independent spec review
Diffstat:
7 files changed, 1221 insertions(+), 110 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
@@ -13,6 +13,7 @@ import org.harvestcircle.ffi.classifyNostrReference
import org.harvestcircle.ffi.compatibilityDescriptor
import org.harvestcircle.testbridge.ffi.HarvestCircleTestBridge
import org.harvestcircle.testbridge.ffi.TestBridgeException
+import org.harvestcircle.testbridge.ffi.TestLifecycle
import java.nio.file.Files
import java.nio.file.Path
import kotlin.io.path.readBytes
@@ -21,11 +22,65 @@ import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
+import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.test.fail
class NativeRuntimeIntegrationTest {
@Test
+ fun generatedActorObserverDiscardsStoppedQueueBeforeRestartAndFullShutdown() {
+ val dataRoot = Files.createTempDirectory("harvestcircle-generated-observer-close-")
+ try {
+ val bridge = HarvestCircleTestBridge.open(dataRoot.toString())
+ try {
+ bridge.bootstrap()
+ val publicKeys =
+ listOf(
+ "00000000-0000-7000-8000-000000000033",
+ "00000000-0000-7000-8000-000000000034",
+ ).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()
+ }
+ }
+ val initial = bridge.selectIdentity(publicKeys[1])
+ bridge.startObserver()
+ assertEquals(initial, assertNotNull(bridge.nextObservedSnapshot(2_000UL)))
+ repeat(18) { offset ->
+ val selected = bridge.selectIdentity(publicKeys[offset % publicKeys.size])
+ assertEquals(initial.revision + offset.toULong() + 1UL, selected.revision)
+ }
+ val finalSnapshot = bridge.snapshot()
+ assertTrue(bridge.stopObserver())
+ assertFalse(bridge.stopObserver())
+
+ bridge.startObserver()
+ assertEquals(finalSnapshot, assertNotNull(bridge.nextObservedSnapshot(2_000UL)))
+ assertNull(bridge.nextObservedSnapshot(0UL))
+ repeat(18) { offset ->
+ val selected = bridge.selectIdentity(publicKeys[offset % publicKeys.size])
+ assertEquals(finalSnapshot.revision + offset.toULong() + 1UL, selected.revision)
+ }
+ val closed = bridge.shutdown()
+
+ assertEquals(TestLifecycle.CLOSED, closed.lifecycle)
+ assertEquals(finalSnapshot.revision + 18UL, closed.revision)
+ assertFalse(bridge.stopObserver())
+ assertFailsWith<TestBridgeException.Failure> { bridge.startObserver() }
+ } finally {
+ bridge.close()
+ }
+ } finally {
+ deleteTree(dataRoot)
+ }
+ }
+
+ @Test
fun generatedActorObserverRecoversFinalTailAfterSaturation() {
val dataRoot = Files.createTempDirectory("harvestcircle-generated-observer-")
try {
diff --git a/app/desktop/src/main/kotlin/org/harvestcircle/application/NativeHarvestCircleRuntime.kt b/app/desktop/src/main/kotlin/org/harvestcircle/application/NativeHarvestCircleRuntime.kt
@@ -60,24 +60,27 @@ class NativeHarvestCircleRuntime internal constructor(
val deliveryFailure = CompletableDeferred<ApplicationFailure>()
val initialChange = CompletableDeferred<ApplicationChange>()
val latestChanges = Channel<ApplicationChange>(Channel.CONFLATED)
- val subscription =
- callNative {
- native.subscribe { change ->
- if (!acceptingCallbacks.get()) return@subscribe
- val mapped =
- try {
- change.toApplicationChange()
- } catch (_: Exception) {
- deliveryFailure.complete(observerDeliveryFailure())
- return@subscribe
+ var subscription: NativeSubscriptionHandle? = null
+ try {
+ withContext(NonCancellable) {
+ subscription =
+ callNative {
+ native.subscribe { change ->
+ if (!acceptingCallbacks.get()) return@subscribe
+ val mapped =
+ try {
+ change.toApplicationChange()
+ } catch (_: Exception) {
+ deliveryFailure.complete(observerDeliveryFailure())
+ return@subscribe
+ }
+ if (initialChange.complete(mapped)) return@subscribe
+ if (latestChanges.trySend(mapped).isFailure && acceptingCallbacks.get()) {
+ deliveryFailure.complete(observerDeliveryFailure())
+ }
}
- if (initialChange.complete(mapped)) return@subscribe
- if (latestChanges.trySend(mapped).isFailure && acceptingCallbacks.get()) {
- deliveryFailure.complete(observerDeliveryFailure())
}
- }
}
- try {
val initial =
select<ApplicationChange> {
deliveryFailure.onAwait { throw it }
@@ -95,7 +98,7 @@ class NativeHarvestCircleRuntime internal constructor(
} finally {
acceptingCallbacks.set(false)
latestChanges.close()
- withContext(NonCancellable) { subscription.unsubscribe() }
+ withContext(NonCancellable) { subscription?.unsubscribe() }
}
}
diff --git a/app/desktop/src/test/kotlin/org/harvestcircle/application/NativeRuntimeMappingsTest.kt b/app/desktop/src/test/kotlin/org/harvestcircle/application/NativeRuntimeMappingsTest.kt
@@ -4,11 +4,13 @@ package org.harvestcircle.application
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.async
+import kotlinx.coroutines.cancelAndJoin
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.launch
import kotlinx.coroutines.test.UnconfinedTestDispatcher
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
@@ -467,6 +469,124 @@ class NativeHarvestCircleRuntimeTest {
assertEquals(RecoveryAction.RestartApplication, failure.problem.recoveryAction)
assertEquals(1, port.subscriptionCloseCalls)
}
+
+ @Test
+ fun cancelledRegistrationBeforeInstallationReleasesAdmission() =
+ runTest {
+ val port = FakeNativeCorePort()
+ val release = CompletableDeferred<Unit>()
+ port.beforeRegistration = release
+ val runtime = NativeHarvestCircleRuntime(port)
+ val delivered = mutableListOf<ApplicationChange>()
+ val collector = launch { runtime.changes().collect { delivered.add(it) } }
+ runCurrent()
+ assertEquals(1, port.subscriptionCalls)
+ assertEquals(0, port.activeSubscriptionCount)
+
+ collector.cancel()
+ runCurrent()
+ release.complete(Unit)
+ collector.join()
+
+ assertEquals(0, port.activeSubscriptionCount)
+ assertTrue(delivered.isEmpty())
+ runtime.shutdown()
+ }
+
+ @Test
+ fun cancelledRegistrationAfterInstallationUnsubscribesBeforeCollectorFinishes() =
+ runTest {
+ val port = FakeNativeCorePort()
+ val release = CompletableDeferred<Unit>()
+ port.afterRegistration = release
+ val runtime = NativeHarvestCircleRuntime(port)
+ val delivered = mutableListOf<ApplicationChange>()
+ val collector = launch { runtime.changes().collect { delivered.add(it) } }
+ runCurrent()
+ assertTrue(port.registrationInstalled.isCompleted)
+ assertEquals(1, port.activeSubscriptionCount)
+
+ collector.cancel()
+ runCurrent()
+ release.complete(Unit)
+ collector.join()
+ port.emitToSubscription(0, SnapshotChangeDto(populatedSnapshot(3UL), 2UL))
+ runCurrent()
+ val retainedAdmission = port.activeSubscriptionCount
+ val unsubscribeCalls = port.subscriptionCloseCalls
+ runtime.shutdown()
+
+ assertTrue(delivered.isEmpty())
+ assertEquals(0, retainedAdmission, "cancelled installed registration must release native admission")
+ assertEquals(1, unsubscribeCalls)
+ }
+
+ @Test
+ fun collectorCancellationWaitsForUnsubscribeAndIgnoresStaleCallback() =
+ runTest {
+ val port = FakeNativeCorePort()
+ val release = CompletableDeferred<Unit>()
+ port.unsubscribeRelease = release
+ val runtime = NativeHarvestCircleRuntime(port)
+ val delivered = mutableListOf<ApplicationChange>()
+ val collector = launch { runtime.changes().collect { delivered.add(it) } }
+ runCurrent()
+ port.emit(SnapshotChangeDto(populatedSnapshot(2UL), 1UL))
+ runCurrent()
+
+ collector.cancel()
+ runCurrent()
+ assertTrue(port.unsubscribeStarted.isCompleted)
+ val completedBeforeUnsubscribe = collector.isCompleted
+ port.emitToSubscription(0, SnapshotChangeDto(populatedSnapshot(3UL), 2UL))
+ runCurrent()
+ release.complete(Unit)
+ collector.join()
+ collector.cancelAndJoin()
+ runtime.shutdown()
+
+ assertFalse(completedBeforeUnsubscribe)
+ assertEquals(listOf(2UL), delivered.map { it.snapshot.revision.value })
+ assertEquals(1, port.subscriptionCloseCalls)
+ assertEquals(0, port.activeSubscriptionCount)
+ }
+
+ @Test
+ fun cancelledCollectorCannotDeliverStaleAccountSnapshotToReplacementCollector() =
+ runTest {
+ val port = FakeNativeCorePort()
+ val runtime = NativeHarvestCircleRuntime(port)
+ val oldChanges = mutableListOf<ApplicationChange>()
+ val oldCollector = launch { runtime.changes().collect { oldChanges.add(it) } }
+ runCurrent()
+ port.emitToSubscription(0, SnapshotChangeDto(populatedSnapshot(2UL), 1UL))
+ runCurrent()
+ oldCollector.cancelAndJoin()
+
+ val replacement = async { runtime.changes().first() }
+ runCurrent()
+ val switchedIdentity = nativeIdentity().copy(publicKeyHex = "02".repeat(32))
+ val switchedSnapshot =
+ populatedSnapshot(3UL).copy(
+ identities = listOf(switchedIdentity),
+ selectedPublicKeyHex = switchedIdentity.publicKeyHex,
+ sessionSubjectPublicKeyHex = switchedIdentity.publicKeyHex,
+ activeIdentity = populatedSnapshot(3UL).activeIdentity?.copy(identity = switchedIdentity),
+ )
+ port.emitToSubscription(0, SnapshotChangeDto(populatedSnapshot(4UL), 3UL))
+ runCurrent()
+ val staleCompletedReplacement = replacement.isCompleted
+ port.emitToSubscription(1, SnapshotChangeDto(switchedSnapshot, 2UL))
+ val current = replacement.await()
+ runtime.shutdown()
+
+ assertFalse(staleCompletedReplacement)
+ assertEquals(listOf(2UL), oldChanges.map { it.snapshot.revision.value })
+ assertEquals(IdentityId.fromPublicKeyHex(switchedIdentity.publicKeyHex), current.snapshot.selectedIdentityId)
+ assertEquals(SnapshotRevision(3UL), current.snapshot.revision)
+ assertEquals(2, port.subscriptionCloseCalls)
+ assertEquals(0, port.activeSubscriptionCount)
+ }
}
private class FakeNativeCorePort(
@@ -481,6 +601,14 @@ private class FakeNativeCorePort(
var subscriptionClosed = false
var subscriptionCloseCalls = 0
var subscriptionCalls = 0
+ var beforeRegistration: CompletableDeferred<Unit>? = null
+ var afterRegistration: CompletableDeferred<Unit>? = null
+ var unsubscribeRelease: CompletableDeferred<Unit>? = null
+ val registrationInstalled = CompletableDeferred<Unit>()
+ val unsubscribeStarted = CompletableDeferred<Unit>()
+ private val observers = mutableListOf<(SnapshotChangeDto) -> Unit>()
+ private val activeSubscriptions = mutableSetOf<Int>()
+ val activeSubscriptionCount: Int get() = activeSubscriptions.size
private var observer: ((SnapshotChangeDto) -> Unit)? = null
private val snapshot = populatedSnapshot(2UL)
@@ -490,11 +618,20 @@ private class FakeNativeCorePort(
override suspend fun subscribe(onChange: (SnapshotChangeDto) -> Unit): NativeSubscriptionHandle {
subscriptionCalls += 1
+ beforeRegistration?.await()
registrationFailure?.let { throw it }
observer = onChange
+ val id = observers.size
+ observers.add(onChange)
+ activeSubscriptions.add(id)
+ registrationInstalled.complete(Unit)
registrationChanges.forEach(onChange)
+ afterRegistration?.await()
return NativeSubscriptionHandle {
subscriptionCloseCalls += 1
+ unsubscribeStarted.complete(Unit)
+ unsubscribeRelease?.await()
+ activeSubscriptions.remove(id)
subscriptionClosed = true
}
}
@@ -503,6 +640,13 @@ private class FakeNativeCorePort(
checkNotNull(observer)(change)
}
+ fun emitToSubscription(
+ id: Int,
+ change: SnapshotChangeDto,
+ ) {
+ observers[id](change)
+ }
+
override suspend fun beginGeneratedIdentity(): NativeGeneratedRecoveryHandle = generated
override suspend fun acknowledgeGeneratedIdentity(
diff --git a/core/crates/harvestcircle_application/src/change_stream.rs b/core/crates/harvestcircle_application/src/change_stream.rs
@@ -163,6 +163,8 @@ impl OrderedSnapshotChanges {
}
pub fn publish(&mut self, snapshot: AppSnapshot) {
+ self.subscribers
+ .retain(|_, subscriber| !subscriber.sender.is_closed());
if self.closed || snapshot.revision() <= self.latest.revision() {
return;
}
@@ -468,6 +470,27 @@ mod tests {
assert!(next_delivery(&mut receiver).await.is_none());
}
+ #[tokio::test]
+ async fn unchanged_publish_reclaims_abandoned_receiver_registration() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (id, receiver) = changes.subscribe(NonZeroUsize::MIN).expect("subscription");
+ drop(receiver);
+
+ changes.publish(snapshot(0));
+
+ assert!(
+ !changes.unsubscribe(id),
+ "unchanged observation must reclaim a closed receiver"
+ );
+ assert_eq!(changes.last_revision(), revision(0));
+ let (_, mut replacement) = changes.subscribe(NonZeroUsize::MIN).expect("replacement");
+ let initial = next_change(&mut replacement, "replacement initial snapshot").await;
+ assert_eq!(initial.revision(), revision(0));
+ assert!(initial.previous_revision().is_none());
+ changes.close();
+ assert!(next_delivery(&mut replacement).await.is_none());
+ }
+
async fn next_delivery(receiver: &mut SnapshotChangeReceiver) -> Option<SnapshotChange> {
tokio::time::timeout(RECEIVE_TIMEOUT, receiver.receive())
.await
diff --git a/core/crates/harvestcircle_ffi/src/commands.rs b/core/crates/harvestcircle_ffi/src/commands.rs
@@ -355,9 +355,16 @@ pub(crate) struct RuntimeCore {
pub(crate) observers: Mutex<
BTreeMap<
harvestcircle_application::ChangeSubscriptionId,
- Option<tokio::task::JoinHandle<()>>,
+ Arc<crate::observer::ObserverTask>,
>,
>,
+ pub(crate) retired_observers: Mutex<
+ BTreeMap<
+ harvestcircle_application::ChangeSubscriptionId,
+ Arc<crate::observer::ObserverTask>,
+ >,
+ >,
+ pub(crate) observer_admission: Arc<tokio::sync::Semaphore>,
pub(crate) close_state: AtomicU8,
pub(crate) close_gate: tokio::sync::Mutex<()>,
#[cfg(test)]
@@ -763,6 +770,10 @@ impl HarvestCircleAppCore {
host_runtime: Some(runtime),
keyring: Some(keyring),
observers: Mutex::new(BTreeMap::new()),
+ retired_observers: Mutex::new(BTreeMap::new()),
+ observer_admission: Arc::new(tokio::sync::Semaphore::new(
+ crate::observer::MAX_OBSERVERS,
+ )),
close_state: AtomicU8::new(0),
close_gate: tokio::sync::Mutex::new(()),
#[cfg(test)]
@@ -1107,6 +1118,10 @@ mod tests {
host_runtime: None,
keyring: None,
observers: std::sync::Mutex::new(std::collections::BTreeMap::new()),
+ retired_observers: std::sync::Mutex::new(std::collections::BTreeMap::new()),
+ observer_admission: Arc::new(tokio::sync::Semaphore::new(
+ crate::observer::MAX_OBSERVERS,
+ )),
close_state: std::sync::atomic::AtomicU8::new(0),
close_gate: tokio::sync::Mutex::new(()),
_test_directory: Some(directory),
diff --git a/core/crates/harvestcircle_ffi/src/observer.rs b/core/crates/harvestcircle_ffi/src/observer.rs
@@ -9,7 +9,19 @@ use crate::commands::RuntimeCore;
use crate::{AppSnapshotDto, HarvestCircleAppCore, HarvestCircleError};
const OBSERVER_CHANGE_CAPACITY: NonZeroUsize = NonZeroUsize::MIN.saturating_add(63);
-const MAX_OBSERVERS: usize = 32;
+pub(crate) const MAX_OBSERVERS: usize = 32;
+
+pub(crate) struct ObserverTask {
+ handle: tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>,
+ stop: Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
+ _admission: Arc<tokio::sync::OwnedSemaphorePermit>,
+}
+
+struct ObserverResources {
+ // Field drop order keeps callback destruction inside its admission reservation.
+ observer: Box<dyn HarvestCircleChangeObserver>,
+ admission: Arc<tokio::sync::OwnedSemaphorePermit>,
+}
#[derive(Clone, Debug, Eq, PartialEq)]
#[cfg_attr(not(coverage_nightly), derive(uniffi::Record))]
@@ -36,30 +48,153 @@ pub struct ObserverSubscription {
id: Mutex<Option<ChangeSubscriptionId>>,
}
+impl Drop for ObserverSubscription {
+ fn drop(&mut self) {
+ let Ok(retained_id) = self.id.get_mut() else {
+ return;
+ };
+ let (Some(core), Some(id)) = (self.core.upgrade(), *retained_id) else {
+ return;
+ };
+ let task = {
+ let Ok(observers) = core.observers.lock() else {
+ return;
+ };
+ let Ok(retired) = core.retired_observers.lock() else {
+ return;
+ };
+ observers.get(&id).or_else(|| retired.get(&id)).cloned()
+ };
+ if let Some(task) = task
+ && let Ok(mut stop) = task.stop.lock()
+ {
+ drop(stop.take());
+ }
+ }
+}
+
#[cfg_attr(not(coverage_nightly), uniffi::export)]
impl ObserverSubscription {
pub async fn unsubscribe(&self) {
let id = {
- let Ok(mut retained_id) = self.id.lock() else {
+ let Ok(retained_id) = self.id.lock() else {
return;
};
- retained_id.take()
+ *retained_id
};
let (Some(core), Some(id)) = (self.core.upgrade(), id) else {
return;
};
- let task = {
- let Ok(mut observers) = core.observers.lock() else {
- return;
- };
- observers.remove(&id)
- };
- if let Some(Some(task)) = task {
- task.abort();
- let _ = task.await;
+ let _close = core.close_gate.lock().await;
+ if finish_observers(&core, Some(id)).await.is_err() {
+ return;
}
let _ = core.actor.unsubscribe_changes(id).await;
+ if let Ok(mut retained_id) = self.id.lock() {
+ *retained_id = None;
+ }
+ }
+}
+
+async fn finish_observers(
+ core: &RuntimeCore,
+ selected: Option<ChangeSubscriptionId>,
+) -> Result<(), HarvestCircleError> {
+ let tasks = {
+ let observers = core
+ .observers
+ .lock()
+ .map_err(|_| crate::commands::internal_state_unavailable())?;
+ let retired = core
+ .retired_observers
+ .lock()
+ .map_err(|_| crate::commands::internal_state_unavailable())?;
+ observers
+ .iter()
+ .chain(retired.iter())
+ .filter(|(id, _)| selected.is_none_or(|selected| selected == **id))
+ .map(|(id, task)| (*id, Arc::clone(task)))
+ .collect::<Vec<_>>()
+ };
+ for (_, task) in &tasks {
+ if let Some(task) = task.handle.lock().await.as_ref() {
+ task.abort();
+ }
}
+ for (id, task) in tasks {
+ {
+ let mut retained = task.handle.lock().await;
+ if let Some(task) = retained.as_mut() {
+ // The registry retains the exact handle if this awaiting future is cancelled.
+ let _ = task.await;
+ }
+ drop(retained.take());
+ }
+ // Admission stays reserved until the join, including callback destruction, finishes.
+ let mut observers = core
+ .observers
+ .lock()
+ .map_err(|_| crate::commands::internal_state_unavailable())?;
+ let mut retired = core
+ .retired_observers
+ .lock()
+ .map_err(|_| crate::commands::internal_state_unavailable())?;
+ observers.remove(&id);
+ retired.remove(&id);
+ }
+ Ok(())
+}
+
+fn retire_observer(core: &RuntimeCore, id: ChangeSubscriptionId) {
+ let Ok(mut observers) = core.observers.lock() else {
+ return;
+ };
+ let Ok(mut retired) = core.retired_observers.lock() else {
+ return;
+ };
+ if let Some(task) = observers.remove(&id) {
+ retired.insert(id, task);
+ }
+}
+
+async fn forward_observer(
+ resources: ObserverResources,
+ runtime_core: Weak<RuntimeCore>,
+ mut subscription: harvestcircle_runtime::RuntimeChangeSubscription,
+ mut stopped: tokio::sync::oneshot::Receiver<()>,
+) {
+ let id = subscription.id();
+ loop {
+ let change = tokio::select! {
+ biased;
+ _ = &mut stopped => break,
+ change = subscription.receive() => change,
+ };
+ let Some(change) = change else {
+ break;
+ };
+ let Some(runtime_core) = runtime_core.upgrade() else {
+ break;
+ };
+ let delivery = SnapshotChangeDto {
+ snapshot: AppSnapshotDto::from_runtime(
+ change.snapshot(),
+ runtime_core.effective_lifecycle(),
+ ),
+ previous_revision: change
+ .previous_revision()
+ .map(harvestcircle_application::SnapshotRevision::value),
+ };
+ if catch_unwind(AssertUnwindSafe(|| resources.observer.on_change(delivery))).is_err() {
+ break;
+ }
+ }
+ if let Some(runtime_core) = runtime_core.upgrade() {
+ let _ = runtime_core.actor.unsubscribe_changes(id).await;
+ retire_observer(&runtime_core, id);
+ }
+ // Consume the whole bundle so async capture cannot separate the callback and reservation.
+ drop(resources);
}
#[cfg_attr(not(coverage_nightly), uniffi::export)]
@@ -73,28 +208,70 @@ impl HarvestCircleAppCore {
&self,
observer: Box<dyn HarvestCircleChangeObserver>,
) -> Result<Arc<ObserverSubscription>, HarvestCircleError> {
- if !self.inner.is_open() {
- return Err(closed_error());
- }
- let mut subscription = self
+ let resources = {
+ let _observers = self
+ .inner
+ .observers
+ .lock()
+ .map_err(|_| observer_registration_error())?;
+ let mut retired = self
+ .inner
+ .retired_observers
+ .lock()
+ .map_err(|_| observer_registration_error())?;
+ retired.retain(|_, task| match task.handle.try_lock() {
+ Ok(retained) => retained.as_ref().is_some_and(|task| !task.is_finished()),
+ Err(_) => true,
+ });
+ if !self.inner.is_open() {
+ return Err(closed_error());
+ }
+ let admission = Arc::clone(&self.inner.observer_admission)
+ .try_acquire_owned()
+ .map_err(|_| observer_registration_error())?;
+ ObserverResources {
+ observer,
+ admission: Arc::new(admission),
+ }
+ };
+ let subscription = self
.inner
.actor
.subscribe_changes(OBSERVER_CHANGE_CAPACITY)
.await
.map_err(HarvestCircleError::from)?;
let id = subscription.id();
- let observer: Arc<dyn HarvestCircleChangeObserver> = Arc::from(observer);
let runtime_core = Arc::downgrade(&self.inner);
+ let (stop, stopped) = tokio::sync::oneshot::channel();
let admitted = {
let mut observers = self
.inner
.observers
.lock()
.map_err(|_| observer_registration_error())?;
- if !self.inner.is_open() || observers.len() >= MAX_OBSERVERS {
+ let _retired = self
+ .inner
+ .retired_observers
+ .lock()
+ .map_err(|_| observer_registration_error())?;
+ if !self.inner.is_open() {
false
} else {
- observers.insert(id, None);
+ let admission = Arc::clone(&resources.admission);
+ let task = self.inner.runtime.spawn(forward_observer(
+ resources,
+ runtime_core,
+ subscription,
+ stopped,
+ ));
+ observers.insert(
+ id,
+ Arc::new(ObserverTask {
+ handle: tokio::sync::Mutex::new(Some(task)),
+ stop: Mutex::new(Some(stop)),
+ _admission: admission,
+ }),
+ );
true
}
};
@@ -106,48 +283,6 @@ impl HarvestCircleAppCore {
.map_err(HarvestCircleError::from)?;
return Err(observer_registration_error());
}
- let task = self.inner.runtime.spawn(async move {
- while let Some(change) = subscription.receive().await {
- let Some(runtime_core) = runtime_core.upgrade() else {
- break;
- };
- let delivery = SnapshotChangeDto {
- snapshot: AppSnapshotDto::from_runtime(
- change.snapshot(),
- runtime_core.effective_lifecycle(),
- ),
- previous_revision: change
- .previous_revision()
- .map(harvestcircle_application::SnapshotRevision::value),
- };
- if catch_unwind(AssertUnwindSafe(|| observer.on_change(delivery))).is_err() {
- break;
- }
- }
- if let Some(runtime_core) = runtime_core.upgrade() {
- let _ = runtime_core.actor.unsubscribe_changes(id).await;
- if let Ok(mut observers) = runtime_core.observers.lock() {
- observers.remove(&id);
- }
- }
- });
- let retained = {
- let mut observers = self.inner.observers.lock().map_err(|_| {
- task.abort();
- observer_registration_error()
- })?;
- if let Some(slot) = observers.get_mut(&id) {
- *slot = Some(task);
- true
- } else {
- task.abort();
- false
- }
- };
- if !retained {
- let _ = self.inner.actor.unsubscribe_changes(id).await;
- return Err(closed_error());
- }
Ok(Arc::new(ObserverSubscription {
core: Arc::downgrade(&self.inner),
id: Mutex::new(Some(id)),
@@ -164,10 +299,22 @@ impl HarvestCircleAppCore {
///
/// Returns a safe closed or timeout error when shutdown cannot complete.
pub async fn shutdown_v2(&self) -> Result<ShutdownReceiptDto, HarvestCircleError> {
- let _ = self
- .inner
- .close_state
- .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire);
+ {
+ let _observers = self
+ .inner
+ .observers
+ .lock()
+ .map_err(|_| crate::commands::internal_state_unavailable())?;
+ let _retired = self
+ .inner
+ .retired_observers
+ .lock()
+ .map_err(|_| crate::commands::internal_state_unavailable())?;
+ let _ =
+ self.inner
+ .close_state
+ .compare_exchange(0, 1, Ordering::AcqRel, Ordering::Acquire);
+ }
let _close = self.inner.close_gate.lock().await;
if self.inner.close_state.load(Ordering::Acquire) == 2 {
return Ok(ShutdownReceiptDto {
@@ -175,23 +322,7 @@ impl HarvestCircleAppCore {
closed: true,
});
}
- let handles = std::mem::take(
- &mut *self
- .inner
- .observers
- .lock()
- .map_err(|_| crate::commands::internal_state_unavailable())?,
- );
- let tasks = handles
- .into_values()
- .flatten()
- .collect::<Vec<tokio::task::JoinHandle<()>>>();
- for task in &tasks {
- task.abort();
- }
- for task in tasks {
- let _ = task.await;
- }
+ finish_observers(&self.inner, None).await?;
self.inner
.actor
.close()
@@ -241,7 +372,10 @@ fn observer_registration_error() -> HarvestCircleError {
#[cfg(test)]
#[cfg_attr(coverage_nightly, coverage(off))]
mod tests {
+ use std::future::Future;
+ use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
+ use std::task::Poll;
use std::time::Duration;
use harvestcircle_application::{
@@ -268,6 +402,36 @@ mod tests {
struct PanickingObserver;
+ #[derive(Default)]
+ struct ObserverDropState {
+ finished: AtomicBool,
+ released: AtomicBool,
+ }
+
+ struct GatedDropObserver {
+ entered: Option<tokio::sync::oneshot::Sender<()>>,
+ release: Mutex<std::sync::mpsc::Receiver<()>>,
+ state: Arc<ObserverDropState>,
+ }
+
+ impl HarvestCircleChangeObserver for GatedDropObserver {
+ fn on_change(&self, _change: SnapshotChangeDto) {}
+ }
+
+ impl Drop for GatedDropObserver {
+ fn drop(&mut self) {
+ if let Some(entered) = self.entered.take() {
+ let _ = entered.send(());
+ }
+ let released = self
+ .release
+ .get_mut()
+ .is_ok_and(|release| release.recv_timeout(OBSERVER_DELIVERY_TIMEOUT).is_ok());
+ self.state.released.store(released, Ordering::Release);
+ self.state.finished.store(true, Ordering::Release);
+ }
+ }
+
struct GatedObserver {
changes: Mutex<Vec<SnapshotChangeDto>>,
entered: Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
@@ -331,6 +495,8 @@ mod tests {
host_runtime: None,
keyring: None,
observers: Mutex::new(std::collections::BTreeMap::new()),
+ retired_observers: Mutex::new(std::collections::BTreeMap::new()),
+ observer_admission: Arc::new(tokio::sync::Semaphore::new(super::MAX_OBSERVERS)),
close_state: std::sync::atomic::AtomicU8::new(0),
close_gate: tokio::sync::Mutex::new(()),
_test_directory: Some(directory),
@@ -349,6 +515,8 @@ mod tests {
host_runtime: Some(host_runtime),
keyring: None,
observers: Mutex::new(std::collections::BTreeMap::new()),
+ retired_observers: Mutex::new(std::collections::BTreeMap::new()),
+ observer_admission: Arc::new(tokio::sync::Semaphore::new(super::MAX_OBSERVERS)),
close_state: std::sync::atomic::AtomicU8::new(0),
close_gate: tokio::sync::Mutex::new(()),
_test_directory: Some(directory),
@@ -384,6 +552,545 @@ mod tests {
});
}
+ #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+ async fn cancelled_unsubscribe_retry_waits_for_the_retained_callback_task() {
+ let core = core().await;
+ let (observer, entered, release) = gated_observer();
+ let subscription = core
+ .subscribe_changes_v2(Box::new(Arc::clone(&observer)))
+ .await
+ .expect("subscribe gated callback");
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered)
+ .await
+ .expect("callback entry deadline")
+ .expect("callback entered");
+ let task = observer_abort_handle(&core, &subscription);
+ let first_was_pending = {
+ let mut first = Box::pin(subscription.unsubscribe());
+ std::future::poll_fn(|context| Poll::Ready(first.as_mut().poll(context)))
+ .await
+ .is_pending()
+ };
+ let mut retry = Box::pin(subscription.unsubscribe());
+ let retry_was_pending =
+ std::future::poll_fn(|context| Poll::Ready(retry.as_mut().poll(context)))
+ .await
+ .is_pending();
+ let finished_before_release = task.is_finished();
+
+ release
+ .send(())
+ .expect("release callback before assertions");
+ if retry_was_pending {
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, retry)
+ .await
+ .expect("unsubscribe retry completion");
+ }
+ wait_for_aborted_task(&task).await;
+ core.shutdown_v2().await.expect("cleanup runtime");
+
+ assert!(
+ first_was_pending,
+ "first unsubscribe must wait for callback exit"
+ );
+ assert!(!finished_before_release, "callback was explicitly gated");
+ assert!(
+ retry_was_pending,
+ "retry must retain ownership of callback cleanup"
+ );
+ assert!(task.is_finished());
+ assert_eq!(observer.changes.lock().expect("changes").len(), 1);
+ subscription.unsubscribe().await;
+ }
+
+ #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+ async fn cancelled_shutdown_retry_finishes_observers_before_terminal_close() {
+ let core = core().await;
+ let (observer, entered, release) = gated_observer();
+ let subscription = core
+ .subscribe_changes_v2(Box::new(Arc::clone(&observer)))
+ .await
+ .expect("subscribe gated callback");
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered)
+ .await
+ .expect("callback entry deadline")
+ .expect("callback entered");
+ let task = observer_abort_handle(&core, &subscription);
+ let first_was_pending = {
+ let mut first = Box::pin(core.shutdown_v2());
+ std::future::poll_fn(|context| Poll::Ready(first.as_mut().poll(context)))
+ .await
+ .is_pending()
+ };
+ let mut retry = Box::pin(core.shutdown_v2());
+ let initial_retry =
+ std::future::poll_fn(|context| Poll::Ready(retry.as_mut().poll(context))).await;
+ // This command orders actor work after any close submitted by the retry.
+ let _ = core.inner.actor.bootstrap().await;
+ let retried = if initial_retry.is_pending() {
+ std::future::poll_fn(|context| Poll::Ready(retry.as_mut().poll(context))).await
+ } else {
+ initial_retry
+ };
+ let closed_before_callback_exit = retried.is_ready();
+ let finished_before_release = task.is_finished();
+
+ release
+ .send(())
+ .expect("release callback before assertions");
+ let receipt = match retried {
+ Poll::Ready(receipt) => receipt.expect("retry close receipt"),
+ Poll::Pending => tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, retry)
+ .await
+ .expect("resumed shutdown completion")
+ .expect("resumed shutdown receipt"),
+ };
+ wait_for_aborted_task(&task).await;
+ subscription.unsubscribe().await;
+ assert_eq!(core.shutdown_v2().await.expect("repeated close"), receipt);
+
+ assert!(
+ first_was_pending,
+ "first shutdown must wait for callback exit"
+ );
+ assert!(!finished_before_release, "callback was explicitly gated");
+ assert!(
+ !closed_before_callback_exit,
+ "terminal close must retain observer cleanup ownership"
+ );
+ assert!(receipt.closed);
+ assert!(task.is_finished());
+ assert!(core.inner.observers.lock().expect("observers").is_empty());
+ assert_eq!(observer.changes.lock().expect("changes").len(), 1);
+ }
+
+ #[tokio::test]
+ async fn abandoned_observer_handles_release_bounded_registration_admission() {
+ let core = core().await;
+ let observer = Arc::new(RecordingObserver::default());
+ for _ in 0..super::MAX_OBSERVERS {
+ let subscription = core
+ .subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))
+ .await
+ .expect("admit observer before abandonment");
+ drop(subscription);
+ }
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("actor ordering barrier");
+ let cleanup_completed = tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async {
+ while !core.inner.observers.lock().expect("observers").is_empty() {
+ tokio::task::yield_now().await;
+ }
+ })
+ .await
+ .is_ok();
+ let registered_after_abandonment = core.inner.observers.lock().expect("observers").len();
+ let replacement = core
+ .subscribe_changes_v2(Box::new(ArcObserver(observer)))
+ .await;
+ let replacement_was_admitted = replacement.is_ok();
+ if let Ok(replacement) = replacement {
+ replacement.unsubscribe().await;
+ }
+ core.shutdown_v2()
+ .await
+ .expect("cleanup abandoned observers");
+
+ assert!(
+ cleanup_completed,
+ "bounded observer task cleanup must finish after abandonment"
+ );
+ assert_eq!(
+ registered_after_abandonment, 0,
+ "abandoned handles must release admission"
+ );
+ assert!(
+ replacement_was_admitted,
+ "abandonment must not exhaust the fixed observer limit"
+ );
+ }
+
+ #[tokio::test]
+ async fn cancelled_native_registration_before_and_after_actor_reply_never_delivers() {
+ let core = core().await;
+ for cancel_after_actor_reply in [false, true] {
+ let observer = Arc::new(RecordingObserver::default());
+ let mut registration =
+ Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))));
+ let first_poll =
+ std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context)))
+ .await;
+ assert!(
+ first_poll.is_pending(),
+ "registration queued on the current-thread actor"
+ );
+ if cancel_after_actor_reply {
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("actor installed registration before cancellation");
+ }
+ drop(registration);
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("actor cancellation ordering barrier");
+ assert!(core.inner.observers.lock().expect("observers").is_empty());
+ assert!(observer.snapshots.lock().expect("snapshots").is_empty());
+ }
+ core.shutdown_v2().await.expect("cleanup runtime");
+ }
+
+ #[tokio::test]
+ async fn pending_native_registrations_share_the_fixed_observer_admission_limit() {
+ let core = core().await;
+ let observer = Arc::new(RecordingObserver::default());
+ let mut pending = Vec::with_capacity(super::MAX_OBSERVERS);
+ for _ in 0..super::MAX_OBSERVERS {
+ let mut registration =
+ Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))));
+ let first_poll =
+ std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context)))
+ .await;
+ assert!(
+ first_poll.is_pending(),
+ "admitted registration awaits the actor"
+ );
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("actor installed the unpolled reply");
+ pending.push(registration);
+ }
+ assert!(core.inner.observers.lock().expect("observers").is_empty());
+ let mut excess =
+ Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))));
+ let refused =
+ std::future::poll_fn(|context| Poll::Ready(excess.as_mut().poll(context))).await;
+ let refused_before_actor_allocation = matches!(
+ refused,
+ Poll::Ready(Err(crate::HarvestCircleError::Failure {
+ code: crate::WireErrorCode::ObserverRegistrationFailed,
+ ..
+ }))
+ );
+ drop(excess);
+ drop(pending);
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("cancelled reply cleanup barrier");
+ let replacement = core
+ .subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer))))
+ .await
+ .expect("replacement after pending cancellation");
+ replacement.unsubscribe().await;
+ core.shutdown_v2().await.expect("cleanup runtime");
+
+ assert!(
+ refused_before_actor_allocation,
+ "pending replies must consume the same fixed observer slots"
+ );
+ assert_eq!(
+ core.inner.observer_admission.available_permits(),
+ super::MAX_OBSERVERS
+ );
+ }
+
+ #[tokio::test]
+ async fn cancelled_pending_registration_retains_admission_until_callback_destruction_finishes()
+ {
+ let core = core().await;
+ let state = Arc::new(ObserverDropState::default());
+ let (entered_sender, entered_receiver) = tokio::sync::oneshot::channel();
+ let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1);
+ let pending_core = Arc::clone(&core);
+ let drop_state = Arc::clone(&state);
+ let mut gated_registration = Box::pin(async move {
+ pending_core
+ .subscribe_changes_v2(Box::new(GatedDropObserver {
+ entered: Some(entered_sender),
+ release: Mutex::new(release_receiver),
+ state: drop_state,
+ }))
+ .await
+ });
+ let first_poll =
+ std::future::poll_fn(|context| Poll::Ready(gated_registration.as_mut().poll(context)))
+ .await;
+ assert!(
+ first_poll.is_pending(),
+ "gated registration awaits the actor"
+ );
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("gated actor reply ready but unpolled");
+
+ let observer = Arc::new(RecordingObserver::default());
+ let mut pending = Vec::with_capacity(super::MAX_OBSERVERS - 1);
+ for _ in 1..super::MAX_OBSERVERS {
+ let mut registration =
+ Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))));
+ let first_poll =
+ std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context)))
+ .await;
+ assert!(
+ first_poll.is_pending(),
+ "peer registration awaits the actor"
+ );
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("peer actor reply ready but unpolled");
+ pending.push(registration);
+ }
+ assert_eq!(core.inner.observer_admission.available_permits(), 0);
+
+ let dropping = std::thread::spawn(move || drop(gated_registration));
+ let entered = tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered_receiver)
+ .await
+ .is_ok_and(|entered| entered.is_ok());
+ let destructor_was_gated = !state.finished.load(Ordering::Acquire);
+ let mut replacement =
+ Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))));
+ let replacement_poll =
+ std::future::poll_fn(|context| Poll::Ready(replacement.as_mut().poll(context))).await;
+ let refused_during_destruction = matches!(
+ replacement_poll,
+ Poll::Ready(Err(crate::HarvestCircleError::Failure {
+ code: crate::WireErrorCode::ObserverRegistrationFailed,
+ ..
+ }))
+ );
+
+ let released = release_sender.send(()).is_ok();
+ let joined = dropping.join().is_ok();
+ drop(replacement);
+ drop(pending);
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("cancelled registration cleanup barrier");
+ let resumed = core
+ .subscribe_changes_v2(Box::new(ArcObserver(observer)))
+ .await
+ .expect("registration after callback destruction");
+ resumed.unsubscribe().await;
+ core.shutdown_v2()
+ .await
+ .expect("cleanup runtime before assertions");
+
+ assert!(entered, "callback destructor entry was observed");
+ assert!(
+ destructor_was_gated,
+ "replacement was probed during callback destruction"
+ );
+ assert!(
+ released && joined,
+ "the destructor gate was released and its owned thread joined"
+ );
+ assert!(state.released.load(Ordering::Acquire));
+ assert!(state.finished.load(Ordering::Acquire));
+ assert!(
+ refused_during_destruction,
+ "callback destruction must retain its pending admission slot"
+ );
+ assert_eq!(
+ core.inner.observer_admission.available_permits(),
+ super::MAX_OBSERVERS
+ );
+ }
+
+ #[tokio::test]
+ async fn ready_native_registration_reply_cannot_install_after_terminal_close() {
+ let core = core().await;
+ let observer = Arc::new(RecordingObserver::default());
+ let mut registration =
+ Box::pin(core.subscribe_changes_v2(Box::new(ArcObserver(Arc::clone(&observer)))));
+ let first_poll =
+ std::future::poll_fn(|context| Poll::Ready(registration.as_mut().poll(context))).await;
+ assert!(
+ first_poll.is_pending(),
+ "registration queued on the current-thread actor"
+ );
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("actor reply ready before close");
+
+ let closed = core
+ .shutdown_v2()
+ .await
+ .expect("close before host registration resumes");
+ assert!(
+ registration.await.is_err(),
+ "closed registration must not install a callback task"
+ );
+
+ assert!(closed.closed);
+ assert!(observer.snapshots.lock().expect("snapshots").is_empty());
+ assert!(core.inner.observers.lock().expect("observers").is_empty());
+ assert!(
+ core.inner
+ .retired_observers
+ .lock()
+ .expect("retired observers")
+ .is_empty()
+ );
+ assert_eq!(
+ core.inner.observer_admission.available_permits(),
+ super::MAX_OBSERVERS
+ );
+ }
+
+ #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+ async fn full_native_callback_queue_terminates_before_shutdown_returns() {
+ let core = core().await;
+ let mut public_keys = Vec::new();
+ for request_id in [
+ "01890f3e-7b1c-7000-8000-000000000052",
+ "01890f3e-7b1c-7000-8000-000000000053",
+ ] {
+ 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("isolated memory-store import");
+ public_keys.push(imported.identity().public_key());
+ }
+ core.inner
+ .actor
+ .select_identity(public_keys[1])
+ .await
+ .expect("initial selection");
+ let initial_revision = core.snapshot().revision;
+ let (observer, entered, release) = gated_observer();
+ let subscription = core
+ .subscribe_changes_v2(Box::new(Arc::clone(&observer)))
+ .await
+ .expect("gated observer registration");
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, entered)
+ .await
+ .expect("callback entry deadline")
+ .expect("callback entered");
+ let task = observer_abort_handle(&core, &subscription);
+ let publications = super::OBSERVER_CHANGE_CAPACITY.get() + 2;
+ for offset in 0..publications {
+ core.inner
+ .actor
+ .select_identity(public_keys[offset % 2])
+ .await
+ .expect("revision-changing selection");
+ }
+ let final_revision = core.snapshot().revision;
+ let mut closing = Box::pin(core.shutdown_v2());
+ let first_close_poll =
+ std::future::poll_fn(|context| Poll::Ready(closing.as_mut().poll(context))).await;
+ let closing_was_pending = first_close_poll.is_pending();
+ release
+ .send(())
+ .expect("release saturated callback before assertions");
+ let receipt = match first_close_poll {
+ Poll::Ready(receipt) => receipt.expect("shutdown receipt"),
+ Poll::Pending => tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, closing)
+ .await
+ .expect("bounded full-queue close")
+ .expect("shutdown receipt"),
+ };
+ let callback_count_at_close = observer.changes.lock().expect("changes").len();
+ subscription.unsubscribe().await;
+ subscription.unsubscribe().await;
+ let repeated = core.shutdown_v2().await.expect("repeated close");
+ let refused = core
+ .subscribe_changes_v2(Box::new(PanickingObserver))
+ .await
+ .is_err();
+
+ assert!(
+ closing_was_pending,
+ "close must wait for the running callback"
+ );
+ assert_eq!(
+ final_revision,
+ initial_revision + u64::try_from(publications).expect("publication count")
+ );
+ assert_eq!(receipt.final_revision, final_revision);
+ assert!(receipt.closed);
+ assert_eq!(repeated, receipt);
+ assert!(task.is_finished());
+ assert!(refused);
+ assert!(core.inner.observers.lock().expect("observers").is_empty());
+ assert_eq!(
+ observer.changes.lock().expect("changes").len(),
+ callback_count_at_close
+ );
+ }
+
+ fn gated_observer() -> (
+ Arc<GatedObserver>,
+ tokio::sync::oneshot::Receiver<()>,
+ std::sync::mpsc::SyncSender<()>,
+ ) {
+ let (entered_sender, entered_receiver) = tokio::sync::oneshot::channel();
+ let (release_sender, release_receiver) = std::sync::mpsc::sync_channel(1);
+ (
+ Arc::new(GatedObserver {
+ changes: Mutex::new(Vec::new()),
+ entered: Mutex::new(Some(entered_sender)),
+ release: Mutex::new(release_receiver),
+ delivered: tokio::sync::Notify::new(),
+ }),
+ entered_receiver,
+ release_sender,
+ )
+ }
+
+ fn observer_abort_handle(
+ core: &HarvestCircleAppCore,
+ subscription: &crate::ObserverSubscription,
+ ) -> tokio::task::AbortHandle {
+ let id = subscription.id.lock().expect("id").expect("registered id");
+ let task = core
+ .inner
+ .observers
+ .lock()
+ .expect("observers")
+ .get(&id)
+ .expect("registered observer")
+ .clone();
+ let retained = task.handle.try_lock().expect("observer join not started");
+ retained.as_ref().expect("observer task").abort_handle()
+ }
+
+ async fn wait_for_aborted_task(task: &tokio::task::AbortHandle) {
+ tokio::time::timeout(OBSERVER_DELIVERY_TIMEOUT, async {
+ while !task.is_finished() {
+ tokio::task::yield_now().await;
+ }
+ })
+ .await
+ .expect("released callback task finishes");
+ }
+
#[test]
fn slow_callback_recovers_final_tail_after_actor_queue_saturation() {
let runtime = tokio::runtime::Builder::new_multi_thread()
@@ -520,16 +1227,20 @@ mod tests {
.lock()
.expect("subscription id")
.expect("active subscription id");
- let handle = core
+ let task = core
.inner
.observers
.lock()
.expect("observers")
- .get_mut(&id)
+ .get(&id)
.expect("registered observer")
- .take()
- .expect("observer task");
- handle.abort();
+ .clone();
+ task.handle
+ .lock()
+ .await
+ .as_ref()
+ .expect("observer task")
+ .abort();
let first = core.shutdown_v2().await.expect("shutdown");
let repeated = core.shutdown_v2().await.expect("repeated shutdown");
diff --git a/core/crates/harvestcircle_runtime/src/runtime_actor.rs b/core/crates/harvestcircle_runtime/src/runtime_actor.rs
@@ -837,6 +837,9 @@ impl RuntimeActor {
completion_sender: &mpsc::Sender<ProfileCompletion>,
) -> bool {
let (context, command, reply) = envelope.into_parts();
+ if matches!(command, RuntimeCommand::SubscribeChanges(_)) && reply.is_closed() {
+ return true;
+ }
if let Some(result) = self.preflight(context, &command) {
let _ = reply.send(CommandReceipt::new(context.request_id(), result));
return true;
@@ -871,7 +874,12 @@ impl RuntimeActor {
if matches!(result, CommandResult::Completed(_)) {
self.changes.publish(self.adapter.core().snapshot());
}
- let _ = reply.send(CommandReceipt::new(context.request_id(), result));
+ if let Err(receipt) = reply.send(CommandReceipt::new(context.request_id(), result))
+ && let CommandResult::Completed(RuntimeCommandValue::Subscription(subscription)) =
+ receipt.into_result()
+ {
+ let _ = self.changes.unsubscribe(subscription.id());
+ }
true
}
@@ -1485,9 +1493,11 @@ mod tests {
use std::time::{Duration, Instant};
use harvestcircle_application::{
- ActiveSessionBinding, BoxFuture, Clock, DurableRequestId, FailureSecretStore,
- InMemorySecretStore, NostrClient, ProfileFetchResult, RelayConfiguration, RuntimeLifecycle,
- SecretStore, SecretStoreOperation, SessionGeneration, SessionState, SnapshotRevision,
+ ActiveSessionBinding, ActorMailbox, BoxFuture, Clock, CommandContext, CommandSubmission,
+ DurableRequestId, FailureSecretStore, GeneratedKeyStage, InMemorySecretStore, NostrClient,
+ OrderedSnapshotChanges, ProfileFetchResult, RelayConfiguration, RequestId,
+ RuntimeLifecycle, SecretStore, SecretStoreOperation, SessionGeneration, SessionState,
+ SnapshotRevision,
};
use harvestcircle_domain::{
LocalKeyringBinding, NostrIdentityReference, PublicKey, SafeError, SafeErrorCode,
@@ -1496,8 +1506,9 @@ mod tests {
use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy};
use super::{
- DEFAULT_COMMAND_TIMEOUT, RuntimeActorHandle, RuntimeDependencies, command_unavailable,
- test_migration_build_identity, test_runtime_context,
+ DEFAULT_COMMAND_TIMEOUT, RuntimeActor, RuntimeActorHandle, RuntimeCommand,
+ RuntimeDependencies, command_unavailable, test_migration_build_identity,
+ test_runtime_context,
};
use crate::{InstallationIdentity, InstallationIdentitySource, UuidInstallationIdentitySource};
@@ -2321,6 +2332,155 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn cancelled_registration_before_installation_does_not_retain_actor_admission() {
+ let (handle, secrets) = actor().await;
+ let snapshot = handle.snapshot();
+ // Obtain the first identifier from an independent probe, without relying on its value.
+ let mut probe = OrderedSnapshotChanges::new(snapshot.clone());
+ let (probe_id, probe_receiver) = probe.subscribe(NonZeroUsize::MIN).expect("probe id");
+ drop(probe_receiver);
+ let mut owner = RuntimeActor {
+ adapter: Arc::clone(&handle.adapter),
+ secrets,
+ clock: Arc::new(FixedClock),
+ nostr: Arc::new(OfflineNostr),
+ lifecycle: Arc::clone(&handle.lifecycle),
+ runtime: tokio::runtime::Handle::current(),
+ session_generation: handle.session_generation(),
+ published_session_generation: Arc::clone(&handle.session_generation),
+ profile_tasks: std::collections::BTreeMap::new(),
+ changes: OrderedSnapshotChanges::new(snapshot),
+ published_foreground_session: Arc::clone(&handle.foreground_session),
+ generated_key_stage: GeneratedKeyStage::default(),
+ };
+ let (mailbox, mut receiver) = ActorMailbox::bounded(NonZeroUsize::MIN);
+ let context = CommandContext::new(
+ RequestId::new(1).expect("request id"),
+ None,
+ Instant::now() + DEFAULT_COMMAND_TIMEOUT,
+ );
+ let ticket = match mailbox
+ .submit(context, RuntimeCommand::SubscribeChanges(NonZeroUsize::MIN))
+ {
+ CommandSubmission::Accepted(ticket) => ticket,
+ CommandSubmission::Rejected(_) => panic!("empty test mailbox must admit registration"),
+ };
+ drop(ticket);
+ let envelope = receiver.recv().await.expect("queued registration");
+ let (completions, _completed) = tokio::sync::mpsc::channel(1);
+ assert!(owner.handle_command(envelope, &completions).await);
+ let retained_abandoned_registration = owner.changes.unsubscribe(probe_id);
+ drop(owner);
+ handle.close().await.expect("cleanup runtime");
+
+ assert!(
+ !retained_abandoned_registration,
+ "cancelled queued registration must not retain admission"
+ );
+ }
+
+ #[tokio::test]
+ async fn abandoned_installed_subscription_releases_actor_admission_without_mutation() {
+ let (actor, _) = actor().await;
+ let subscription = actor
+ .subscribe_changes(NonZeroUsize::MIN)
+ .await
+ .expect("installed actor subscription");
+ let id = subscription.id();
+ let unchanged_revision = actor.snapshot().revision();
+ drop(subscription);
+ // Awaiting an Observe command orders any cleanup through the same actor mailbox.
+ let snapshot = actor.bootstrap().await.expect("serialized cleanup barrier");
+ let retained_abandoned_registration = actor
+ .unsubscribe_changes(id)
+ .await
+ .expect("registration probe");
+ actor.close().await.expect("cleanup runtime");
+
+ assert_eq!(snapshot.revision(), unchanged_revision);
+ assert!(
+ !retained_abandoned_registration,
+ "receiver abandonment must release admission without a later state mutation"
+ );
+ }
+
+ #[tokio::test]
+ async fn full_actor_queues_preserve_initial_and_final_tail_then_terminate_on_close() {
+ let (actor, _) = actor().await;
+ let mut public_keys = Vec::new();
+ for value in [
+ "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7",
+ "6e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7",
+ ] {
+ public_keys.push(
+ actor
+ .import_secret_key_test(secret(value))
+ .await
+ .expect("memory-store import")
+ .identity()
+ .public_key(),
+ );
+ }
+ actor
+ .select_identity(public_keys[1])
+ .await
+ .expect("initial selection");
+ let initial_revision = actor.snapshot().revision();
+ let mut first = actor
+ .subscribe_changes(NonZeroUsize::MIN)
+ .await
+ .expect("first consumer");
+ let mut second = actor
+ .subscribe_changes(NonZeroUsize::MIN)
+ .await
+ .expect("second consumer");
+ for offset in 0..4 {
+ let selected = actor
+ .select_identity(public_keys[offset % 2])
+ .await
+ .expect("revision-changing selection");
+ assert_eq!(
+ selected.revision().value(),
+ initial_revision.value() + u64::try_from(offset).expect("offset") + 1
+ );
+ }
+ let final_snapshot = actor.snapshot();
+ tokio::time::timeout(Duration::from_secs(1), actor.close())
+ .await
+ .expect("bounded full-queue close")
+ .expect("close");
+ actor.close().await.expect("repeated close");
+ for subscription in [&mut first, &mut second] {
+ let initial = tokio::time::timeout(Duration::from_secs(1), subscription.receive())
+ .await
+ .expect("initial deadline")
+ .expect("initial");
+ assert_eq!(initial.revision(), initial_revision);
+ assert!(initial.previous_revision().is_none());
+ let tail = tokio::time::timeout(Duration::from_secs(1), subscription.receive())
+ .await
+ .expect("tail deadline")
+ .expect("final tail");
+ assert_eq!(tail.snapshot(), &final_snapshot);
+ assert_eq!(
+ tail.previous_revision()
+ .expect("producer predecessor")
+ .value(),
+ final_snapshot.revision().value() - 1
+ );
+ for _ in 0..2 {
+ assert!(
+ tokio::time::timeout(Duration::from_secs(1), subscription.receive())
+ .await
+ .expect("terminal delivery deadline")
+ .is_none()
+ );
+ }
+ }
+ assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed);
+ }
+
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn expired_queued_shutdown_does_not_close_runtime_later() {
let secrets = Arc::new(BlockingSecretStore::new());