From 69e210de7141d962deeed92eb4a06b19214bfec9 Mon Sep 17 00:00:00 2001 From: Orual Date: Sat, 18 Oct 2025 20:15:19 -0400 Subject: [PATCH] thing to allow reconfiguration using the original sub parameters additionally carried through the actual subscription parameters and drove the trait that way allows easier reconnect logic, etc. --- crates/jacquard-common/src/websocket.rs | 20 +++++++ .../jacquard-common/src/xrpc/subscription.rs | 56 +++++++++++++++++++ 2 files changed, 76 insertions(+) diff --git a/crates/jacquard-common/src/websocket.rs b/crates/jacquard-common/src/websocket.rs index 376efe739..6b21f4837 100644 --- a/crates/jacquard-common/src/websocket.rs +++ b/crates/jacquard-common/src/websocket.rs @@ -409,6 +409,26 @@ impl WsSink { pub fn into_inner(self) -> Pin>> { self.0 } + + /// get a mutable reference to the inner boxed sink + #[cfg(not(target_arch = "wasm32"))] + pub fn get_mut( + &mut self, + ) -> &mut Pin + Send>> { + use std::borrow::BorrowMut; + + self.0.borrow_mut() + } + + /// get a mutable reference to the inner boxed sink + #[cfg(target_arch = "wasm32")] + pub fn get_mut( + &mut self, + ) -> &mut Pin + 'static>> { + use std::borrow::BorrowMut; + + self.0.borrow_mut() + } } impl fmt::Debug for WsSink { diff --git a/crates/jacquard-common/src/xrpc/subscription.rs b/crates/jacquard-common/src/xrpc/subscription.rs index 784bb1a10..26aa3fbb1 100644 --- a/crates/jacquard-common/src/xrpc/subscription.rs +++ b/crates/jacquard-common/src/xrpc/subscription.rs @@ -216,6 +216,62 @@ where } } +/// Websocket subscriber-sent control message +/// +/// Note: this is not meaningful for atproto event stream endpoints as +/// those do not support control after the fact. Jetstream does, however. +/// +/// If you wish to control an ongoing Jetstream connection, wrap the [`WsSink`] +/// returned from one of the `into_*` methods of the [`SubscriptionStream`] +/// in a [`SubscriptionController`] with the corresponding message implementing +/// this trait as a generic parameter. +pub trait SubscriptionControlMessage: Serialize { + /// The subscription this is associated with + type Subscription: XrpcSubscription; + + /// Encode the control message for transmission + /// + /// Defaults to json text (matches Jetstream) + fn encode(&self) -> Result { + Ok(WsMessage::from( + serde_json::to_string(&self).map_err(StreamError::encode)?, + )) + } + + /// Decode the control message + fn decode<'de>(frame: &'de [u8]) -> Result + where + Self: Deserialize<'de>, + { + Ok(serde_json::from_slice(frame).map_err(StreamError::decode)?) + } +} + +/// Control a websocket stream with a given subscription control message +pub struct SubscriptionController { + controller: WsSink, + _marker: PhantomData S>, +} + +impl SubscriptionController { + /// Create a new subscription controller from a WebSocket sink. + pub fn new(controller: WsSink) -> Self { + Self { + controller, + _marker: PhantomData, + } + } + + /// Configure the upstream connection via the websocket + pub async fn configure(&mut self, params: &S) -> Result<(), StreamError> { + let message = params.encode()?; + + n0_future::SinkExt::send(self.controller.get_mut(), message) + .await + .map_err(StreamError::transport) + } +} + /// Typed subscription stream wrapping a WebSocket connection. /// /// Analogous to `Response` for XRPC but for subscription streams. -- 2.51.2