1mod fan_in;
12
13use std::io::Write;
14
15use crate::event::{Event, JsonlWriter};
16
17pub use fan_in::{EventFanIn, MergedEvents, TaggedEvent, TaggedSink};
18
19pub trait EventSink: Send + 'static {
29 fn emit(&mut self, event: Event) -> std::io::Result<()>;
30}
31
32impl EventSink for Box<dyn EventSink> {
38 fn emit(&mut self, event: Event) -> std::io::Result<()> {
39 (**self).emit(event)
40 }
41}
42
43impl<W: Write + Send + 'static> EventSink for JsonlWriter<W> {
44 fn emit(&mut self, event: Event) -> std::io::Result<()> {
45 self.write(event).map(|_| ())
46 }
47}
48
49#[derive(Debug, Default, Clone, PartialEq)]
52pub struct CollectingSink {
53 events: Vec<Event>,
54}
55
56impl CollectingSink {
57 pub fn new() -> Self {
58 Self::default()
59 }
60
61 pub fn events(&self) -> &[Event] {
62 &self.events
63 }
64
65 pub fn into_events(self) -> Vec<Event> {
66 self.events
67 }
68}
69
70impl EventSink for CollectingSink {
71 fn emit(&mut self, event: Event) -> std::io::Result<()> {
72 self.events.push(event);
73 Ok(())
74 }
75}
76
77#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
79pub struct NullSink;
80
81impl EventSink for NullSink {
82 fn emit(&mut self, _event: Event) -> std::io::Result<()> {
83 Ok(())
84 }
85}
86
87pub struct FnSink<F>(F)
89where
90 F: FnMut(Event) -> std::io::Result<()> + Send + 'static;
91
92impl<F> FnSink<F>
93where
94 F: FnMut(Event) -> std::io::Result<()> + Send + 'static,
95{
96 pub fn new(callback: F) -> Self {
97 Self(callback)
98 }
99}
100
101impl<F> EventSink for FnSink<F>
102where
103 F: FnMut(Event) -> std::io::Result<()> + Send + 'static,
104{
105 fn emit(&mut self, event: Event) -> std::io::Result<()> {
106 (self.0)(event)
107 }
108}
109
110#[cfg(test)]
111mod tests {
112 use super::*;
113 use crate::event::RunOutcome;
114
115 fn delta(text: &str) -> Event {
116 Event::AssistantDelta {
117 text: text.to_string(),
118 }
119 }
120
121 #[test]
122 fn collecting_sink_keeps_order() {
123 let mut sink = CollectingSink::new();
124 sink.emit(delta("a")).expect("emits");
125 sink.emit(delta("b")).expect("emits");
126
127 assert_eq!(sink.into_events(), vec![delta("a"), delta("b")]);
128 }
129
130 #[test]
131 fn null_sink_accepts_everything() {
132 let mut sink = NullSink;
133
134 assert!(
135 sink.emit(Event::RunFinished {
136 outcome: RunOutcome::Ok,
137 stopped_by: None
138 })
139 .is_ok()
140 );
141 }
142
143 #[test]
144 fn fn_sink_forwards_to_the_closure() {
145 let (tx, rx) = std::sync::mpsc::channel();
146 let mut sink = FnSink::new(move |event| {
147 tx.send(event).expect("receiver alive");
148 Ok(())
149 });
150
151 sink.emit(delta("x")).expect("emits");
152
153 assert_eq!(rx.recv().expect("an event"), delta("x"));
154 }
155
156 #[test]
157 fn a_jsonl_writer_is_a_sink() {
158 let mut sink = JsonlWriter::new(Vec::new());
159 sink.emit(delta("hi")).expect("emits");
160
161 let written = String::from_utf8(sink.into_inner()).expect("utf-8");
162 assert!(written.contains("\"type\":\"assistant_delta\""));
163 }
164
165 #[test]
166 fn a_sink_chosen_at_runtime_is_still_a_sink() {
167 let mut sink: Box<dyn EventSink> = if cfg!(test) {
170 Box::new(CollectingSink::new())
171 } else {
172 Box::new(NullSink)
173 };
174
175 sink.emit(delta("boxed")).expect("emits");
176 }
177}