commit 410566bba0fd7bd2c840dbeefe117b46c1e39e64
parent a68f55e4670afc9ed1cf57ebc108bb9b41f00707
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:
2 files changed, 156 insertions(+), 25 deletions(-)
diff --git a/core/crates/application/src/change_stream.rs b/core/crates/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/core/crates/storage/src/runtime_actor.rs b/core/crates/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")
}