Skip to main content

devicerail_protocol/
stream.rs

1use std::fmt;
2
3use serde::{Deserialize, Deserializer, Serialize};
4use uuid::Uuid;
5
6use crate::{
7    ErrorInfo, EventSequence, JsonRpcVersion, RpcResponse, SessionId, SessionState, TestEvent,
8};
9
10/// Identifies one daemon process lifetime. A cursor from another epoch must
11/// never be accepted as a resumable position.
12#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
13#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Deserialize, Serialize)]
14#[serde(transparent)]
15pub struct EventStreamEpoch(pub Uuid);
16
17impl EventStreamEpoch {
18    pub fn new() -> Self {
19        Self(Uuid::new_v4())
20    }
21}
22
23impl Default for EventStreamEpoch {
24    fn default() -> Self {
25        Self::new()
26    }
27}
28
29impl fmt::Display for EventStreamEpoch {
30    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
31        self.0.fmt(formatter)
32    }
33}
34
35#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
36#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Deserialize, Serialize)]
37#[serde(transparent)]
38pub struct EventSubscriptionId(pub Uuid);
39
40impl EventSubscriptionId {
41    pub fn new() -> Self {
42        Self(Uuid::new_v4())
43    }
44}
45
46impl Default for EventSubscriptionId {
47    fn default() -> Self {
48        Self::new()
49    }
50}
51
52impl fmt::Display for EventSubscriptionId {
53    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
54        self.0.fmt(formatter)
55    }
56}
57
58/// A Session-scoped, daemon-epoch-scoped application acknowledgement.
59#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
60#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)]
61#[serde(rename_all = "camelCase", deny_unknown_fields)]
62pub struct EventStreamCursor {
63    pub stream_epoch: EventStreamEpoch,
64    pub session_id: SessionId,
65    pub sequence: EventSequence,
66}
67
68/// Browser Origin policy bound to a single-use stream capability.
69#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
70#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)]
71#[serde(tag = "kind", rename_all = "camelCase", deny_unknown_fields)]
72pub enum EventStreamOriginPolicy {
73    /// Node and other non-browser clients must omit Origin entirely.
74    Absent {},
75    /// Browser clients must send exactly this Origin value.
76    Exact { origin: String },
77}
78
79/// A short-lived bearer URL. Debug intentionally never exposes its contents.
80#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
81#[derive(Clone, PartialEq, Eq, Serialize)]
82#[serde(transparent)]
83pub struct EventStreamEndpoint(
84    #[cfg_attr(feature = "schema", schemars(length(min = 1, max = 2048)))] String,
85);
86
87impl EventStreamEndpoint {
88    pub const MAX_BYTES: usize = 2_048;
89
90    pub fn new(value: String) -> Option<Self> {
91        (!value.is_empty() && value.len() <= Self::MAX_BYTES).then_some(Self(value))
92    }
93
94    /// Explicitly exposes the bearer URL to the transport that will consume it.
95    pub fn expose_secret(&self) -> &str {
96        &self.0
97    }
98}
99
100impl fmt::Debug for EventStreamEndpoint {
101    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
102        formatter.write_str("EventStreamEndpoint(<redacted>)")
103    }
104}
105
106impl<'de> Deserialize<'de> for EventStreamEndpoint {
107    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
108    where
109        D: Deserializer<'de>,
110    {
111        let value = String::deserialize(deserializer)?;
112        Self::new(value).ok_or_else(|| {
113            serde::de::Error::custom(format!(
114                "event stream endpoint must contain 1..={} bytes",
115                Self::MAX_BYTES
116            ))
117        })
118    }
119}
120
121#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
122#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)]
123#[serde(rename_all = "camelCase", deny_unknown_fields)]
124pub struct EventsStreamOpenParams {
125    pub session_id: SessionId,
126    pub origin_policy: EventStreamOriginPolicy,
127}
128
129#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
130#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)]
131#[serde(rename_all = "camelCase", deny_unknown_fields)]
132pub struct EventsStreamOpenResult {
133    pub endpoint: EventStreamEndpoint,
134    pub stream_epoch: EventStreamEpoch,
135    #[serde(
136        serialize_with = "crate::wire_integer::serialize_js_safe_u64",
137        deserialize_with = "crate::wire_integer::deserialize_js_safe_u64"
138    )]
139    #[cfg_attr(
140        feature = "schema",
141        schemars(range(min = 0_u64, max = 9_007_199_254_740_991_u64))
142    )]
143    pub expires_at_ms: u64,
144}
145
146#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
147#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)]
148#[serde(rename_all = "camelCase", deny_unknown_fields)]
149pub struct EventsSubscribeParams {
150    pub session_id: SessionId,
151    #[serde(default, skip_serializing_if = "Option::is_none")]
152    pub after_cursor: Option<EventStreamCursor>,
153}
154
155#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
156#[derive(Clone, Debug, PartialEq, Eq, Deserialize, Serialize)]
157#[serde(rename_all = "camelCase", deny_unknown_fields)]
158pub struct EventsSubscribeResult {
159    pub subscription_id: EventSubscriptionId,
160    pub session_id: SessionId,
161    pub replay_through: EventStreamCursor,
162    pub session_state: SessionState,
163}
164
165#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
166#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize, Serialize)]
167pub enum EventsStreamEventMethod {
168    #[serde(rename = "events.stream.event")]
169    Event,
170}
171
172#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
173#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
174#[serde(rename_all = "camelCase", deny_unknown_fields)]
175pub struct EventsStreamEventParams {
176    pub subscription_id: EventSubscriptionId,
177    pub cursor: EventStreamCursor,
178    pub event: TestEvent,
179}
180
181#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
182#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
183#[serde(rename_all = "camelCase", deny_unknown_fields)]
184pub struct EventsStreamEventNotification {
185    pub jsonrpc: JsonRpcVersion,
186    pub method: EventsStreamEventMethod,
187    pub params: EventsStreamEventParams,
188}
189
190impl EventsStreamEventNotification {
191    pub fn new(params: EventsStreamEventParams) -> Self {
192        Self {
193            jsonrpc: JsonRpcVersion::V2,
194            method: EventsStreamEventMethod::Event,
195            params,
196        }
197    }
198}
199
200/// The terminal reason is a closed tagged union so reason/error combinations
201/// cannot become ambiguous on the wire.
202#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
203#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
204#[serde(tag = "reason", rename_all = "camelCase", deny_unknown_fields)]
205pub enum EventsStreamTermination {
206    SessionEnded,
207    Cancelled,
208    SlowConsumer { error: ErrorInfo },
209    SessionDeleted { error: ErrorInfo },
210    ServerShutdown { error: ErrorInfo },
211    SequenceGap { error: ErrorInfo },
212    EventTooLarge { error: ErrorInfo },
213    InternalError { error: ErrorInfo },
214}
215
216#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
217#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
218#[serde(rename_all = "camelCase", deny_unknown_fields)]
219pub struct EventsStreamTerminalParams {
220    pub subscription_id: EventSubscriptionId,
221    pub session_id: SessionId,
222    #[serde(default, skip_serializing_if = "Option::is_none")]
223    pub last_emitted_cursor: Option<EventStreamCursor>,
224    pub termination: EventsStreamTermination,
225}
226
227#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
228#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize, Serialize)]
229pub enum EventsStreamTerminalMethod {
230    #[serde(rename = "events.stream.terminal")]
231    Terminal,
232}
233
234#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
235#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
236#[serde(rename_all = "camelCase", deny_unknown_fields)]
237pub struct EventsStreamTerminalNotification {
238    pub jsonrpc: JsonRpcVersion,
239    pub method: EventsStreamTerminalMethod,
240    pub params: EventsStreamTerminalParams,
241}
242
243impl EventsStreamTerminalNotification {
244    pub fn new(params: EventsStreamTerminalParams) -> Self {
245        Self {
246            jsonrpc: JsonRpcVersion::V2,
247            method: EventsStreamTerminalMethod::Terminal,
248            params,
249        }
250    }
251}
252
253#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
254#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
255#[serde(untagged)]
256pub enum RpcServerNotification {
257    Event(EventsStreamEventNotification),
258    Terminal(EventsStreamTerminalNotification),
259}
260
261#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
262#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
263#[serde(untagged)]
264pub enum RpcServerMessage {
265    Response(RpcResponse),
266    Notification(RpcServerNotification),
267}
268
269#[cfg(test)]
270mod tests {
271    use serde_json::json;
272
273    use super::{
274        EventStreamEndpoint, EventStreamEpoch, EventStreamOriginPolicy, EventsStreamOpenParams,
275        EventsStreamTermination,
276    };
277    use crate::SessionId;
278
279    #[test]
280    fn endpoint_debug_is_redacted_and_input_is_bounded() {
281        let endpoint =
282            EventStreamEndpoint::new("ws://127.0.0.1:1234/v/secret-capability".to_owned())
283                .expect("valid endpoint");
284        assert_eq!(format!("{endpoint:?}"), "EventStreamEndpoint(<redacted>)");
285        assert!(!format!("{endpoint:?}").contains("secret-capability"));
286        assert!(EventStreamEndpoint::new(String::new()).is_none());
287        assert!(EventStreamEndpoint::new("x".repeat(EventStreamEndpoint::MAX_BYTES + 1)).is_none());
288    }
289
290    #[test]
291    fn origin_policy_and_termination_are_closed_tagged_unions() {
292        let params: EventsStreamOpenParams = serde_json::from_value(json!({
293            "sessionId": SessionId::new(),
294            "originPolicy": { "kind": "absent" }
295        }))
296        .expect("absent origin");
297        assert_eq!(params.origin_policy, EventStreamOriginPolicy::Absent {});
298        assert!(
299            serde_json::from_value::<EventsStreamOpenParams>(json!({
300                "sessionId": SessionId::new(),
301                "originPolicy": { "kind": "absent", "origin": "null" }
302            }))
303            .is_err()
304        );
305
306        assert_eq!(
307            serde_json::to_value(EventsStreamTermination::SessionEnded)
308                .expect("serialize termination"),
309            json!({ "reason": "sessionEnded" })
310        );
311        let _ = EventStreamEpoch::new();
312    }
313}