link_cli/protocol/
protocols.rs1use 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
16pub trait LinoProtocol: Debug + Send + Sync {
18 fn write_document(
20 &self,
21 writer: &mut dyn Write,
22 document: &[LiNo<String>],
23 ) -> ProtocolResult<()>;
24
25 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#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
50pub struct TextLinoProtocol {
51 pub limits: DecodeLimits,
53}
54
55impl TextLinoProtocol {
56 pub fn new() -> Self {
58 Self::default()
59 }
60
61 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 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 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#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
147pub struct BinaryLinoProtocol {
148 pub options: BinaryLinoOptions,
150 pub limits: DecodeLimits,
152}
153
154impl BinaryLinoProtocol {
155 pub fn new() -> Self {
157 Self::default()
158 }
159
160 pub fn with_options(options: BinaryLinoOptions) -> Self {
162 Self {
163 options,
164 ..Self::default()
165 }
166 }
167
168 pub fn encode(&self, document: &[LiNo<String>]) -> ProtocolResult<Vec<u8>> {
170 encode_document(document, self.options)?.to_bytes()
171 }
172
173 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#[derive(Clone, Copy, Debug, PartialEq, Eq)]
199pub enum MessageFormat {
200 Text,
202 Binary(BinaryLinoOptions),
204}
205
206impl MessageFormat {
207 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
216pub fn is_binary_start(byte: u8) -> bool {
218 byte & 0xF0 == BINARY_VERSION_1
219}
220
221pub 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#[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 pub fn new(stream: S, protocol: P) -> Self {
256 Self {
257 stream: BufReader::new(stream),
258 protocol,
259 }
260 }
261
262 pub fn protocol(&self) -> &P {
264 &self.protocol
265 }
266
267 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 pub fn receive(&mut self) -> ProtocolResult<Option<LinoDocument>> {
279 self.protocol.read_document(&mut self.stream)
280 }
281
282 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 pub fn into_inner(self) -> S {
291 self.stream.into_inner()
292 }
293}