Skip to main content

link_cli/protocol/
protocols.rs

1//! Interchangeable message protocols for LiNo documents.
2//!
3//! Both [`TextLinoProtocol`] and [`BinaryLinoProtocol`] implement
4//! [`LinoProtocol`], so code written against the trait switches between them
5//! by swapping one value. [`LinoConnection`] decorates any byte stream (a
6//! `TcpStream`, a pipe, an in-memory buffer) with a protocol.
7
8use super::error::{ProtocolError, ProtocolResult};
9use super::format::{format_document, parse_document};
10use super::mapping::{decode_document, encode_document, BinaryLinoOptions, LinoDocument};
11use super::packet::{DecodeLimits, LinksPacket, BINARY_VERSION_1};
12use links_notation::LiNo;
13use std::fmt::Debug;
14use std::io::{BufRead, BufReader, Read, Write};
15
16/// Reads and writes whole LiNo documents on a byte stream.
17pub trait LinoProtocol: Debug + Send + Sync {
18    /// Writes one message carrying `document`.
19    fn write_document(
20        &self,
21        writer: &mut dyn Write,
22        document: &[LiNo<String>],
23    ) -> ProtocolResult<()>;
24
25    /// Reads one message; `Ok(None)` when the stream ended cleanly.
26    fn read_document(&self, reader: &mut dyn BufRead) -> ProtocolResult<Option<LinoDocument>>;
27}
28
29impl<P: LinoProtocol + ?Sized> LinoProtocol for Box<P> {
30    fn write_document(
31        &self,
32        writer: &mut dyn Write,
33        document: &[LiNo<String>],
34    ) -> ProtocolResult<()> {
35        (**self).write_document(writer, document)
36    }
37
38    fn read_document(&self, reader: &mut dyn BufRead) -> ProtocolResult<Option<LinoDocument>> {
39        (**self).read_document(reader)
40    }
41}
42
43/// UTF-8 LiNo text, one message per block of lines ended by a line holding
44/// only `.`. Lines starting with `.` get one extra `.` (SMTP dot-stuffing).
45///
46/// Lines may end with `\n` or `\r\n`; a carriage return right before a line
47/// feed is dropped, so a quoted reference holding `\r\n` arrives as `\n`
48/// (use [`BinaryLinoProtocol`] to carry such references exactly).
49#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
50pub struct TextLinoProtocol {
51    /// Limits applied when reading.
52    pub limits: DecodeLimits,
53}
54
55impl TextLinoProtocol {
56    /// A text protocol with default limits.
57    pub fn new() -> Self {
58        Self::default()
59    }
60
61    /// Frames already formatted LiNo text as one message.
62    pub fn write_text(writer: &mut dyn Write, text: &str) -> ProtocolResult<()> {
63        let mut out = String::with_capacity(text.len() + 4);
64        if !text.is_empty() {
65            for line in text.split('\n') {
66                if line.starts_with('.') {
67                    out.push('.');
68                }
69                out.push_str(line);
70                out.push('\n');
71            }
72        }
73        out.push_str(".\n");
74        writer.write_all(out.as_bytes())?;
75        writer.flush()?;
76        Ok(())
77    }
78
79    /// Reads one framed message as raw text, without parsing it.
80    pub fn read_text(&self, reader: &mut dyn BufRead) -> ProtocolResult<Option<String>> {
81        let mut text = Vec::new();
82        let mut line = Vec::new();
83        let mut first = true;
84        loop {
85            line.clear();
86            let budget = self.limits.max_text_bytes.saturating_sub(text.len()) as u64 + 2;
87            let read = reader.take(budget).read_until(b'\n', &mut line)?;
88            if read == 0 {
89                if first {
90                    return Ok(None);
91                }
92                return Err(ProtocolError::malformed(
93                    "stream ended before the '.' terminator line",
94                ));
95            }
96            if line.last() != Some(&b'\n') {
97                if read as u64 >= budget {
98                    return Err(ProtocolError::LimitExceeded(format!(
99                        "text message longer than {} bytes",
100                        self.limits.max_text_bytes
101                    )));
102                }
103                return Err(ProtocolError::malformed(
104                    "stream ended before the '.' terminator line",
105                ));
106            }
107            line.pop();
108            // CRLF line endings (telnet, netcat on Windows) are accepted.
109            if line.last() == Some(&b'\r') {
110                line.pop();
111            }
112            if line == b"." {
113                break;
114            }
115            if !first {
116                text.push(b'\n');
117            }
118            first = false;
119            let content = line.strip_prefix(b".").unwrap_or(&line);
120            text.extend_from_slice(content);
121        }
122        String::from_utf8(text)
123            .map(Some)
124            .map_err(|_| ProtocolError::malformed("text message is not valid UTF-8"))
125    }
126}
127
128impl LinoProtocol for TextLinoProtocol {
129    fn write_document(
130        &self,
131        writer: &mut dyn Write,
132        document: &[LiNo<String>],
133    ) -> ProtocolResult<()> {
134        Self::write_text(writer, &format_document(document))
135    }
136
137    fn read_document(&self, reader: &mut dyn BufRead) -> ProtocolResult<Option<LinoDocument>> {
138        self.read_text(reader)?
139            .map(|text| parse_document(&text))
140            .transpose()
141    }
142}
143
144/// Binary links packets (see [`crate::protocol::packet`]); every message is
145/// self-delimiting, so no extra framing is needed.
146#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
147pub struct BinaryLinoProtocol {
148    /// Optional features used when writing.
149    pub options: BinaryLinoOptions,
150    /// Limits applied when reading.
151    pub limits: DecodeLimits,
152}
153
154impl BinaryLinoProtocol {
155    /// A binary protocol with every optional feature off.
156    pub fn new() -> Self {
157        Self::default()
158    }
159
160    /// A binary protocol with the given optional features.
161    pub fn with_options(options: BinaryLinoOptions) -> Self {
162        Self {
163            options,
164            ..Self::default()
165        }
166    }
167
168    /// Encodes a document into packet bytes.
169    pub fn encode(&self, document: &[LiNo<String>]) -> ProtocolResult<Vec<u8>> {
170        encode_document(document, self.options)?.to_bytes()
171    }
172
173    /// Decodes packet bytes into a document.
174    pub fn decode(&self, bytes: &[u8]) -> ProtocolResult<LinoDocument> {
175        decode_document(&LinksPacket::from_bytes(bytes, &self.limits)?, &self.limits)
176    }
177}
178
179impl LinoProtocol for BinaryLinoProtocol {
180    fn write_document(
181        &self,
182        writer: &mut dyn Write,
183        document: &[LiNo<String>],
184    ) -> ProtocolResult<()> {
185        writer.write_all(&self.encode(document)?)?;
186        writer.flush()?;
187        Ok(())
188    }
189
190    fn read_document(&self, reader: &mut dyn BufRead) -> ProtocolResult<Option<LinoDocument>> {
191        LinksPacket::read_from(reader, &self.limits)?
192            .map(|packet| decode_document(&packet, &self.limits))
193            .transpose()
194    }
195}
196
197/// The wire format a message arrived in, so a reply can use the same one.
198#[derive(Clone, Copy, Debug, PartialEq, Eq)]
199pub enum MessageFormat {
200    /// [`TextLinoProtocol`].
201    Text,
202    /// [`BinaryLinoProtocol`] with the options read from the packet header.
203    Binary(BinaryLinoOptions),
204}
205
206impl MessageFormat {
207    /// A protocol that writes messages in this format.
208    pub fn protocol(self, limits: DecodeLimits) -> Box<dyn LinoProtocol> {
209        match self {
210            MessageFormat::Text => Box::new(TextLinoProtocol { limits }),
211            MessageFormat::Binary(options) => Box::new(BinaryLinoProtocol { options, limits }),
212        }
213    }
214}
215
216/// True when `byte` starts a binary message rather than a text one.
217pub fn is_binary_start(byte: u8) -> bool {
218    byte & 0xF0 == BINARY_VERSION_1
219}
220
221/// Reads one message in whichever protocol the peer used.
222pub fn read_any_document(
223    reader: &mut dyn BufRead,
224    limits: &DecodeLimits,
225) -> ProtocolResult<Option<(LinoDocument, MessageFormat)>> {
226    let Some(&first) = reader.fill_buf()?.first() else {
227        return Ok(None);
228    };
229    if !is_binary_start(first) {
230        let text = TextLinoProtocol { limits: *limits };
231        return Ok(text
232            .read_document(reader)?
233            .map(|document| (document, MessageFormat::Text)));
234    }
235    LinksPacket::read_from(reader, limits)?
236        .map(|packet| {
237            let document = decode_document(&packet, limits)?;
238            Ok((
239                document,
240                MessageFormat::Binary(BinaryLinoOptions::of_packet(&packet)),
241            ))
242        })
243        .transpose()
244}
245
246/// A byte stream decorated with a [`LinoProtocol`].
247#[derive(Debug)]
248pub struct LinoConnection<S: Read + Write, P: LinoProtocol> {
249    stream: BufReader<S>,
250    protocol: P,
251}
252
253impl<S: Read + Write, P: LinoProtocol> LinoConnection<S, P> {
254    /// Wraps `stream`.
255    pub fn new(stream: S, protocol: P) -> Self {
256        Self {
257            stream: BufReader::new(stream),
258            protocol,
259        }
260    }
261
262    /// The protocol in use.
263    pub fn protocol(&self) -> &P {
264        &self.protocol
265    }
266
267    /// Sends one document.
268    pub fn send(&mut self, document: &[LiNo<String>]) -> ProtocolResult<()> {
269        let mut buffer = Vec::new();
270        self.protocol.write_document(&mut buffer, document)?;
271        let stream = self.stream.get_mut();
272        stream.write_all(&buffer)?;
273        stream.flush()?;
274        Ok(())
275    }
276
277    /// Receives one document; `Ok(None)` when the peer closed the stream.
278    pub fn receive(&mut self) -> ProtocolResult<Option<LinoDocument>> {
279        self.protocol.read_document(&mut self.stream)
280    }
281
282    /// Sends `document` and waits for the reply.
283    pub fn request(&mut self, document: &[LiNo<String>]) -> ProtocolResult<LinoDocument> {
284        self.send(document)?;
285        self.receive()?
286            .ok_or_else(|| ProtocolError::malformed("connection closed before the reply"))
287    }
288
289    /// Unwraps the underlying stream.
290    pub fn into_inner(self) -> S {
291        self.stream.into_inner()
292    }
293}