1use crate::value_objects::{BackpressureSignal, SessionId, StreamId};
4use chrono::{DateTime, Utc};
5use serde::{Deserialize, Serialize};
6
7mod serde_session_id {
9 use crate::value_objects::SessionId;
10 use serde::{Deserialize, Deserializer, Serialize, Serializer};
11
12 pub fn serialize<S>(id: &SessionId, serializer: S) -> Result<S::Ok, S::Error>
13 where
14 S: Serializer,
15 {
16 id.as_uuid().serialize(serializer)
17 }
18
19 pub fn deserialize<'de, D>(deserializer: D) -> Result<SessionId, D::Error>
20 where
21 D: Deserializer<'de>,
22 {
23 let uuid = uuid::Uuid::deserialize(deserializer)?;
24 Ok(SessionId::from_uuid(uuid))
25 }
26}
27
28mod serde_stream_id {
30 use crate::value_objects::StreamId;
31 use serde::{Deserialize, Deserializer, Serialize, Serializer};
32
33 pub fn serialize<S>(id: &StreamId, serializer: S) -> Result<S::Ok, S::Error>
34 where
35 S: Serializer,
36 {
37 id.as_uuid().serialize(serializer)
38 }
39
40 pub fn deserialize<'de, D>(deserializer: D) -> Result<StreamId, D::Error>
41 where
42 D: Deserializer<'de>,
43 {
44 let uuid = uuid::Uuid::deserialize(deserializer)?;
45 Ok(StreamId::from_uuid(uuid))
46 }
47}
48
49mod serde_option_stream_id {
51 use crate::value_objects::StreamId;
52 use serde::{Deserialize, Deserializer, Serialize, Serializer};
53
54 pub fn serialize<S>(id: &Option<StreamId>, serializer: S) -> Result<S::Ok, S::Error>
55 where
56 S: Serializer,
57 {
58 id.map(|i| i.as_uuid()).serialize(serializer)
59 }
60
61 pub fn deserialize<'de, D>(deserializer: D) -> Result<Option<StreamId>, D::Error>
62 where
63 D: Deserializer<'de>,
64 {
65 let uuid: Option<uuid::Uuid> = Option::deserialize(deserializer)?;
66 Ok(uuid.map(StreamId::from_uuid))
67 }
68}
69
70#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
72#[non_exhaustive]
73pub enum SessionState {
74 Initializing,
76 Active,
78 Closing,
80 Completed,
82 Failed,
84}
85
86impl SessionState {
87 pub fn as_str(&self) -> &'static str {
89 match self {
90 SessionState::Initializing => "Initializing",
91 SessionState::Active => "Active",
92 SessionState::Closing => "Closing",
93 SessionState::Completed => "Completed",
94 SessionState::Failed => "Failed",
95 }
96 }
97}
98
99#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
110#[serde(tag = "event_type", rename_all = "snake_case")]
111#[non_exhaustive]
112pub enum DomainEvent {
113 SessionActivated {
115 #[serde(with = "serde_session_id")]
117 session_id: SessionId,
118 timestamp: DateTime<Utc>,
120 },
121
122 SessionClosed {
124 #[serde(with = "serde_session_id")]
126 session_id: SessionId,
127 timestamp: DateTime<Utc>,
129 },
130
131 SessionExpired {
133 #[serde(with = "serde_session_id")]
135 session_id: SessionId,
136 timestamp: DateTime<Utc>,
138 },
139
140 SessionTimedOut {
142 #[serde(with = "serde_session_id")]
144 session_id: SessionId,
145 original_state: SessionState,
147 timeout_duration: u64,
149 timestamp: DateTime<Utc>,
151 },
152
153 SessionTimeoutExtended {
155 #[serde(with = "serde_session_id")]
157 session_id: SessionId,
158 additional_seconds: u64,
160 new_expires_at: DateTime<Utc>,
162 timestamp: DateTime<Utc>,
164 },
165
166 StreamCreated {
168 #[serde(with = "serde_session_id")]
170 session_id: SessionId,
171 #[serde(with = "serde_stream_id")]
173 stream_id: StreamId,
174 timestamp: DateTime<Utc>,
176 },
177
178 StreamStarted {
180 #[serde(with = "serde_session_id")]
182 session_id: SessionId,
183 #[serde(with = "serde_stream_id")]
185 stream_id: StreamId,
186 timestamp: DateTime<Utc>,
188 },
189
190 StreamCompleted {
192 #[serde(with = "serde_session_id")]
194 session_id: SessionId,
195 #[serde(with = "serde_stream_id")]
197 stream_id: StreamId,
198 timestamp: DateTime<Utc>,
200 },
201
202 StreamFailed {
204 #[serde(with = "serde_session_id")]
206 session_id: SessionId,
207 #[serde(with = "serde_stream_id")]
209 stream_id: StreamId,
210 error: String,
212 timestamp: DateTime<Utc>,
214 },
215
216 StreamCancelled {
218 #[serde(with = "serde_session_id")]
220 session_id: SessionId,
221 #[serde(with = "serde_stream_id")]
223 stream_id: StreamId,
224 timestamp: DateTime<Utc>,
226 },
227
228 SkeletonGenerated {
230 #[serde(with = "serde_session_id")]
232 session_id: SessionId,
233 #[serde(with = "serde_stream_id")]
235 stream_id: StreamId,
236 frame_size_bytes: u64,
238 timestamp: DateTime<Utc>,
240 },
241
242 PatchFramesGenerated {
244 #[serde(with = "serde_session_id")]
246 session_id: SessionId,
247 #[serde(with = "serde_stream_id")]
249 stream_id: StreamId,
250 frame_count: usize,
252 total_bytes: u64,
254 highest_priority: u8,
256 timestamp: DateTime<Utc>,
258 },
259
260 FramesBatched {
262 #[serde(with = "serde_session_id")]
264 session_id: SessionId,
265 frame_count: usize,
267 timestamp: DateTime<Utc>,
269 },
270
271 PriorityThresholdAdjusted {
273 #[serde(with = "serde_session_id")]
275 session_id: SessionId,
276 old_threshold: u8,
278 new_threshold: u8,
280 reason: String,
282 timestamp: DateTime<Utc>,
284 },
285
286 StreamConfigUpdated {
288 #[serde(with = "serde_session_id")]
290 session_id: SessionId,
291 #[serde(with = "serde_stream_id")]
293 stream_id: StreamId,
294 timestamp: DateTime<Utc>,
296 },
297
298 PerformanceMetricsRecorded {
300 #[serde(with = "serde_session_id")]
302 session_id: SessionId,
303 metrics: PerformanceMetrics,
305 timestamp: DateTime<Utc>,
307 },
308
309 BackpressureReceived {
311 #[serde(with = "serde_session_id")]
313 session_id: SessionId,
314 #[serde(default, skip_serializing_if = "Option::is_none")]
316 #[serde(with = "serde_option_stream_id")]
317 stream_id: Option<StreamId>,
318 signal: BackpressureSignal,
320 timestamp: DateTime<Utc>,
322 },
323
324 StreamPaused {
326 #[serde(with = "serde_session_id")]
328 session_id: SessionId,
329 #[serde(with = "serde_stream_id")]
331 stream_id: StreamId,
332 reason: String,
334 timestamp: DateTime<Utc>,
336 },
337
338 StreamResumed {
340 #[serde(with = "serde_session_id")]
342 session_id: SessionId,
343 #[serde(with = "serde_stream_id")]
345 stream_id: StreamId,
346 paused_duration_ms: u64,
348 timestamp: DateTime<Utc>,
350 },
351
352 CreditsUpdated {
354 #[serde(with = "serde_session_id")]
356 session_id: SessionId,
357 available_credits: usize,
359 max_credits: usize,
361 timestamp: DateTime<Utc>,
363 },
364}
365
366#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
368pub struct PerformanceMetrics {
369 pub frames_per_second: f64,
371 pub bytes_per_second: f64,
373 pub average_frame_size: f64,
375 pub priority_distribution: PriorityDistribution,
377 pub latency_ms: Option<u64>,
379}
380
381#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
383pub struct PriorityDistribution {
384 pub critical_frames: u64,
386 pub high_frames: u64,
388 pub medium_frames: u64,
390 pub low_frames: u64,
392 pub background_frames: u64,
394}
395
396impl PriorityDistribution {
397 pub fn new() -> Self {
399 Self::default()
400 }
401
402 pub fn total_frames(&self) -> u64 {
404 self.critical_frames
405 + self.high_frames
406 + self.medium_frames
407 + self.low_frames
408 + self.background_frames
409 }
410
411 pub fn as_percentages(&self) -> PriorityPercentages {
413 let total = self.total_frames() as f64;
414 if total == 0.0 {
415 return PriorityPercentages::default();
416 }
417
418 PriorityPercentages {
419 critical: self.critical_frames as f64 / total,
420 high: self.high_frames as f64 / total,
421 medium: self.medium_frames as f64 / total,
422 low: self.low_frames as f64 / total,
423 background: self.background_frames as f64 / total,
424 }
425 }
426
427 pub fn from_counts(
429 critical_count: u64,
430 high_count: u64,
431 medium_count: u64,
432 low_count: u64,
433 background_count: u64,
434 ) -> Self {
435 Self {
436 critical_frames: critical_count,
437 high_frames: high_count,
438 medium_frames: medium_count,
439 low_frames: low_count,
440 background_frames: background_count,
441 }
442 }
443}
444
445#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
447pub struct PriorityPercentages {
448 pub critical: f64,
450 pub high: f64,
452 pub medium: f64,
454 pub low: f64,
456 pub background: f64,
458}
459
460impl Default for PriorityPercentages {
461 fn default() -> Self {
462 Self {
463 critical: 0.0,
464 high: 0.0,
465 medium: 0.0,
466 low: 0.0,
467 background: 0.0,
468 }
469 }
470}
471
472impl DomainEvent {
473 pub fn session_id(&self) -> SessionId {
475 match self {
476 Self::SessionActivated { session_id, .. } => *session_id,
477 Self::SessionClosed { session_id, .. } => *session_id,
478 Self::SessionExpired { session_id, .. } => *session_id,
479 Self::StreamCreated { session_id, .. } => *session_id,
480 Self::StreamStarted { session_id, .. } => *session_id,
481 Self::StreamCompleted { session_id, .. } => *session_id,
482 Self::StreamFailed { session_id, .. } => *session_id,
483 Self::StreamCancelled { session_id, .. } => *session_id,
484 Self::SkeletonGenerated { session_id, .. } => *session_id,
485 Self::PatchFramesGenerated { session_id, .. } => *session_id,
486 Self::FramesBatched { session_id, .. } => *session_id,
487 Self::PriorityThresholdAdjusted { session_id, .. } => *session_id,
488 Self::StreamConfigUpdated { session_id, .. } => *session_id,
489 Self::PerformanceMetricsRecorded { session_id, .. } => *session_id,
490 Self::SessionTimedOut { session_id, .. } => *session_id,
491 Self::SessionTimeoutExtended { session_id, .. } => *session_id,
492 Self::BackpressureReceived { session_id, .. } => *session_id,
493 Self::StreamPaused { session_id, .. } => *session_id,
494 Self::StreamResumed { session_id, .. } => *session_id,
495 Self::CreditsUpdated { session_id, .. } => *session_id,
496 }
497 }
498
499 pub fn stream_id(&self) -> Option<StreamId> {
501 match self {
502 Self::StreamCreated { stream_id, .. } => Some(*stream_id),
503 Self::StreamStarted { stream_id, .. } => Some(*stream_id),
504 Self::StreamCompleted { stream_id, .. } => Some(*stream_id),
505 Self::StreamFailed { stream_id, .. } => Some(*stream_id),
506 Self::StreamCancelled { stream_id, .. } => Some(*stream_id),
507 Self::SkeletonGenerated { stream_id, .. } => Some(*stream_id),
508 Self::PatchFramesGenerated { stream_id, .. } => Some(*stream_id),
509 Self::StreamConfigUpdated { stream_id, .. } => Some(*stream_id),
510 Self::StreamPaused { stream_id, .. } => Some(*stream_id),
511 Self::StreamResumed { stream_id, .. } => Some(*stream_id),
512 Self::BackpressureReceived { stream_id, .. } => *stream_id,
513 Self::SessionActivated { .. }
514 | Self::SessionClosed { .. }
515 | Self::SessionExpired { .. }
516 | Self::SessionTimedOut { .. }
517 | Self::SessionTimeoutExtended { .. }
518 | Self::FramesBatched { .. }
519 | Self::PriorityThresholdAdjusted { .. }
520 | Self::PerformanceMetricsRecorded { .. }
521 | Self::CreditsUpdated { .. } => None,
522 }
523 }
524
525 pub fn timestamp(&self) -> DateTime<Utc> {
527 match self {
528 Self::SessionActivated { timestamp, .. } => *timestamp,
529 Self::SessionClosed { timestamp, .. } => *timestamp,
530 Self::SessionExpired { timestamp, .. } => *timestamp,
531 Self::StreamCreated { timestamp, .. } => *timestamp,
532 Self::StreamStarted { timestamp, .. } => *timestamp,
533 Self::StreamCompleted { timestamp, .. } => *timestamp,
534 Self::StreamFailed { timestamp, .. } => *timestamp,
535 Self::StreamCancelled { timestamp, .. } => *timestamp,
536 Self::SkeletonGenerated { timestamp, .. } => *timestamp,
537 Self::PatchFramesGenerated { timestamp, .. } => *timestamp,
538 Self::FramesBatched { timestamp, .. } => *timestamp,
539 Self::PriorityThresholdAdjusted { timestamp, .. } => *timestamp,
540 Self::StreamConfigUpdated { timestamp, .. } => *timestamp,
541 Self::PerformanceMetricsRecorded { timestamp, .. } => *timestamp,
542 Self::SessionTimedOut { timestamp, .. } => *timestamp,
543 Self::SessionTimeoutExtended { timestamp, .. } => *timestamp,
544 Self::BackpressureReceived { timestamp, .. } => *timestamp,
545 Self::StreamPaused { timestamp, .. } => *timestamp,
546 Self::StreamResumed { timestamp, .. } => *timestamp,
547 Self::CreditsUpdated { timestamp, .. } => *timestamp,
548 }
549 }
550
551 pub fn event_type(&self) -> &'static str {
553 match self {
554 Self::SessionActivated { .. } => "session_activated",
555 Self::SessionClosed { .. } => "session_closed",
556 Self::SessionExpired { .. } => "session_expired",
557 Self::StreamCreated { .. } => "stream_created",
558 Self::StreamStarted { .. } => "stream_started",
559 Self::StreamCompleted { .. } => "stream_completed",
560 Self::StreamFailed { .. } => "stream_failed",
561 Self::StreamCancelled { .. } => "stream_cancelled",
562 Self::SkeletonGenerated { .. } => "skeleton_generated",
563 Self::PatchFramesGenerated { .. } => "patch_frames_generated",
564 Self::FramesBatched { .. } => "frames_batched",
565 Self::PriorityThresholdAdjusted { .. } => "priority_threshold_adjusted",
566 Self::StreamConfigUpdated { .. } => "stream_config_updated",
567 Self::PerformanceMetricsRecorded { .. } => "performance_metrics_recorded",
568 Self::SessionTimedOut { .. } => "session_timed_out",
569 Self::SessionTimeoutExtended { .. } => "session_timeout_extended",
570 Self::BackpressureReceived { .. } => "backpressure_received",
571 Self::StreamPaused { .. } => "stream_paused",
572 Self::StreamResumed { .. } => "stream_resumed",
573 Self::CreditsUpdated { .. } => "credits_updated",
574 }
575 }
576
577 pub fn is_critical(&self) -> bool {
579 matches!(
580 self,
581 Self::StreamFailed { .. } | Self::SessionExpired { .. }
582 )
583 }
584
585 pub fn is_error(&self) -> bool {
587 matches!(self, Self::StreamFailed { .. })
588 }
589
590 pub fn is_completion(&self) -> bool {
592 matches!(
593 self,
594 Self::StreamCompleted { .. } | Self::SessionClosed { .. }
595 )
596 }
597}
598
599pub trait EventStore {
601 fn append_events(&mut self, events: Vec<DomainEvent>) -> Result<(), String>;
607
608 fn get_events_for_session(&self, session_id: SessionId) -> Result<Vec<DomainEvent>, String>;
614
615 fn get_events_for_stream(&self, stream_id: StreamId) -> Result<Vec<DomainEvent>, String>;
621
622 fn get_events_since(&self, since: DateTime<Utc>) -> Result<Vec<DomainEvent>, String>;
628}
629
630#[derive(Debug, Clone, Default)]
632pub struct InMemoryEventStore {
633 events: Vec<DomainEvent>,
634}
635
636impl InMemoryEventStore {
637 pub fn new() -> Self {
639 Self::default()
640 }
641
642 pub fn all_events(&self) -> &[DomainEvent] {
644 &self.events
645 }
646
647 pub fn event_count(&self) -> usize {
649 self.events.len()
650 }
651}
652
653impl EventStore for InMemoryEventStore {
654 fn append_events(&mut self, mut events: Vec<DomainEvent>) -> Result<(), String> {
655 self.events.append(&mut events);
656 Ok(())
657 }
658
659 fn get_events_for_session(&self, session_id: SessionId) -> Result<Vec<DomainEvent>, String> {
660 Ok(self
661 .events
662 .iter()
663 .filter(|e| e.session_id() == session_id)
664 .cloned()
665 .collect())
666 }
667
668 fn get_events_for_stream(&self, stream_id: StreamId) -> Result<Vec<DomainEvent>, String> {
669 Ok(self
670 .events
671 .iter()
672 .filter(|e| e.stream_id() == Some(stream_id))
673 .cloned()
674 .collect())
675 }
676
677 fn get_events_since(&self, since: DateTime<Utc>) -> Result<Vec<DomainEvent>, String> {
678 Ok(self
679 .events
680 .iter()
681 .filter(|e| e.timestamp() > since)
682 .cloned()
683 .collect())
684 }
685}
686
687#[cfg(test)]
688mod tests {
689 use super::*;
690 use crate::value_objects::{SessionId, StreamId};
691
692 #[test]
693 fn test_domain_event_properties() {
694 let session_id = SessionId::new();
695 let stream_id = StreamId::new();
696 let timestamp = Utc::now();
697
698 let event = DomainEvent::StreamCreated {
699 session_id,
700 stream_id,
701 timestamp,
702 };
703
704 assert_eq!(event.session_id(), session_id);
705 assert_eq!(event.stream_id(), Some(stream_id));
706 assert_eq!(event.timestamp(), timestamp);
707 assert_eq!(event.event_type(), "stream_created");
708 assert!(!event.is_critical());
709 assert!(!event.is_error());
710 }
711
712 #[test]
713 fn test_critical_events() {
714 let session_id = SessionId::new();
715 let stream_id = StreamId::new();
716
717 let error_event = DomainEvent::StreamFailed {
718 session_id,
719 stream_id,
720 error: "Connection lost".to_string(),
721 timestamp: Utc::now(),
722 };
723
724 assert!(error_event.is_critical());
725 assert!(error_event.is_error());
726 assert!(!error_event.is_completion());
727 }
728
729 #[test]
730 fn test_event_store() {
731 let mut store = InMemoryEventStore::new();
732 let session_id = SessionId::new();
733 let stream_id = StreamId::new();
734
735 let events = vec![
736 DomainEvent::SessionActivated {
737 session_id,
738 timestamp: Utc::now(),
739 },
740 DomainEvent::StreamCreated {
741 session_id,
742 stream_id,
743 timestamp: Utc::now(),
744 },
745 ];
746
747 store
748 .append_events(events.clone())
749 .expect("Failed to append events to store in test");
750 assert_eq!(store.event_count(), 2);
751
752 let session_events = store
753 .get_events_for_session(session_id)
754 .expect("Failed to retrieve session events in test");
755 assert_eq!(session_events.len(), 2);
756
757 let stream_events = store
758 .get_events_for_stream(stream_id)
759 .expect("Failed to retrieve stream events in test");
760 assert_eq!(stream_events.len(), 1);
761 }
762
763 #[test]
764 fn test_event_serialization() {
765 let session_id = SessionId::new();
766 let event = DomainEvent::SessionActivated {
767 session_id,
768 timestamp: Utc::now(),
769 };
770
771 let serialized = serde_json::to_string(&event).expect("Failed to serialize event in test");
772 let deserialized: DomainEvent =
773 serde_json::from_str(&serialized).expect("Failed to deserialize event in test");
774
775 assert_eq!(event, deserialized);
776 }
777}
778
779#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
781pub struct EventId(uuid::Uuid);
782
783impl EventId {
784 pub fn new() -> Self {
786 Self(uuid::Uuid::new_v4())
787 }
788
789 pub fn from_uuid(uuid: uuid::Uuid) -> Self {
791 Self(uuid)
792 }
793
794 pub fn inner(&self) -> uuid::Uuid {
796 self.0
797 }
798}
799
800impl std::fmt::Display for EventId {
801 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
802 write!(f, "{}", self.0)
803 }
804}
805
806impl Default for EventId {
807 fn default() -> Self {
808 Self::new()
809 }
810}
811
812pub trait EventSubscriber {
814 type HandleFuture<'a>: std::future::Future<Output = crate::DomainResult<()>> + Send + 'a
816 where
817 Self: 'a;
818
819 fn handle(&self, event: &DomainEvent) -> Self::HandleFuture<'_>;
821}
822
823impl DomainEvent {
825 pub fn occurred_at(&self) -> DateTime<Utc> {
827 self.timestamp()
828 }
829
830 pub fn metadata(&self) -> std::collections::HashMap<String, String> {
832 let mut metadata = std::collections::HashMap::new();
833 metadata.insert("event_type".to_string(), self.event_type().to_string());
834 metadata.insert("session_id".to_string(), self.session_id().to_string());
835 metadata.insert("timestamp".to_string(), self.timestamp().to_rfc3339());
836
837 if let Some(stream_id) = self.stream_id() {
838 metadata.insert("stream_id".to_string(), stream_id.to_string());
839 }
840
841 metadata
842 }
843}