use std::io::{BufRead, BufReader, Write};
use std::net::{TcpStream, ToSocketAddrs};
use crate::{Error, MessageEvent};
const TERMINATOR: &[u8] = b"\r\n";
#[derive(Debug)]
pub struct LineSink {
stream: TcpStream,
}
impl LineSink {
pub fn connect(address: impl ToSocketAddrs) -> Result<Self, Error> {
Ok(Self {
stream: TcpStream::connect(address)?,
})
}
pub fn send(&mut self, event: &MessageEvent) -> Result<(), Error> {
let mut line = serde_json::to_vec(event)?;
line.extend_from_slice(TERMINATOR);
self.stream.write_all(&line)?;
self.stream.flush()?;
Ok(())
}
}
#[derive(Debug)]
pub struct LineSource {
reader: BufReader<TcpStream>,
}
impl LineSource {
#[must_use]
pub fn new(stream: TcpStream) -> Self {
Self {
reader: BufReader::new(stream),
}
}
pub fn next_event(&mut self) -> Result<Option<MessageEvent>, Error> {
let mut line = String::new();
loop {
line.clear();
if self.reader.read_line(&mut line)? == 0 {
return Ok(None);
}
let trimmed = line.trim_end_matches(['\r', '\n']);
if trimmed.is_empty() {
continue;
}
return Ok(Some(serde_json::from_str(trimmed)?));
}
}
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use super::TERMINATOR;
use crate::MessageEvent;
fn event() -> MessageEvent {
MessageEvent {
code: 10115,
numeric: 0.0,
text: "context".to_owned(),
payload: Some(r#"{"Locals":{"CustomStatus":"ok"}}"#.to_owned()),
synchronous: true,
execution_id: Some(1),
}
}
#[test]
fn a_frame_is_one_line_and_ends_with_crlf() {
let mut line = serde_json::to_vec(&event()).unwrap_or_default();
assert!(
!line.contains(&b'\n'),
"a serialized event must not contain a raw newline"
);
line.extend_from_slice(TERMINATOR);
assert!(line.ends_with(b"\r\n"), "a frame ends with CRLF");
}
#[test]
fn an_event_survives_the_round_trip_as_text() {
let text = serde_json::to_string(&event()).unwrap_or_default();
let back: MessageEvent = serde_json::from_str(&text).unwrap_or_else(|_| MessageEvent {
code: -1,
numeric: 0.0,
text: String::new(),
payload: None,
synchronous: false,
execution_id: None,
});
assert_eq!(back, event());
}
#[test]
fn an_embedded_newline_is_escaped_rather_than_emitted() {
let mut multiline = event();
multiline.text = "first\r\nsecond".to_owned();
let mut container = BTreeMap::new();
container.insert("note".to_owned(), "a\nb".to_owned());
multiline.payload = serde_json::to_string(&container).ok();
let line = serde_json::to_string(&multiline).unwrap_or_default();
assert!(!line.contains('\n'), "framing would break: {line}");
assert!(!line.contains('\r'), "framing would break: {line}");
}
}