1use serde::{Deserialize, Serialize};
8use std::io::{self, Write};
9use std::time::SystemTime;
10use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
11use uuid::Uuid;
12
13#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
17pub struct Event<L> {
18 pub time: SystemTime,
20 pub episode: Uuid,
22 pub payload: L,
24}
25
26impl<L> Event<L> {
27 pub fn now(episode: Uuid, payload: L) -> Self {
29 Self {
30 time: SystemTime::now(),
31 episode,
32 payload,
33 }
34 }
35}
36
37impl<L: Serialize> Event<L> {
38 pub fn write<W: Write>(&self, mut sink: W) -> io::Result<()> {
47 serde_json::to_writer(&mut sink, self)?;
48 sink.write_all(b"\n")
49 }
50}
51
52pub type Logger<L> = UnboundedSender<Event<L>>;
54
55pub async fn drain<L: Serialize, W: Write>(
63 mut events: UnboundedReceiver<Event<L>>,
64 mut sink: W,
65) -> io::Result<()> {
66 while let Some(event) = events.recv().await {
67 event.write(&mut sink)?;
68 }
69 sink.flush()
70}
71
72pub async fn console_log<L: Serialize>(events: UnboundedReceiver<Event<L>>) -> io::Result<()> {
78 drain(events, io::stderr()).await
79}
80
81#[cfg(test)]
82mod tests {
83 use super::*;
84 use std::time::{Duration, UNIX_EPOCH};
85 use tokio::sync::mpsc::unbounded_channel;
86
87 fn episode() -> Uuid {
88 Uuid::from_u128(1)
89 }
90
91 #[test]
92 fn an_event_is_stamped_when_it_is_made() {
93 let before = SystemTime::now();
94 let event = Event::now(episode(), "something happened");
95 let after = SystemTime::now();
96
97 assert!(before <= event.time && event.time <= after, "{event:?}");
98 assert_eq!(event.episode, episode());
99 assert_eq!(event.payload, "something happened");
100 }
101
102 #[tokio::test]
103 async fn drain_writes_each_event_as_a_line_of_json_until_every_logger_is_gone() {
104 let (logger, events) = unbounded_channel();
105 let at = UNIX_EPOCH + Duration::from_millis(1_700_000_000_123);
106 let one = Event {
107 time: at,
108 episode: episode(),
109 payload: "one".to_string(),
110 };
111 let two = Event {
112 time: at + Duration::from_millis(1_000),
113 episode: episode(),
114 payload: "two".to_string(),
115 };
116 logger.send(one.clone()).unwrap();
117 logger.send(two.clone()).unwrap();
118 drop(logger);
119
120 let mut sink = Vec::new();
121 drain(events, &mut sink).await.unwrap();
122
123 let written = String::from_utf8(sink).unwrap();
124 let lines: Vec<&str> = written.lines().collect();
125 assert_eq!(
126 lines[0],
127 concat!(
128 r#"{"time":{"secs_since_epoch":1700000000,"nanos_since_epoch":123000000},"#,
129 r#""episode":"00000000-0000-0000-0000-000000000001","payload":"one"}"#
130 )
131 );
132 let read: Vec<Event<String>> = lines
133 .iter()
134 .map(|line| serde_json::from_str(line).unwrap())
135 .collect();
136 assert_eq!(read, [one, two]);
137 }
138
139 #[tokio::test]
140 async fn a_time_before_the_epoch_cannot_be_written() {
141 let (logger, events) = unbounded_channel();
142 logger
143 .send(Event {
144 time: UNIX_EPOCH - Duration::from_secs(1),
145 episode: episode(),
146 payload: "long ago",
147 })
148 .unwrap();
149 drop(logger);
150
151 let mut sink = Vec::new();
152 assert!(drain(events, &mut sink).await.is_err());
153 }
154
155 #[tokio::test]
156 async fn console_log_writes_to_standard_error_until_every_logger_is_gone() {
157 let (logger, events) = unbounded_channel();
158 logger
159 .send(Event::now(episode(), "to standard error"))
160 .unwrap();
161 drop(logger);
162
163 console_log(events).await.unwrap();
164 }
165}