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#[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#[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#[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 Absent {},
75 Exact { origin: String },
77}
78
79#[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 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#[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}