commit 4f8c9958c2046e6a4e70d7f4efe0e850ca639a42
parent cf67b4594b67907805d95abe0c4a76536dbc4dde
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 21:52:42 +0000
runtime: add ordered snapshot change stream
- publish actor-owned changes only for advancing snapshot revisions
- deliver bounded ordered events to independent runtime consumers
- preserve explicit revision gaps for slow-consumer recovery
- cover monotonic multi-consumer delivery and nonblocking overflow
Diffstat:
3 files changed, 197 insertions(+), 3 deletions(-)
diff --git a/crates/studio_application/src/change_stream.rs b/crates/studio_application/src/change_stream.rs
@@ -0,0 +1,181 @@
+use std::collections::BTreeMap;
+use std::num::{NonZeroU64, NonZeroUsize};
+
+use tokio::sync::mpsc;
+
+use crate::{AppSnapshot, SnapshotRevision};
+
+#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
+pub struct ChangeSubscriptionId(NonZeroU64);
+
+impl ChangeSubscriptionId {
+ #[must_use]
+ pub const fn value(self) -> u64 {
+ self.0.get()
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SnapshotChange {
+ snapshot: AppSnapshot,
+}
+
+impl SnapshotChange {
+ #[must_use]
+ pub const fn revision(&self) -> SnapshotRevision {
+ self.snapshot.revision()
+ }
+
+ #[must_use]
+ pub const fn snapshot(&self) -> &AppSnapshot {
+ &self.snapshot
+ }
+
+ #[must_use]
+ pub fn into_snapshot(self) -> AppSnapshot {
+ self.snapshot
+ }
+}
+
+pub struct SnapshotChangeReceiver {
+ receiver: mpsc::Receiver<SnapshotChange>,
+}
+
+impl SnapshotChangeReceiver {
+ pub async fn receive(&mut self) -> Option<SnapshotChange> {
+ self.receiver.recv().await
+ }
+}
+
+pub struct OrderedSnapshotChanges {
+ last_revision: SnapshotRevision,
+ next_subscription: u64,
+ subscribers: BTreeMap<ChangeSubscriptionId, mpsc::Sender<SnapshotChange>>,
+}
+
+impl OrderedSnapshotChanges {
+ #[must_use]
+ pub fn new(initial_revision: SnapshotRevision) -> Self {
+ Self {
+ last_revision: initial_revision,
+ next_subscription: 1,
+ subscribers: BTreeMap::new(),
+ }
+ }
+
+ #[must_use]
+ pub const fn last_revision(&self) -> SnapshotRevision {
+ self.last_revision
+ }
+
+ /// Registers a bounded consumer for future changes.
+ ///
+ /// # Errors
+ ///
+ /// Returns `None` if the subscription identifier space is exhausted.
+ pub fn subscribe(
+ &mut self,
+ capacity: NonZeroUsize,
+ ) -> Option<(ChangeSubscriptionId, SnapshotChangeReceiver)> {
+ let id = ChangeSubscriptionId(NonZeroU64::new(self.next_subscription)?);
+ self.next_subscription = self.next_subscription.checked_add(1)?;
+ let (sender, receiver) = mpsc::channel(capacity.get());
+ self.subscribers.insert(id, sender);
+ Some((id, SnapshotChangeReceiver { receiver }))
+ }
+
+ #[must_use]
+ pub fn unsubscribe(&mut self, id: ChangeSubscriptionId) -> bool {
+ self.subscribers.remove(&id).is_some()
+ }
+
+ pub fn publish(&mut self, snapshot: AppSnapshot) {
+ if snapshot.revision() <= self.last_revision {
+ return;
+ }
+ self.last_revision = snapshot.revision();
+ let change = SnapshotChange { snapshot };
+ self.subscribers
+ .retain(|_, sender| match sender.try_send(change.clone()) {
+ Ok(()) | Err(mpsc::error::TrySendError::Full(_)) => true,
+ Err(mpsc::error::TrySendError::Closed(_)) => false,
+ });
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use std::num::NonZeroUsize;
+
+ use crate::{
+ AppSnapshot, OrderedSnapshotChanges, RelayConfiguration, SessionState, SnapshotRevision,
+ };
+
+ #[tokio::test]
+ async fn change_stream_publishes_monotonic_revisions_to_multiple_consumers() {
+ let mut changes = OrderedSnapshotChanges::new(revision(0));
+ let (_, mut first) = changes
+ .subscribe(NonZeroUsize::new(4).expect("capacity"))
+ .expect("first subscription");
+ let (_, mut second) = changes
+ .subscribe(NonZeroUsize::new(4).expect("capacity"))
+ .expect("second subscription");
+
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(2));
+
+ for receiver in [&mut first, &mut second] {
+ assert_eq!(
+ receiver.receive().await.expect("revision 1").revision(),
+ revision(1)
+ );
+ assert_eq!(
+ receiver.receive().await.expect("revision 2").revision(),
+ revision(2)
+ );
+ }
+ assert_eq!(changes.last_revision(), revision(2));
+ }
+
+ #[tokio::test]
+ async fn slow_consumers_expose_a_revision_gap_without_blocking_publication() {
+ let mut changes = OrderedSnapshotChanges::new(revision(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)
+ );
+ changes.publish(snapshot(3));
+ assert_eq!(
+ receiver.receive().await.expect("gap").revision(),
+ revision(3)
+ );
+ }
+
+ fn revision(value: u64) -> SnapshotRevision {
+ SnapshotRevision::from_value(value)
+ }
+
+ fn snapshot(value: u64) -> AppSnapshot {
+ if value == 0 {
+ AppSnapshot::booting()
+ } else {
+ AppSnapshot::ready(
+ revision(value),
+ RelayConfiguration::default(),
+ Vec::new(),
+ None,
+ SessionState::SignedOut,
+ None,
+ None,
+ )
+ .expect("snapshot")
+ }
+ }
+}
diff --git a/crates/studio_application/src/lib.rs b/crates/studio_application/src/lib.rs
@@ -3,6 +3,7 @@
pub mod accounts;
pub mod actor;
pub mod app_core;
+mod change_stream;
pub mod config;
pub mod nostr_client;
pub mod ports;
@@ -23,6 +24,9 @@ pub use actor::{
RuntimeLifecycle, SessionGeneration, TaskCorrelation,
};
pub use app_core::{AppCore, AppObserver, ObserverHandle, RemovalConfirmationToken};
+pub use change_stream::{
+ ChangeSubscriptionId, OrderedSnapshotChanges, SnapshotChange, SnapshotChangeReceiver,
+};
pub use config::{
RelayRuntimeMode, relay_configuration_from_environment, relay_configuration_from_value,
};
diff --git a/crates/studio_storage/src/runtime_actor.rs b/crates/studio_storage/src/runtime_actor.rs
@@ -8,9 +8,9 @@ use std::time::{Duration, Instant};
use radroots_studio_application::{
ActorMailbox, AppObserver, AppSnapshot, Clock, CommandContext, CommandEnvelope, CommandReceipt,
CommandResult, CommandSubmission, GenerateAccountReceipt, ImportAccountReceipt, LifecycleGate,
- NostrClient, ObserverHandle, ProfileRefreshPlan, RelayConfiguration, RemovalConfirmationToken,
- RequestId, RuntimeCommandClass, RuntimeLifecycle, SecretStore, SessionGeneration,
- SnapshotRevision, TaskCorrelation,
+ NostrClient, ObserverHandle, OrderedSnapshotChanges, ProfileRefreshPlan, RelayConfiguration,
+ RemovalConfirmationToken, RequestId, RuntimeCommandClass, RuntimeLifecycle, SecretStore,
+ SessionGeneration, SnapshotRevision, TaskCorrelation,
};
use radroots_studio_domain::{
Kind0ProfileCandidate, PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput,
@@ -71,6 +71,7 @@ struct RuntimeActor {
session_generation: SessionGeneration,
published_session_generation: Arc<AtomicU64>,
profile_tasks: BTreeMap<RequestId, PendingProfileTask>,
+ changes: OrderedSnapshotChanges,
}
struct PendingProfileTask {
@@ -163,6 +164,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 actor = RuntimeActor {
adapter: Arc::clone(&adapter),
secrets,
@@ -173,6 +175,7 @@ impl RuntimeActorHandle {
session_generation: SessionGeneration::initial(),
published_session_generation: Arc::clone(&session_generation),
profile_tasks: BTreeMap::new(),
+ changes,
};
drop(runtime.spawn(actor.run(receiver)));
Ok(Self {
@@ -473,6 +476,9 @@ impl RuntimeActor {
if changes_session && matches!(result, CommandResult::Completed(_)) {
self.advance_session_generation();
}
+ if matches!(result, CommandResult::Completed(_)) {
+ self.changes.publish(self.adapter.core().snapshot());
+ }
let _ = reply.send(CommandReceipt::new(context.request_id(), result));
}
@@ -644,6 +650,9 @@ impl RuntimeActor {
} else {
CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(current)))
};
+ if matches!(result, CommandResult::Completed(_)) {
+ self.changes.publish(self.adapter.core().snapshot());
+ }
let _ = task
.reply
.send(CommandReceipt::new(task.correlation.request_id(), result));