lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 53cc1d73cdb51ae268ecc7224d398c2ebf31d569
parent 4f8c9958c2046e6a4e70d7f4efe0e850ca639a42
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 21:55:46 +0000

runtime: make change subscription atomic

- register subscriptions through the serialized actor command boundary
- deliver the current full snapshot before any later revision
- retain predecessor revisions so consumers can detect and recover gaps
- cover initial delivery ordering, mutation delivery, and unsubscribe

Diffstat:
Mcrates/studio_application/src/change_stream.rs | 64+++++++++++++++++++++++++++++++++++++++++++++++-----------------
Mcrates/studio_storage/src/runtime_actor.rs | 117+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
2 files changed, 156 insertions(+), 25 deletions(-)

diff --git a/crates/studio_application/src/change_stream.rs b/crates/studio_application/src/change_stream.rs @@ -18,6 +18,7 @@ impl ChangeSubscriptionId { #[derive(Clone, Debug, Eq, PartialEq)] pub struct SnapshotChange { snapshot: AppSnapshot, + previous_revision: Option<SnapshotRevision>, } impl SnapshotChange { @@ -35,6 +36,17 @@ impl SnapshotChange { pub fn into_snapshot(self) -> AppSnapshot { self.snapshot } + + #[must_use] + pub const fn previous_revision(&self) -> Option<SnapshotRevision> { + self.previous_revision + } + + #[must_use] + pub fn recovers_gap_after(&self, observed: SnapshotRevision) -> bool { + self.previous_revision + .is_some_and(|previous| previous != observed) + } } pub struct SnapshotChangeReceiver { @@ -48,16 +60,16 @@ impl SnapshotChangeReceiver { } pub struct OrderedSnapshotChanges { - last_revision: SnapshotRevision, + latest: AppSnapshot, next_subscription: u64, subscribers: BTreeMap<ChangeSubscriptionId, mpsc::Sender<SnapshotChange>>, } impl OrderedSnapshotChanges { #[must_use] - pub fn new(initial_revision: SnapshotRevision) -> Self { + pub fn new(initial_snapshot: AppSnapshot) -> Self { Self { - last_revision: initial_revision, + latest: initial_snapshot, next_subscription: 1, subscribers: BTreeMap::new(), } @@ -65,7 +77,7 @@ impl OrderedSnapshotChanges { #[must_use] pub const fn last_revision(&self) -> SnapshotRevision { - self.last_revision + self.latest.revision() } /// Registers a bounded consumer for future changes. @@ -80,6 +92,12 @@ impl OrderedSnapshotChanges { let id = ChangeSubscriptionId(NonZeroU64::new(self.next_subscription)?); self.next_subscription = self.next_subscription.checked_add(1)?; let (sender, receiver) = mpsc::channel(capacity.get()); + sender + .try_send(SnapshotChange { + snapshot: self.latest.clone(), + previous_revision: None, + }) + .ok()?; self.subscribers.insert(id, sender); Some((id, SnapshotChangeReceiver { receiver })) } @@ -90,11 +108,14 @@ impl OrderedSnapshotChanges { } pub fn publish(&mut self, snapshot: AppSnapshot) { - if snapshot.revision() <= self.last_revision { + if snapshot.revision() <= self.latest.revision() { return; } - self.last_revision = snapshot.revision(); - let change = SnapshotChange { snapshot }; + let change = SnapshotChange { + previous_revision: Some(self.latest.revision()), + snapshot, + }; + self.latest = change.snapshot.clone(); self.subscribers .retain(|_, sender| match sender.try_send(change.clone()) { Ok(()) | Err(mpsc::error::TrySendError::Full(_)) => true, @@ -113,7 +134,7 @@ mod tests { #[tokio::test] async fn change_stream_publishes_monotonic_revisions_to_multiple_consumers() { - let mut changes = OrderedSnapshotChanges::new(revision(0)); + let mut changes = OrderedSnapshotChanges::new(snapshot(0)); let (_, mut first) = changes .subscribe(NonZeroUsize::new(4).expect("capacity")) .expect("first subscription"); @@ -121,6 +142,14 @@ mod tests { .subscribe(NonZeroUsize::new(4).expect("capacity")) .expect("second subscription"); + assert_eq!( + first.receive().await.expect("initial").revision(), + revision(0) + ); + assert_eq!( + second.receive().await.expect("initial").revision(), + revision(0) + ); changes.publish(snapshot(1)); changes.publish(snapshot(1)); changes.publish(snapshot(2)); @@ -140,22 +169,23 @@ mod tests { #[tokio::test] async fn slow_consumers_expose_a_revision_gap_without_blocking_publication() { - let mut changes = OrderedSnapshotChanges::new(revision(0)); + let mut changes = OrderedSnapshotChanges::new(snapshot(0)); let (_, mut receiver) = changes .subscribe(NonZeroUsize::new(1).expect("capacity")) .expect("subscription"); - changes.publish(snapshot(1)); - changes.publish(snapshot(2)); assert_eq!( - receiver.receive().await.expect("first").revision(), - revision(1) + receiver.receive().await.expect("initial").revision(), + revision(0) ); + changes.publish(snapshot(1)); + changes.publish(snapshot(2)); + let first = receiver.receive().await.expect("first"); + assert_eq!(first.revision(), revision(1)); changes.publish(snapshot(3)); - assert_eq!( - receiver.receive().await.expect("gap").revision(), - revision(3) - ); + let recovered = receiver.receive().await.expect("gap recovery"); + assert_eq!(recovered.revision(), revision(3)); + assert!(recovered.recovers_gap_after(first.revision())); } fn revision(value: u64) -> SnapshotRevision { diff --git a/crates/studio_storage/src/runtime_actor.rs b/crates/studio_storage/src/runtime_actor.rs @@ -6,11 +6,12 @@ use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use radroots_studio_application::{ - ActorMailbox, AppObserver, AppSnapshot, Clock, CommandContext, CommandEnvelope, CommandReceipt, - CommandResult, CommandSubmission, GenerateAccountReceipt, ImportAccountReceipt, LifecycleGate, - NostrClient, ObserverHandle, OrderedSnapshotChanges, ProfileRefreshPlan, RelayConfiguration, - RemovalConfirmationToken, RequestId, RuntimeCommandClass, RuntimeLifecycle, SecretStore, - SessionGeneration, SnapshotRevision, TaskCorrelation, + 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, }; use radroots_studio_domain::{ Kind0ProfileCandidate, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, @@ -33,6 +34,8 @@ enum RuntimeCommand { RefreshActiveProfile, RequestAccountRemoval(PublicKey), ConfirmAccountRemoval(RemovalConfirmationToken), + SubscribeChanges(NonZeroUsize), + UnsubscribeChanges(ChangeSubscriptionId), Close, } @@ -41,13 +44,17 @@ enum RuntimeCommandValue { Generated(GenerateAccountReceipt), Imported(ImportAccountReceipt), RemovalRequest(RemovalConfirmationToken), + Subscription(RuntimeChangeSubscription), + Unsubscribed(bool), Closed, } impl RuntimeCommand { const fn class(&self) -> RuntimeCommandClass { match self { - Self::Snapshot => RuntimeCommandClass::Observe, + Self::Snapshot | Self::SubscribeChanges(_) | Self::UnsubscribeChanges(_) => { + RuntimeCommandClass::Observe + } Self::GenerateAccount | Self::ImportSecretKey(_) | Self::ActivateAccount(_) @@ -95,6 +102,22 @@ pub struct RuntimeActorHandle { session_generation: Arc<AtomicU64>, } +pub struct RuntimeChangeSubscription { + id: ChangeSubscriptionId, + receiver: SnapshotChangeReceiver, +} + +impl RuntimeChangeSubscription { + #[must_use] + pub const fn id(&self) -> ChangeSubscriptionId { + self.id + } + + pub async fn receive(&mut self) -> Option<SnapshotChange> { + self.receiver.receive().await + } +} + impl RuntimeActorHandle { /// Opens, migrates, recovers, and starts one actor-owned file-backed runtime. /// @@ -164,7 +187,7 @@ impl RuntimeActorHandle { let lifecycle = Arc::new(Mutex::new(gate)); let (mailbox, receiver) = ActorMailbox::bounded(capacity); let session_generation = Arc::new(AtomicU64::new(SessionGeneration::initial().value())); - let changes = OrderedSnapshotChanges::new(adapter.core().snapshot().revision()); + let changes = OrderedSnapshotChanges::new(adapter.core().snapshot()); let actor = RuntimeActor { adapter: Arc::clone(&adapter), secrets, @@ -349,6 +372,39 @@ impl RuntimeActorHandle { } } + /// Atomically registers a bounded ordered change consumer with its initial snapshot. + /// + /// # Errors + /// + /// Returns a safe actor or subscription error. + pub async fn subscribe_changes( + &self, + capacity: NonZeroUsize, + ) -> Result<RuntimeChangeSubscription, SafeError> { + match self + .dispatch(RuntimeCommand::SubscribeChanges(capacity), None) + .await? + { + RuntimeCommandValue::Subscription(subscription) => Ok(subscription), + _ => Err(invalid_actor_response()), + } + } + + /// Removes a change consumer through the serialized actor boundary. + /// + /// # Errors + /// + /// Returns a safe actor error. + pub async fn unsubscribe_changes(&self, id: ChangeSubscriptionId) -> Result<bool, SafeError> { + match self + .dispatch(RuntimeCommand::UnsubscribeChanges(id), None) + .await? + { + RuntimeCommandValue::Unsubscribed(removed) => Ok(removed), + _ => Err(invalid_actor_response()), + } + } + async fn dispatch( &self, command: RuntimeCommand, @@ -511,7 +567,7 @@ impl RuntimeActor { None } - fn execute_sync(&self, command: RuntimeCommand) -> CommandResult<RuntimeCommandValue> { + fn execute_sync(&mut self, command: RuntimeCommand) -> CommandResult<RuntimeCommandValue> { let result = match command { RuntimeCommand::Snapshot => Ok(RuntimeCommandValue::Snapshot(Box::new( self.adapter.core().snapshot(), @@ -548,6 +604,16 @@ impl RuntimeActor { .confirm_account_removal(token, self.secrets.as_ref(), self.clock.as_ref()) .map(Box::new) .map(RuntimeCommandValue::Snapshot), + RuntimeCommand::SubscribeChanges(capacity) => self + .changes + .subscribe(capacity) + .map(|(id, receiver)| { + RuntimeCommandValue::Subscription(RuntimeChangeSubscription { id, receiver }) + }) + .ok_or_else(observer_registration_failed), + RuntimeCommand::UnsubscribeChanges(id) => Ok(RuntimeCommandValue::Unsubscribed( + self.changes.unsubscribe(id), + )), RuntimeCommand::Close | RuntimeCommand::RefreshActiveProfile => { Err(invalid_actor_response()) } @@ -738,6 +804,13 @@ const fn invalid_actor_response() -> SafeError { ) } +const fn observer_registration_failed() -> SafeError { + SafeError::new( + SafeErrorCode::ObserverRegistrationFailed, + SafeMessage::new("The application change subscription could not be registered."), + ) +} + #[cfg(test)] mod tests { use std::num::NonZeroUsize; @@ -1074,6 +1147,34 @@ mod tests { } } + #[tokio::test(flavor = "multi_thread")] + async fn actor_subscription_atomically_delivers_initial_then_ordered_changes() { + let (actor, _) = actor(); + let mut subscription = actor + .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) + .await + .expect("subscribe"); + let initial = subscription.receive().await.expect("initial snapshot"); + assert_eq!(initial.revision(), actor.snapshot().revision()); + assert!(initial.previous_revision().is_none()); + + actor + .import_secret_key(secret( + "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", + )) + .await + .expect("import"); + let changed = subscription.receive().await.expect("change"); + assert!(changed.revision() > initial.revision()); + assert_eq!(changed.previous_revision(), Some(initial.revision())); + assert!( + actor + .unsubscribe_changes(subscription.id()) + .await + .expect("unsubscribe") + ); + } + fn secret(value: &str) -> SecretKeyInput { SecretKeyInput::parse(value.to_owned()).expect("valid test secret") }