use std::collections::VecDeque;
use super::TransportError;
pub(crate) const SCHEME_SEPARATOR: &str = "://";
pub(crate) const CONTROL_LINK_SAME: &str = "*";
pub(crate) const MAX_LINE_BYTES: usize = 8 * 1024 * 1024;
#[derive(Debug, Default)]
pub(crate) struct LineBuffer {
partial: Vec<u8>,
scanned: usize,
ready: VecDeque<Vec<u8>>,
}
impl LineBuffer {
pub(crate) fn push_bytes(&mut self, bytes: &[u8]) -> Result<(), TransportError> {
self.partial.extend_from_slice(bytes);
let mut consumed = 0;
let mut cursor = self.scanned;
while let Some(offset) = self
.partial
.get(cursor..)
.and_then(|rest| rest.iter().position(|&byte| byte == b'\n'))
{
let end = cursor.saturating_add(offset);
let raw = self.partial.get(consumed..end).unwrap_or_default();
let line = raw.strip_suffix(b"\r").unwrap_or(raw);
if !line.is_empty() {
self.ready.push_back(line.to_vec());
}
consumed = end.saturating_add(1);
cursor = consumed;
}
if consumed > 0 {
self.partial.drain(..consumed);
}
self.scanned = self.partial.len();
if self.partial.len() > MAX_LINE_BYTES {
let held = self.partial.len();
self.clear();
return Err(TransportError::Capacity {
limit_bytes: MAX_LINE_BYTES,
reason: format!("{held} bytes were received with no line terminator"),
});
}
Ok(())
}
#[inline]
pub(crate) fn push_message(&mut self, text: &str) -> Result<(), TransportError> {
self.push_bytes(text.as_bytes())
}
pub(crate) fn pop_line(&mut self) -> Option<Result<String, TransportError>> {
let bytes = self.ready.pop_front()?;
Some(
String::from_utf8(bytes).map_err(|error| TransportError::MalformedFrame {
reason: format!("a line was not valid UTF-8: {}", error.utf8_error()),
}),
)
}
pub(crate) fn flush_partial(&mut self) -> Option<Result<String, TransportError>> {
self.scanned = 0;
if self.partial.is_empty() {
return None;
}
let mut raw = std::mem::take(&mut self.partial);
if raw.last() == Some(&b'\r') {
raw.pop();
}
if raw.is_empty() {
return None;
}
Some(
String::from_utf8(raw).map_err(|error| TransportError::MalformedFrame {
reason: format!("a line was not valid UTF-8: {}", error.utf8_error()),
}),
)
}
pub(crate) fn discard_partial(&mut self) -> Option<usize> {
let held = self.partial.len();
self.partial.clear();
self.scanned = 0;
(held > 0).then_some(held)
}
pub(crate) fn clear(&mut self) {
self.partial.clear();
self.scanned = 0;
self.ready.clear();
}
}
#[must_use]
pub(crate) fn redact(uri: &str) -> String {
let Some((scheme, rest)) = uri.split_once(SCHEME_SEPARATOR) else {
return uri.to_owned();
};
let authority_end = rest.find('/').unwrap_or(rest.len());
let authority = rest.get(..authority_end).unwrap_or(rest);
let path = rest.get(authority_end..).unwrap_or("");
match authority.rsplit_once('@') {
Some((_, host)) => format!("{scheme}{SCHEME_SEPARATOR}***@{host}{path}"),
None => uri.to_owned(),
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used, clippy::expect_used)]
use super::*;
fn push(buffer: &mut LineBuffer, bytes: &[u8]) {
assert!(
buffer.push_bytes(bytes).is_ok(),
"the push is well within the configured limits"
);
}
fn line(buffer: &mut LineBuffer) -> Option<String> {
buffer.pop_line().map(|result| result.expect("valid UTF-8"))
}
#[test]
fn test_line_buffer_single_line_yields_that_line() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"WSOK\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("WSOK"));
assert!(buffer.pop_line().is_none());
}
#[test]
fn test_line_buffer_multi_line_push_yields_every_line_in_order() {
let mut buffer = LineBuffer::default();
push(
&mut buffer,
b"CONOK,S1aa6c792585db57aT1726545,50000,5000,*\r\n\
SERVNAME,Lightstreamer HTTP Server\r\n\
CLIENTIP,127.0.0.1\r\n\
CONS,unlimited\r\n",
);
assert_eq!(
line(&mut buffer).as_deref(),
Some("CONOK,S1aa6c792585db57aT1726545,50000,5000,*")
);
assert_eq!(
line(&mut buffer).as_deref(),
Some("SERVNAME,Lightstreamer HTTP Server")
);
assert_eq!(line(&mut buffer).as_deref(), Some("CLIENTIP,127.0.0.1"));
assert_eq!(line(&mut buffer).as_deref(), Some("CONS,unlimited"));
assert!(buffer.pop_line().is_none());
}
#[test]
fn test_line_buffer_line_split_across_pushes_is_reassembled() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"U,1,1,20.4");
assert!(buffer.pop_line().is_none(), "no terminator yet");
push(&mut buffer, b"2|EUR\r\nU,1,2,");
assert_eq!(line(&mut buffer).as_deref(), Some("U,1,1,20.42|EUR"));
assert!(buffer.pop_line().is_none());
push(&mut buffer, b"3\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("U,1,2,3"));
}
#[test]
fn test_line_buffer_utf8_split_across_pushes_is_held_until_complete() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"SERVNAME,\xE2\x82");
assert!(buffer.pop_line().is_none(), "the character is incomplete");
push(&mut buffer, b"\xAC\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("SERVNAME,\u{20AC}"));
}
#[test]
fn test_line_buffer_bare_lf_is_tolerated() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"PROBE\nLOOP,0\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("PROBE"));
assert_eq!(line(&mut buffer).as_deref(), Some("LOOP,0"));
}
#[test]
fn test_line_buffer_empty_lines_are_dropped() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"\r\n\r\nEND,31,bye\r\n\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("END,31,bye"));
assert!(buffer.pop_line().is_none());
}
#[test]
fn test_line_buffer_terminated_push_leaves_no_partial() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"SYNC,12\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("SYNC,12"));
assert_eq!(buffer.discard_partial(), None);
}
#[test]
fn test_line_buffer_unterminated_tail_is_discarded_not_promoted() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"LOOP,0");
assert!(buffer.pop_line().is_none());
assert_eq!(buffer.discard_partial(), Some("LOOP,0".len()));
assert_eq!(buffer.discard_partial(), None);
assert!(
buffer.pop_line().is_none(),
"the fragment must not become a line"
);
}
#[test]
fn test_line_buffer_refuses_an_unterminated_fragment_past_its_limit() {
let mut buffer = LineBuffer::default();
let chunk = vec![b'x'; 1024 * 1024];
let mut outcome = Ok(());
for _ in 0..=(MAX_LINE_BYTES / chunk.len()) {
outcome = buffer.push_bytes(&chunk);
if outcome.is_err() {
break;
}
}
match outcome {
Err(TransportError::Capacity { limit_bytes, .. }) => {
assert_eq!(limit_bytes, MAX_LINE_BYTES);
}
other => panic!("expected a capacity refusal, got {other:?}"),
}
assert!(buffer.pop_line().is_none());
assert_eq!(buffer.discard_partial(), None);
}
#[test]
fn test_line_buffer_handles_many_lines_in_one_push() {
let mut buffer = LineBuffer::default();
let mut message = Vec::new();
for index in 0..10_000 {
message.extend_from_slice(format!("PROBE,{index}\r\n").as_bytes());
}
push(&mut buffer, &message);
for index in 0..10_000 {
assert_eq!(line(&mut buffer), Some(format!("PROBE,{index}")));
}
assert!(buffer.pop_line().is_none());
assert_eq!(buffer.discard_partial(), None);
}
#[test]
fn test_line_buffer_preserves_line_payload_verbatim() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"U,3,1,a%2Cb|#|^Pxyz\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("U,3,1,a%2Cb|#|^Pxyz"));
}
#[test]
fn test_line_buffer_invalid_utf8_line_is_malformed_after_the_valid_one() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"SYNC,1\r\n\xFF\xFE\r\n");
assert_eq!(line(&mut buffer).as_deref(), Some("SYNC,1"));
assert!(matches!(
buffer.pop_line(),
Some(Err(TransportError::MalformedFrame { .. }))
));
}
#[test]
fn test_line_buffer_flush_partial_yields_the_unterminated_tail() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"REQOK,7");
assert!(buffer.pop_line().is_none(), "no terminator yet");
assert_eq!(
buffer
.flush_partial()
.map(|result| result.expect("valid UTF-8")),
Some("REQOK,7".to_owned())
);
assert!(buffer.flush_partial().is_none());
}
#[test]
fn test_line_buffer_flush_partial_ignores_a_bare_terminator() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"\r");
assert!(buffer.flush_partial().is_none());
}
#[test]
fn test_line_buffer_clear_discards_everything() {
let mut buffer = LineBuffer::default();
push(&mut buffer, b"CONOK,S1,50000,5000,*\r\nSERVNAME");
buffer.clear();
assert!(buffer.pop_line().is_none());
assert_eq!(buffer.discard_partial(), None);
}
#[test]
fn test_redact_removes_userinfo() {
assert_eq!(
redact("wss://user:secret@push.example.com/lightstreamer"),
"wss://***@push.example.com/lightstreamer"
);
}
#[test]
fn test_redact_leaves_a_plain_address_alone() {
assert_eq!(
redact("https://push.example.com/lightstreamer"),
"https://push.example.com/lightstreamer"
);
}
}