use std::io::{self, Write};
use crate::canonical::{Delta, Event, ExitClass};
pub trait Sink {
fn write(&mut self, ev: &Event) -> io::Result<()>;
}
pub struct NdjsonSink<W: Write> {
w: W,
}
impl<W: Write> NdjsonSink<W> {
pub fn new(w: W) -> Self {
Self { w }
}
}
impl<W: Write> Sink for NdjsonSink<W> {
fn write(&mut self, ev: &Event) -> io::Result<()> {
#[allow(clippy::expect_used)]
let mut buf = serde_json::to_vec(ev).expect("Event is infallibly serializable");
buf.push(b'\n');
self.w.write_all(&buf)?;
self.w.flush()
}
}
pub struct TextSink<O: Write, E: Write> {
out: O,
err: E,
thinking: bool,
pending_sep: bool,
}
impl<O: Write, E: Write> TextSink<O, E> {
pub fn new(out: O, err: E, thinking: bool) -> Self {
Self {
out,
err,
thinking,
pending_sep: false,
}
}
}
impl<O: Write, E: Write> Sink for TextSink<O, E> {
fn write(&mut self, ev: &Event) -> io::Result<()> {
match ev {
Event::ContentDelta {
delta: Delta::ThinkingDelta(text),
..
} if self.thinking => {
self.out.write_all(text.as_bytes())?;
self.pending_sep = true;
self.out.flush()
}
Event::ContentDelta {
delta: Delta::TextDelta(text),
..
} => {
if self.pending_sep {
self.out.write_all(b"\n")?;
self.pending_sep = false;
}
self.out.write_all(text.as_bytes())?;
self.out.flush()
}
Event::Error(err) => {
writeln!(self.err, "{}", err.message)?;
self.err.flush()
}
_ => Ok(()),
}
}
}
pub struct RawSink<W: Write> {
w: W,
}
impl<W: Write> RawSink<W> {
pub fn new(w: W) -> Self {
Self { w }
}
}
impl<W: Write> Sink for RawSink<W> {
fn write(&mut self, ev: &Event) -> io::Result<()> {
if let Event::Raw(bytes) = ev {
self.w.write_all(bytes)?;
self.w.flush()?;
}
Ok(())
}
}
pub(crate) fn pump<I: IntoIterator<Item = Event>>(events: I, sink: &mut dyn Sink) -> u8 {
let mut exit = ExitClass::Ok.code();
for ev in events {
if let Event::Error(err) = &ev {
exit = err.exit_code();
}
if let Err(io_err) = sink.write(&ev) {
return ExitClass::from_io(&io_err).code();
}
}
exit
}