mod fan_in;
use std::io::Write;
use crate::event::{Event, JsonlWriter};
pub use fan_in::{EventFanIn, MergedEvents, TaggedEvent, TaggedSink};
pub trait EventSink: Send + 'static {
fn emit(&mut self, event: Event) -> std::io::Result<()>;
}
impl EventSink for Box<dyn EventSink> {
fn emit(&mut self, event: Event) -> std::io::Result<()> {
(**self).emit(event)
}
}
impl<W: Write + Send + 'static> EventSink for JsonlWriter<W> {
fn emit(&mut self, event: Event) -> std::io::Result<()> {
self.write(event).map(|_| ())
}
}
#[derive(Debug, Default, Clone, PartialEq)]
pub struct CollectingSink {
events: Vec<Event>,
}
impl CollectingSink {
pub fn new() -> Self {
Self::default()
}
pub fn events(&self) -> &[Event] {
&self.events
}
pub fn into_events(self) -> Vec<Event> {
self.events
}
}
impl EventSink for CollectingSink {
fn emit(&mut self, event: Event) -> std::io::Result<()> {
self.events.push(event);
Ok(())
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct NullSink;
impl EventSink for NullSink {
fn emit(&mut self, _event: Event) -> std::io::Result<()> {
Ok(())
}
}
pub struct FnSink<F>(F)
where
F: FnMut(Event) -> std::io::Result<()> + Send + 'static;
impl<F> FnSink<F>
where
F: FnMut(Event) -> std::io::Result<()> + Send + 'static,
{
pub fn new(callback: F) -> Self {
Self(callback)
}
}
impl<F> EventSink for FnSink<F>
where
F: FnMut(Event) -> std::io::Result<()> + Send + 'static,
{
fn emit(&mut self, event: Event) -> std::io::Result<()> {
(self.0)(event)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::RunOutcome;
fn delta(text: &str) -> Event {
Event::AssistantDelta {
text: text.to_string(),
}
}
#[test]
fn collecting_sink_keeps_order() {
let mut sink = CollectingSink::new();
sink.emit(delta("a")).expect("emits");
sink.emit(delta("b")).expect("emits");
assert_eq!(sink.into_events(), vec![delta("a"), delta("b")]);
}
#[test]
fn null_sink_accepts_everything() {
let mut sink = NullSink;
assert!(
sink.emit(Event::RunFinished {
outcome: RunOutcome::Ok,
stopped_by: None
})
.is_ok()
);
}
#[test]
fn fn_sink_forwards_to_the_closure() {
let (tx, rx) = std::sync::mpsc::channel();
let mut sink = FnSink::new(move |event| {
tx.send(event).expect("receiver alive");
Ok(())
});
sink.emit(delta("x")).expect("emits");
assert_eq!(rx.recv().expect("an event"), delta("x"));
}
#[test]
fn a_jsonl_writer_is_a_sink() {
let mut sink = JsonlWriter::new(Vec::new());
sink.emit(delta("hi")).expect("emits");
let written = String::from_utf8(sink.into_inner()).expect("utf-8");
assert!(written.contains("\"type\":\"assistant_delta\""));
}
#[test]
fn a_sink_chosen_at_runtime_is_still_a_sink() {
let mut sink: Box<dyn EventSink> = if cfg!(test) {
Box::new(CollectingSink::new())
} else {
Box::new(NullSink)
};
sink.emit(delta("boxed")).expect("emits");
}
}