commit 1509ee154f3e4af29d3a87e9424cfbe726e35b3a
parent b4370117a4110d8fb283122fffd57ef0e9f9dec6
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 22:02:47 +0000
runtime: remove unordered observer registry
- remove callback storage and publication from the canonical application core
- bridge FFI callbacks from actor-owned ordered subscriptions
- register Kotlin observation asynchronously before runtime bootstrap
- preserve reentry, deregistration, shutdown, and relay regressions
Diffstat:
11 files changed, 164 insertions(+), 270 deletions(-)
diff --git a/app/desktop/src/main/kotlin/org/radroots/studio/application/StudioAppStore.kt b/app/desktop/src/main/kotlin/org/radroots/studio/application/StudioAppStore.kt
@@ -29,11 +29,7 @@ class StudioAppStore(
) : AutoCloseable {
private val mutableState = mutableStateOf(StudioStoreState(snapshot = gateway.snapshot()))
private var closed = false
- private val subscription = gateway.subscribe { snapshot ->
- scope.launch {
- if (!closed) acceptSnapshot(snapshot)
- }
- }
+ private var subscription: AutoCloseable? = null
private var pendingRemoval: RemovalTicket? = null
private var command: Job? = null
@@ -41,7 +37,19 @@ class StudioAppStore(
get() = mutableState
init {
- runCommand { gateway.bootstrap() }
+ launchCommand {
+ val registered = gateway.subscribe { snapshot ->
+ scope.launch {
+ if (!closed) acceptSnapshot(snapshot)
+ }
+ }
+ if (closed) {
+ registered.close()
+ return@launchCommand
+ }
+ subscription = registered
+ acceptSnapshot(gateway.bootstrap())
+ }
}
fun editImportDraft(value: String) {
@@ -175,7 +183,7 @@ class StudioAppStore(
closed = true
command?.cancel()
pendingRemoval?.close()
- subscription.close()
+ subscription?.close()
gateway.close()
}
}
diff --git a/app/desktop/src/main/kotlin/org/radroots/studio/application/StudioCoreGateway.kt b/app/desktop/src/main/kotlin/org/radroots/studio/application/StudioCoreGateway.kt
@@ -12,7 +12,7 @@ interface RemovalTicket : AutoCloseable
interface StudioCoreGateway : AutoCloseable {
fun snapshot(): AppSnapshotDto
- fun subscribe(onSnapshot: (AppSnapshotDto) -> Unit): AutoCloseable
+ suspend fun subscribe(onSnapshot: (AppSnapshotDto) -> Unit): AutoCloseable
suspend fun bootstrap(): AppSnapshotDto
@@ -38,7 +38,7 @@ class NativeStudioCoreGateway(
) : StudioCoreGateway {
override fun snapshot(): AppSnapshotDto = core.snapshot()
- override fun subscribe(onSnapshot: (AppSnapshotDto) -> Unit): AutoCloseable {
+ override suspend fun subscribe(onSnapshot: (AppSnapshotDto) -> Unit): AutoCloseable {
val subscription = core.subscribe(
object : StudioObserver {
override fun onSnapshotChanged(snapshot: AppSnapshotDto) {
diff --git a/app/desktop/src/test/kotlin/org/radroots/studio/application/RadrootsApplicationTest.kt b/app/desktop/src/test/kotlin/org/radroots/studio/application/RadrootsApplicationTest.kt
@@ -75,7 +75,7 @@ private class ApplicationGateway : StudioCoreGateway {
var closed = false
override fun snapshot() = applicationSnapshot(0UL)
- override fun subscribe(onSnapshot: (org.radroots.studio.ffi.AppSnapshotDto) -> Unit) =
+ override suspend fun subscribe(onSnapshot: (org.radroots.studio.ffi.AppSnapshotDto) -> Unit) =
AutoCloseable {}
override suspend fun bootstrap() = applicationSnapshot(1UL)
override suspend fun generateAccount(): org.radroots.studio.ffi.GeneratedAccountDto =
diff --git a/app/desktop/src/test/kotlin/org/radroots/studio/application/StudioAppStoreTest.kt b/app/desktop/src/test/kotlin/org/radroots/studio/application/StudioAppStoreTest.kt
@@ -130,7 +130,7 @@ private class FakeStudioCoreGateway(
override fun snapshot(): AppSnapshotDto = current
- override fun subscribe(onSnapshot: (AppSnapshotDto) -> Unit): AutoCloseable {
+ override suspend fun subscribe(onSnapshot: (AppSnapshotDto) -> Unit): AutoCloseable {
observer = onSnapshot
return AutoCloseable { subscriptionClosed = true }
}
diff --git a/core/crates/application/src/app_core.rs b/core/crates/application/src/app_core.rs
@@ -1,5 +1,5 @@
use std::collections::BTreeMap;
-use std::sync::{Arc, Mutex, MutexGuard};
+use std::sync::{Mutex, MutexGuard};
use radroots_studio_domain::{SafeError, SafeErrorCode, SafeMessage};
@@ -8,30 +8,14 @@ use crate::{
StateMachine, StateTransition,
};
-pub trait AppObserver: Send + Sync {
- fn on_snapshot_changed(&self, snapshot: AppSnapshot);
-}
-
-#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
-pub struct ObserverHandle(u64);
-
pub struct RemovalConfirmationToken {
id: u64,
public_key: radroots_studio_domain::PublicKey,
revision: SnapshotRevision,
}
-impl ObserverHandle {
- #[must_use]
- pub const fn value(self) -> u64 {
- self.0
- }
-}
-
struct CoreState {
state_machine: StateMachine,
- observers: BTreeMap<ObserverHandle, Arc<dyn AppObserver>>,
- next_observer: u64,
removal_tokens: BTreeMap<u64, (radroots_studio_domain::PublicKey, SnapshotRevision)>,
next_removal_token: u64,
}
@@ -48,8 +32,6 @@ impl AppCore {
relay_configuration,
state: Mutex::new(CoreState {
state_machine: StateMachine::booting(),
- observers: BTreeMap::new(),
- next_observer: 1,
removal_tokens: BTreeMap::new(),
next_removal_token: 1,
}),
@@ -98,53 +80,13 @@ impl AppCore {
self.lock_state().state_machine.snapshot().clone()
}
- /// Registers an observer and immediately supplies the current snapshot.
- ///
- /// # Errors
- ///
- /// Returns a safe observer error if the handle space is exhausted.
- pub fn subscribe(&self, observer: Arc<dyn AppObserver>) -> Result<ObserverHandle, SafeError> {
- let (handle, snapshot, observer) = {
- let mut state = self.lock_state();
- let handle = ObserverHandle(state.next_observer);
- state.next_observer = state
- .next_observer
- .checked_add(1)
- .ok_or_else(observer_registration_failed)?;
- let registered = Arc::clone(&observer);
- state.observers.insert(handle, observer);
- (handle, state.state_machine.snapshot().clone(), registered)
- };
- observer.on_snapshot_changed(snapshot);
- Ok(handle)
- }
-
- #[must_use]
- pub fn unsubscribe(&self, handle: ObserverHandle) -> bool {
- self.lock_state().observers.remove(&handle).is_some()
- }
-
pub(crate) fn apply_transition(
&self,
transition: StateTransition,
) -> Result<AppSnapshot, SafeError> {
- let (snapshot, observers) = {
- let mut state = self.lock_state();
- let previous_revision = state.state_machine.snapshot().revision();
- let snapshot = state
- .state_machine
- .apply(transition, &self.relay_configuration)?;
- let observers = if snapshot.revision() == previous_revision {
- Vec::new()
- } else {
- state.observers.values().cloned().collect::<Vec<_>>()
- };
- (snapshot, observers)
- };
- for observer in observers {
- observer.on_snapshot_changed(snapshot.clone());
- }
- Ok(snapshot)
+ self.lock_state()
+ .state_machine
+ .apply(transition, &self.relay_configuration)
}
pub(crate) fn issue_removal_token(
@@ -199,13 +141,6 @@ impl AppCore {
}
}
-const fn observer_registration_failed() -> SafeError {
- SafeError::new(
- SafeErrorCode::ObserverRegistrationFailed,
- SafeMessage::new("The application observer could not be registered."),
- )
-}
-
const fn invalid_application_state() -> SafeError {
SafeError::new(
SafeErrorCode::InvalidApplicationState,
@@ -222,77 +157,27 @@ const fn account_not_found() -> SafeError {
#[cfg(test)]
mod tests {
- use std::sync::{Arc, Mutex, Weak};
-
- use crate::{AppCore, AppLifecycle, AppObserver, AppSnapshot, RelayConfiguration};
-
- #[derive(Default)]
- struct RecordingObserver {
- snapshots: Mutex<Vec<AppSnapshot>>,
- core: Mutex<Option<Weak<AppCore>>>,
- }
-
- impl AppObserver for RecordingObserver {
- fn on_snapshot_changed(&self, snapshot: AppSnapshot) {
- if let Some(core) = self
- .core
- .lock()
- .expect("observer core lock")
- .as_ref()
- .and_then(Weak::upgrade)
- {
- assert_eq!(core.snapshot().revision(), snapshot.revision());
- }
- self.snapshots
- .lock()
- .expect("observer snapshots lock")
- .push(snapshot);
- }
- }
+ use crate::{AppCore, AppLifecycle, RelayConfiguration};
#[test]
- fn bootstrap_publishes_ready_snapshot_outside_the_core_lock() {
- let core = Arc::new(AppCore::in_memory(RelayConfiguration::default()));
- let observer = Arc::new(RecordingObserver::default());
- *observer.core.lock().expect("observer core lock") = Some(Arc::downgrade(&core));
-
- let handle = core.subscribe(observer.clone()).expect("register observer");
+ fn bootstrap_is_idempotent_and_advances_only_once() {
+ let core = AppCore::in_memory(RelayConfiguration::default());
let ready = core.bootstrap().expect("bootstrap");
let repeated = core.bootstrap().expect("idempotent bootstrap");
- assert_eq!(handle.value(), 1);
assert_eq!(ready.lifecycle(), AppLifecycle::Ready);
assert_eq!(ready.revision().value(), 1);
assert_eq!(repeated, ready);
- assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 2);
- }
-
- #[test]
- fn deregistered_observer_receives_no_later_bootstrap_update() {
- let core = AppCore::in_memory(RelayConfiguration::default());
- let observer = Arc::new(RecordingObserver::default());
- let handle = core.subscribe(observer.clone()).expect("register observer");
-
- assert!(core.unsubscribe(handle));
- assert!(!core.unsubscribe(handle));
- core.bootstrap().expect("bootstrap");
-
- assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1);
}
#[test]
- fn core_instances_never_share_state_or_observers() {
+ fn core_instances_never_share_state() {
let first = AppCore::in_memory(RelayConfiguration::default());
let second = AppCore::in_memory(RelayConfiguration::default());
- let observer = Arc::new(RecordingObserver::default());
- first
- .subscribe(observer.clone())
- .expect("register observer");
first.bootstrap().expect("first bootstrap");
assert_eq!(first.snapshot().revision().value(), 1);
assert_eq!(second.snapshot().revision().value(), 0);
- assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 2);
}
}
diff --git a/core/crates/application/src/lib.rs b/core/crates/application/src/lib.rs
@@ -23,7 +23,7 @@ pub use actor::{
CommandSubmission, CommandTicket, LifecycleGate, RequestId, RuntimeCommandClass,
RuntimeLifecycle, SessionGeneration, TaskCorrelation,
};
-pub use app_core::{AppCore, AppObserver, ObserverHandle, RemovalConfirmationToken};
+pub use app_core::{AppCore, RemovalConfirmationToken};
pub use change_stream::{
ChangeSubscriptionId, OrderedSnapshotChanges, SnapshotChange, SnapshotChangeReceiver,
};
diff --git a/core/crates/application/src/profile_refresh.rs b/core/crates/application/src/profile_refresh.rs
@@ -192,7 +192,7 @@ const fn refresh_status(error: SafeError) -> ProfileRefreshStatus {
#[cfg(test)]
mod tests {
- use std::sync::{Arc, Mutex};
+ use std::sync::Mutex;
use radroots_studio_domain::{
EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, RelayUrl, SafeError,
@@ -200,7 +200,7 @@ mod tests {
};
use crate::{
- AppCore, AppObserver, BoxFuture, CachedProfile, Clock, InMemoryAccountRepository,
+ AppCore, BoxFuture, CachedProfile, Clock, InMemoryAccountRepository,
InMemoryOperationJournal, InMemorySecretStore, NostrClient, ProfileLoadState,
ProfileRefreshStatus, ProfileRepository, RelayConfiguration, RelayConnectionState,
};
@@ -283,19 +283,6 @@ mod tests {
}
}
- #[derive(Default)]
- struct States(Mutex<Vec<(ProfileLoadState, RelayConnectionState)>>);
- impl AppObserver for States {
- fn on_snapshot_changed(&self, snapshot: crate::AppSnapshot) {
- if let Some(active) = snapshot.active_account() {
- self.0
- .lock()
- .expect("states")
- .push((active.profile_state(), active.relay_state()));
- }
- }
- }
-
fn profile(public_key: PublicKey, name: &str, timestamp: i64) -> Kind0ProfileCandidate {
Kind0ProfileCandidate::new(
EventId::from_bytes([u8::try_from(timestamp).expect("small timestamp"); 32]),
@@ -350,29 +337,36 @@ mod tests {
}
#[tokio::test]
- async fn refresh_emits_cache_before_loading_and_fresh_profile() {
+ async fn refresh_transitions_from_cache_through_loading_to_fresh_profile() {
let profiles = MemoryProfiles::default();
let (core, public_key) = active_core(&profiles, Some("Cached"));
- let states = Arc::new(States::default());
- core.subscribe(states.clone()).expect("observer");
- core.refresh_profile_for_active_account(
- &profiles,
- &FixedClient(Ok(Some(profile(public_key, "Fresh", 20)))),
- &FixedClock,
- )
- .await
- .expect("refresh");
-
- let observed = states.0.lock().expect("states");
assert_eq!(
- observed.first().map(|state| state.0),
+ core.snapshot()
+ .active_account()
+ .map(crate::ActiveAccountSnapshot::profile_state),
Some(ProfileLoadState::Cached)
);
- assert!(observed.contains(&(ProfileLoadState::Loading, RelayConnectionState::Connecting)));
+ let plan = core
+ .begin_profile_refresh()
+ .expect("begin refresh")
+ .expect("active refresh");
+ let loading = core.snapshot();
+ assert_eq!(
+ loading
+ .active_account()
+ .map(crate::ActiveAccountSnapshot::profile_state),
+ Some(ProfileLoadState::Loading)
+ );
assert_eq!(
- observed.last(),
- Some(&(ProfileLoadState::Fresh, RelayConnectionState::Connected))
+ loading
+ .active_account()
+ .map(crate::ActiveAccountSnapshot::relay_state),
+ Some(RelayConnectionState::Connecting)
);
+ let client = FixedClient(Ok(Some(profile(public_key, "Fresh", 20))));
+ let result = client.fetch_profile(plan.public_key(), plan.relays()).await;
+ core.complete_profile_refresh(&plan, result, &profiles, &FixedClock)
+ .expect("complete refresh");
assert_eq!(
core.snapshot()
.active_account()
diff --git a/core/crates/ffi/src/commands.rs b/core/crates/ffi/src/commands.rs
@@ -1,4 +1,4 @@
-use std::collections::BTreeSet;
+use std::collections::BTreeMap;
use std::fmt::{self, Display, Formatter};
use std::num::NonZeroUsize;
use std::path::{Path, PathBuf};
@@ -68,7 +68,9 @@ impl RemovalRequest {
pub(crate) struct RuntimeCore {
pub(crate) actor: RuntimeActorHandle,
- pub(crate) observers: Mutex<BTreeSet<radroots_studio_application::ObserverHandle>>,
+ pub(crate) observers: Mutex<
+ BTreeMap<radroots_studio_application::ChangeSubscriptionId, tokio::task::JoinHandle<()>>,
+ >,
pub(crate) closed: AtomicBool,
}
@@ -278,7 +280,7 @@ impl StudioAppCore {
Ok(Arc::new(Self {
inner: Arc::new(RuntimeCore {
actor,
- observers: Mutex::new(BTreeSet::new()),
+ observers: Mutex::new(BTreeMap::new()),
closed: AtomicBool::new(false),
}),
}))
@@ -366,7 +368,7 @@ mod tests {
Arc::new(StudioAppCore {
inner: Arc::new(RuntimeCore {
actor,
- observers: std::sync::Mutex::new(std::collections::BTreeSet::new()),
+ observers: std::sync::Mutex::new(std::collections::BTreeMap::new()),
closed: std::sync::atomic::AtomicBool::new(false),
}),
})
diff --git a/core/crates/ffi/src/observer.rs b/core/crates/ffi/src/observer.rs
@@ -1,48 +1,51 @@
+use std::num::NonZeroUsize;
use std::sync::atomic::Ordering;
use std::sync::{Arc, Mutex, Weak};
-use radroots_studio_application::{AppObserver, AppSnapshot, ObserverHandle};
+use radroots_studio_application::ChangeSubscriptionId;
use crate::commands::RuntimeCore;
use crate::{AppSnapshotDto, StudioAppCore, StudioError};
+const OBSERVER_CHANGE_CAPACITY: NonZeroUsize = match NonZeroUsize::new(64) {
+ Some(capacity) => capacity,
+ None => unreachable!(),
+};
+
#[uniffi::export(callback_interface)]
pub trait StudioObserver: Send + Sync {
fn on_snapshot_changed(&self, snapshot: AppSnapshotDto);
}
-struct ObserverBridge {
- observer: Arc<dyn StudioObserver>,
-}
-
-impl AppObserver for ObserverBridge {
- fn on_snapshot_changed(&self, snapshot: AppSnapshot) {
- self.observer.on_snapshot_changed((&snapshot).into());
- }
-}
-
#[derive(uniffi::Object)]
pub struct ObserverSubscription {
core: Weak<RuntimeCore>,
- handle: Mutex<Option<ObserverHandle>>,
+ id: Mutex<Option<ChangeSubscriptionId>>,
}
#[uniffi::export]
impl ObserverSubscription {
pub fn unsubscribe(&self) {
- let handle = self
- .handle
+ let id = self
+ .id
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take();
- let (Some(core), Some(handle)) = (self.core.upgrade(), handle) else {
+ let (Some(core), Some(id)) = (self.core.upgrade(), id) else {
return;
};
- let _ = core.actor.unsubscribe(handle);
- core.observers
+ if let Some(task) = core
+ .observers
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
- .remove(&handle);
+ .remove(&id)
+ {
+ task.abort();
+ }
+ let actor = core.actor.clone();
+ crate::commands::runtime().spawn(async move {
+ let _ = actor.unsubscribe_changes(id).await;
+ });
}
}
@@ -59,29 +62,34 @@ impl StudioAppCore {
/// # Errors
///
/// Returns a safe observer or lifecycle error.
- pub fn subscribe(
+ pub async fn subscribe(
&self,
observer: Box<dyn StudioObserver>,
) -> Result<Arc<ObserverSubscription>, StudioError> {
if self.inner.closed.load(Ordering::Acquire) {
return Err(closed_error());
}
- let bridge = Arc::new(ObserverBridge {
- observer: Arc::from(observer),
- });
- let handle = self
+ let mut subscription = self
.inner
.actor
- .subscribe(bridge)
+ .subscribe_changes(OBSERVER_CHANGE_CAPACITY)
+ .await
.map_err(StudioError::from)?;
+ let id = subscription.id();
+ let observer: Arc<dyn StudioObserver> = Arc::from(observer);
+ let task = crate::commands::runtime().spawn(async move {
+ while let Some(change) = subscription.receive().await {
+ observer.on_snapshot_changed(change.snapshot().into());
+ }
+ });
self.inner
.observers
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
- .insert(handle);
+ .insert(id, task);
Ok(Arc::new(ObserverSubscription {
core: Arc::downgrade(&self.inner),
- handle: Mutex::new(Some(handle)),
+ id: Mutex::new(Some(id)),
}))
}
@@ -96,8 +104,8 @@ impl StudioAppCore {
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
);
- for handle in handles {
- let _ = self.inner.actor.unsubscribe(handle);
+ for (_, task) in handles {
+ task.abort();
}
let actor = self.inner.actor.clone();
crate::commands::runtime().spawn(async move {
@@ -163,7 +171,7 @@ mod tests {
Arc::new(StudioAppCore {
inner: Arc::new(RuntimeCore {
actor,
- observers: Mutex::new(std::collections::BTreeSet::new()),
+ observers: Mutex::new(std::collections::BTreeMap::new()),
closed: std::sync::atomic::AtomicBool::new(false),
}),
})
@@ -171,37 +179,44 @@ mod tests {
#[test]
fn callbacks_allow_reentry_and_stop_after_subscription_close() {
- let core = core();
- let observer = Arc::new(RecordingObserver::default());
- *observer.core.lock().expect("core") = Some(Arc::clone(&core));
- let subscription = core
- .subscribe(Box::new(ArcObserver(observer.clone())))
- .expect("subscribe");
+ runtime().block_on(async {
+ let core = core();
+ let observer = Arc::new(RecordingObserver::default());
+ *observer.core.lock().expect("core") = Some(Arc::clone(&core));
+ let subscription = core
+ .subscribe(Box::new(ArcObserver(observer.clone())))
+ .await
+ .expect("subscribe");
- assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1);
- runtime()
- .block_on(core.inner.actor.bootstrap())
- .expect("idempotent bootstrap");
- assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1);
- subscription.unsubscribe();
- runtime()
- .block_on(core.inner.actor.sign_out())
- .expect("sign out");
- assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1);
+ wait_for_snapshot_count(&observer, 1).await;
+ core.inner
+ .actor
+ .bootstrap()
+ .await
+ .expect("idempotent bootstrap");
+ assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1);
+ subscription.unsubscribe();
+ core.inner.actor.sign_out().await.expect("sign out");
+ assert_eq!(observer.snapshots.lock().expect("snapshots").len(), 1);
+ });
}
#[test]
fn core_close_deregisters_all_observers_and_rejects_new_subscriptions() {
let core = core();
let observer = Arc::new(RecordingObserver::default());
- let _subscription = core
- .subscribe(Box::new(ArcObserver(observer.clone())))
+ let _subscription = runtime()
+ .block_on(core.subscribe(Box::new(ArcObserver(observer.clone()))))
.expect("subscribe");
core.shutdown();
core.shutdown();
- assert!(core.subscribe(Box::new(ArcObserver(observer))).is_err());
+ assert!(
+ runtime()
+ .block_on(core.subscribe(Box::new(ArcObserver(observer))))
+ .is_err()
+ );
assert!(core.inner.observers.lock().expect("observers").is_empty());
}
@@ -231,6 +246,7 @@ mod tests {
*observer.core.lock().expect("core") = Some(Arc::clone(&core));
let subscription = core
.subscribe(Box::new(ArcObserver(observer.clone())))
+ .await
.expect("subscribe");
let imported = core
.import_secret_key(SECRET_HEX.to_owned())
@@ -240,6 +256,7 @@ mod tests {
core.activate_account(public_key).await.expect("activate");
core.refresh_active_profile().await.expect("refresh");
+ wait_for_fresh_profile(&observer).await;
let snapshots = observer.snapshots.lock().expect("snapshots").clone();
assert!(snapshots.iter().any(|snapshot| {
snapshot.active_account.as_ref().is_some_and(|active| {
@@ -268,4 +285,37 @@ mod tests {
self.0.on_snapshot_changed(snapshot);
}
}
+
+ async fn wait_for_snapshot_count(observer: &RecordingObserver, minimum: usize) {
+ tokio::time::timeout(Duration::from_secs(1), async {
+ while observer.snapshots.lock().expect("snapshots").len() < minimum {
+ tokio::task::yield_now().await;
+ }
+ })
+ .await
+ .expect("snapshot delivery");
+ }
+
+ async fn wait_for_fresh_profile(observer: &RecordingObserver) {
+ tokio::time::timeout(Duration::from_secs(1), async {
+ loop {
+ let fresh = observer
+ .snapshots
+ .lock()
+ .expect("snapshots")
+ .iter()
+ .any(|snapshot| {
+ snapshot.active_account.as_ref().is_some_and(|active| {
+ active.profile_state == ProfileLoadStateDto::Fresh
+ })
+ });
+ if fresh {
+ break;
+ }
+ tokio::task::yield_now().await;
+ }
+ })
+ .await
+ .expect("fresh profile delivery");
+ }
}
diff --git a/core/crates/storage/src/runtime_actor.rs b/core/crates/storage/src/runtime_actor.rs
@@ -6,12 +6,11 @@ use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use radroots_studio_application::{
- ActorMailbox, AppObserver, AppSnapshot, ChangeSubscriptionId, Clock, CommandContext,
- CommandEnvelope, CommandReceipt, CommandResult, CommandSubmission, GenerateAccountReceipt,
- ImportAccountReceipt, LifecycleGate, NostrClient, ObserverHandle, OrderedSnapshotChanges,
- ProfileRefreshPlan, RelayConfiguration, RemovalConfirmationToken, RequestId,
- RuntimeCommandClass, RuntimeLifecycle, SecretStore, SessionGeneration, SnapshotChange,
- SnapshotChangeReceiver, SnapshotRevision, TaskCorrelation,
+ ActorMailbox, AppSnapshot, ChangeSubscriptionId, Clock, CommandContext, CommandEnvelope,
+ CommandReceipt, CommandResult, CommandSubmission, GenerateAccountReceipt, ImportAccountReceipt,
+ LifecycleGate, NostrClient, OrderedSnapshotChanges, ProfileRefreshPlan, RelayConfiguration,
+ RemovalConfirmationToken, RequestId, RuntimeCommandClass, RuntimeLifecycle, SecretStore,
+ SessionGeneration, SnapshotChange, SnapshotChangeReceiver, SnapshotRevision, TaskCorrelation,
};
use radroots_studio_domain::{
Kind0ProfileCandidate, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput,
@@ -228,20 +227,6 @@ impl RuntimeActorHandle {
self.adapter.core().snapshot()
}
- /// Registers a read-only snapshot observer.
- ///
- /// # Errors
- ///
- /// Returns a safe observer-registration error.
- pub fn subscribe(&self, observer: Arc<dyn AppObserver>) -> Result<ObserverHandle, SafeError> {
- self.adapter.core().subscribe(observer)
- }
-
- #[must_use]
- pub fn unsubscribe(&self, handle: ObserverHandle) -> bool {
- self.adapter.core().unsubscribe(handle)
- }
-
/// Returns the ready snapshot through the actor command boundary.
///
/// # Errors
diff --git a/core/crates/storage/tests/local_relay_e2e.rs b/core/crates/storage/tests/local_relay_e2e.rs
@@ -1,12 +1,11 @@
-use std::sync::{Arc, Mutex};
use std::time::Duration;
use nostr::{EventBuilder, Keys, Metadata};
use nostr_relay_builder::MockRelay;
use nostr_sdk::Client;
use radroots_studio_application::{
- AppObserver, Clock, InMemorySecretStore, ProfileLoadState, ProfileRepository,
- RelayConfiguration, RelayConnectionState, SdkNostrClient, SecretStore, SessionState,
+ Clock, InMemorySecretStore, ProfileLoadState, ProfileRepository, RelayConfiguration,
+ RelayConnectionState, SdkNostrClient, SecretStore, SessionState,
};
use radroots_studio_domain::{RelayUrl, SecretKeyInput, UnixTimestamp};
use radroots_studio_storage::PersistentAppCore;
@@ -21,15 +20,6 @@ impl Clock for FixedClock {
}
}
-#[derive(Default)]
-struct RecordingObserver(Mutex<Vec<radroots_studio_application::AppSnapshot>>);
-
-impl AppObserver for RecordingObserver {
- fn on_snapshot_changed(&self, snapshot: radroots_studio_application::AppSnapshot) {
- self.0.lock().expect("observer lock").push(snapshot);
- }
-}
-
#[tokio::test]
async fn local_relay_e2e_imports_activates_refreshes_and_caches_profile() {
let local_relay = MockRelay::run().await.expect("local relay");
@@ -70,11 +60,6 @@ async fn local_relay_e2e_imports_activates_refreshes_and_caches_profile() {
.activate_account(public_key, &secrets, &FixedClock)
.expect("activate account");
- let observer = Arc::new(RecordingObserver::default());
- let handle = adapter
- .core()
- .subscribe(observer.clone())
- .expect("observer");
let refreshed = adapter
.core()
.refresh_active_profile(
@@ -102,21 +87,6 @@ async fn local_relay_e2e_imports_activates_refreshes_and_caches_profile() {
cached.candidate().metadata().preferred_name(),
Some("Farm Account")
);
- let recorded_snapshots = observer.0.lock().expect("observer snapshots").clone();
- assert!(recorded_snapshots.iter().any(|snapshot| {
- snapshot.active_account().is_some_and(|account| {
- account.profile_state() == ProfileLoadState::Loading
- && account.relay_state() == RelayConnectionState::Connecting
- })
- }));
- assert_eq!(
- recorded_snapshots
- .last()
- .map(radroots_studio_application::AppSnapshot::revision),
- Some(refreshed.revision())
- );
- assert!(adapter.core().unsubscribe(handle));
-
let public_debug = format!("{refreshed:?}");
assert!(!public_debug.contains(SECRET_HEX));
assert!(!public_debug.contains("nsec1"));