Skip to main content

tabnas_render/
text.rs

1//! Text output: where rendered fragments go.
2//!
3//! A renderer produces many small fragments (a quote, a field, a comma) and
4//! must never hold a whole document. [`TextOut`] is the boundary: a
5//! fragment in, a failure out. [`WriteOut`] coalesces fragments to a byte
6//! budget before they reach an [`io::Write`], so a renderer can write a
7//! character at a time without paying a system call for each, and it is
8//! where the output limit and the output-bytes metric live, because it is
9//! the one stage that knows what actually left. [`Join`] and
10//! [`ReplaceText`] are the two text combinators whose correctness depends
11//! on the difference between a logical item and a transport chunk, which is
12//! why they live beside the writer rather than in the language that uses
13//! them.
14
15use std::io;
16use std::sync::Arc;
17
18use tabnas_transduce::{Fail, Limits, Metrics};
19
20/// A consumer of text fragments.
21///
22/// Fragments arrive in order and are concatenated; where the boundaries
23/// fall carries no meaning. `flush` pushes everything held so far to the
24/// final destination, and a renderer calls it exactly once, at the end of
25/// the protocol it renders, so that a document that failed half way is not
26/// flushed as if it were whole.
27pub trait TextOut {
28    fn write_str(&mut self, s: &str) -> Result<(), Fail>;
29    fn flush(&mut self) -> Result<(), Fail>;
30
31    /// Whether any text has reached the final destination, so that a
32    /// failure found now leaves partial output behind. A renderer asks this
33    /// when it fails and reports `committed_output` from the answer, which
34    /// is how a host knows to print `output: "partial"` rather than
35    /// `"none"`. The default is the conservative answer for an output that
36    /// cannot tell: whatever the renderer handed over may be out.
37    /// [`WriteOut`] answers exactly, from the bytes its writer received; a
38    /// fragment that is still buffered is not committed, and `into_inner`
39    /// drops it rather than sending it after the fact.
40    fn has_committed(&self) -> bool {
41        true
42    }
43}
44
45impl<O: TextOut + ?Sized> TextOut for &mut O {
46    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
47        (**self).write_str(s)
48    }
49
50    fn flush(&mut self) -> Result<(), Fail> {
51        (**self).flush()
52    }
53
54    fn has_committed(&self) -> bool {
55        (**self).has_committed()
56    }
57}
58
59impl<O: TextOut + ?Sized> TextOut for Box<O> {
60    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
61        (**self).write_str(s)
62    }
63
64    fn flush(&mut self) -> Result<(), Fail> {
65        (**self).flush()
66    }
67
68    fn has_committed(&self) -> bool {
69        (**self).has_committed()
70    }
71}
72
73/// The default coalescing budget of a [`WriteOut`]: large enough that a
74/// write per budget is negligible next to the parse, small enough to be
75/// invisible in a process's memory.
76pub const DEFAULT_BUDGET: usize = 32 * 1024;
77
78/// Coalesces fragments and writes them to an [`io::Write`].
79///
80/// Retention is bounded by the budget: the buffer never holds more than
81/// `budget` bytes, and a fragment at least as large as the budget goes to
82/// the writer directly, after whatever was buffered before it. The output
83/// limit is checked on every fragment BEFORE it is accepted, counting the
84/// bytes buffered as well as the bytes written, so a run that would exceed
85/// `max_output_bytes` fails without emitting the fragment that crossed the
86/// line. `output_bytes` in the shared [`Metrics`] counts bytes the writer
87/// accepted, which is what "written to the output" means to a caller
88/// reading the metrics after a failure; like `committed()`, it is kept
89/// per `write` call, so the part of a buffer a writer took before failing
90/// is counted.
91pub struct WriteOut<W: io::Write> {
92    writer: W,
93    buf: Vec<u8>,
94    budget: usize,
95    max_output_bytes: Option<u64>,
96    metrics: Option<Arc<Metrics>>,
97    /// Bytes accepted: buffered or written.
98    accepted: u64,
99    /// Bytes the writer accepted, counted write by write rather than
100    /// buffer by buffer, so a writer that took part of a buffer and then
101    /// failed is counted for the part it took. Nonzero means a later
102    /// failure finds committed output.
103    committed: u64,
104}
105
106impl<W: io::Write> WriteOut<W> {
107    pub fn new(writer: W) -> Self {
108        WriteOut {
109            writer,
110            buf: Vec::new(),
111            budget: DEFAULT_BUDGET,
112            max_output_bytes: None,
113            metrics: None,
114            accepted: 0,
115            committed: 0,
116        }
117    }
118
119    /// The coalescing budget in bytes. Zero means every fragment is written
120    /// as it arrives, which is what a test of the writer's ordering wants.
121    pub fn with_budget(mut self, budget: usize) -> Self {
122        self.budget = budget;
123        self
124    }
125
126    /// Enforce `limits.max_output_bytes`; the other limits belong to the
127    /// stages upstream.
128    pub fn with_limits(mut self, limits: &Limits) -> Self {
129        self.max_output_bytes = limits.max_output_bytes;
130        self
131    }
132
133    /// Count `output_bytes` into these metrics.
134    pub fn with_metrics(mut self, metrics: Arc<Metrics>) -> Self {
135        self.metrics = Some(metrics);
136        self
137    }
138
139    /// Bytes accepted so far, buffered or written.
140    pub fn accepted(&self) -> u64 {
141        self.accepted
142    }
143
144    /// Bytes the writer accepted, including the bytes of a short write
145    /// that a failure cut off: after `OUTPUT_FAILED` this is what the
146    /// writer holds, not the buffers that were sent whole.
147    pub fn committed(&self) -> u64 {
148        self.committed
149    }
150
151    /// Hand the writer back WITHOUT flushing. Whatever the buffer still
152    /// holds is dropped, so the writer holds exactly the bytes `committed()`
153    /// counts, a short write before a failure included: a document that
154    /// failed before its `End` does not reach the writer on the way out,
155    /// which is what a failure that reported no committed output promised
156    /// the host. A renderer flushes once, at its `End`, and a caller that
157    /// wants a partial output anyway calls `flush` first, knowingly.
158    pub fn into_inner(self) -> W {
159        self.writer
160    }
161
162    fn fail_io(&self, e: io::Error) -> Fail {
163        let f = Fail::output(format!("writing the output failed: {e}"));
164        if self.committed > 0 {
165            f.committed()
166        } else {
167            f
168        }
169    }
170
171    /// Write `bytes` through `write`, as `write_all` does, but counting
172    /// every write the writer accepted before going on to the next. A
173    /// writer may take part of a buffer and then fail (a short write to a
174    /// full disk, or up to a file-size limit): `write_all` would report
175    /// only the error, and a counter kept per buffer would then say the
176    /// writer received nothing while the kernel holds the part it took.
177    /// Counting per write keeps `committed()` equal to what the writer
178    /// holds, whichever write failed. `Interrupted` is retried and `Ok(0)`
179    /// is `WriteZero`, as in `write_all`.
180    fn send(&mut self, mut bytes: &[u8]) -> Result<(), Fail> {
181        while !bytes.is_empty() {
182            match self.writer.write(bytes) {
183                Ok(0) => {
184                    self.buf.clear();
185                    return Err(self.fail_io(io::Error::new(
186                        io::ErrorKind::WriteZero,
187                        "failed to write whole buffer",
188                    )));
189                }
190                Ok(n) => {
191                    self.committed += n as u64;
192                    if let Some(m) = &self.metrics {
193                        Metrics::add(&m.output_bytes, n as u64);
194                    }
195                    bytes = &bytes[n..];
196                }
197                Err(e) if e.kind() == io::ErrorKind::Interrupted => {}
198                Err(e) => {
199                    // The buffer is not retried: a failed writer is done,
200                    // and the caller learns how many bytes it accepted
201                    // before that, the short write included.
202                    self.buf.clear();
203                    return Err(self.fail_io(e));
204                }
205            }
206        }
207        Ok(())
208    }
209
210    fn drain(&mut self) -> Result<(), Fail> {
211        if self.buf.is_empty() {
212            return Ok(());
213        }
214        let pending = std::mem::take(&mut self.buf);
215        let r = self.send(&pending);
216        // Keep the allocation: the buffer refills up to the budget again.
217        if r.is_ok() {
218            self.buf = pending;
219            self.buf.clear();
220        }
221        r
222    }
223}
224
225impl<W: io::Write> TextOut for WriteOut<W> {
226    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
227        let len = s.len() as u64;
228        if let Some(max) = self.max_output_bytes {
229            if self.accepted.saturating_add(len) > max {
230                let f = Fail::limit(
231                    "max_output_bytes",
232                    max,
233                    format!(
234                        "the output would exceed {max} bytes: {} written, {len} more",
235                        self.accepted
236                    ),
237                );
238                return Err(if self.committed > 0 { f.committed() } else { f });
239            }
240        }
241        if self.buf.len() + s.len() > self.budget {
242            self.drain()?;
243        }
244        if s.len() >= self.budget {
245            self.send(s.as_bytes())?;
246        } else {
247            self.buf.extend_from_slice(s.as_bytes());
248        }
249        self.accepted += len;
250        Ok(())
251    }
252
253    fn flush(&mut self) -> Result<(), Fail> {
254        self.drain()?;
255        self.writer.flush().map_err(|e| self.fail_io(e))
256    }
257
258    fn has_committed(&self) -> bool {
259        self.committed > 0
260    }
261}
262
263/// A [`TextOut`] that keeps the text, for tests and small results.
264#[derive(Clone, Debug, Default, PartialEq, Eq)]
265pub struct StringOut(pub String);
266
267impl StringOut {
268    pub fn new() -> Self {
269        StringOut::default()
270    }
271
272    pub fn as_str(&self) -> &str {
273        &self.0
274    }
275
276    pub fn into_string(self) -> String {
277        self.0
278    }
279}
280
281impl TextOut for StringOut {
282    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
283        self.0.push_str(s);
284        Ok(())
285    }
286
287    fn flush(&mut self) -> Result<(), Fail> {
288        Ok(())
289    }
290
291    /// The string is the destination, so its text is committed as soon as
292    /// it is there.
293    fn has_committed(&self) -> bool {
294        !self.0.is_empty()
295    }
296}
297
298/// Writes a separator between logical items.
299///
300/// An item is what lies between [`Join::item_start`] and
301/// [`Join::item_end`]; it may be written in any number of fragments, or in
302/// none, and an empty item is still an item, so `["", ""]` joined with `,`
303/// is `,`. A fragment written outside an item is an item of its own, which
304/// is the common case of one text per element and needs no markers. The
305/// separator goes before every item but the first, never between the
306/// fragments of one item: that distinction is the whole reason this type
307/// exists, because a chunked transport must not change the text.
308pub struct Join<O: TextOut> {
309    out: O,
310    separator: Box<str>,
311    items: u64,
312    in_item: bool,
313}
314
315impl<O: TextOut> Join<O> {
316    pub fn new(out: O, separator: impl Into<Box<str>>) -> Self {
317        Join {
318            out,
319            separator: separator.into(),
320            items: 0,
321            in_item: false,
322        }
323    }
324
325    /// Begin an item: the separator is written now if an item came before.
326    /// Starting an item inside an item is `PROTOCOL_ORDER_ERROR`.
327    pub fn item_start(&mut self) -> Result<(), Fail> {
328        if self.in_item {
329            return Err(Fail::protocol(
330                "join: an item started inside an item that has not ended",
331            ));
332        }
333        if self.items > 0 && !self.separator.is_empty() {
334            self.out.write_str(&self.separator)?;
335        }
336        self.items += 1;
337        self.in_item = true;
338        Ok(())
339    }
340
341    /// End the current item. Ending when no item is open is
342    /// `PROTOCOL_ORDER_ERROR`.
343    pub fn item_end(&mut self) -> Result<(), Fail> {
344        if !self.in_item {
345            return Err(Fail::protocol("join: an item ended when none was open"));
346        }
347        self.in_item = false;
348        Ok(())
349    }
350
351    /// Items begun so far.
352    pub fn items(&self) -> u64 {
353        self.items
354    }
355
356    pub fn into_inner(self) -> O {
357        self.out
358    }
359}
360
361impl<O: TextOut> TextOut for Join<O> {
362    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
363        if self.in_item {
364            return self.out.write_str(s);
365        }
366        self.item_start()?;
367        self.out.write_str(s)?;
368        self.item_end()
369    }
370
371    /// Flushes the output beneath; an open item stays open, since a flush
372    /// is about transport and an item is about meaning.
373    fn flush(&mut self) -> Result<(), Fail> {
374        self.out.flush()
375    }
376
377    fn has_committed(&self) -> bool {
378        self.out.has_committed()
379    }
380}
381
382/// Concatenation: a [`Join`] with no separator, under the name the design
383/// brief and the language give it.
384///
385/// Every fragment is appended as it is. The item markers are accepted and
386/// counted all the same, so the interpreter's `concat` and `join` share
387/// one shape and a program can move between them without the calls around
388/// them changing; that shared shape is the whole reason for a type where a
389/// bare output would do.
390pub struct Concat<O: TextOut>(Join<O>);
391
392impl<O: TextOut> Concat<O> {
393    pub fn new(out: O) -> Self {
394        Concat(Join::new(out, ""))
395    }
396
397    /// Begin an item; see [`Join::item_start`].
398    pub fn item_start(&mut self) -> Result<(), Fail> {
399        self.0.item_start()
400    }
401
402    /// End the current item; see [`Join::item_end`].
403    pub fn item_end(&mut self) -> Result<(), Fail> {
404        self.0.item_end()
405    }
406
407    /// Items begun so far.
408    pub fn items(&self) -> u64 {
409        self.0.items()
410    }
411
412    pub fn into_inner(self) -> O {
413        self.0.into_inner()
414    }
415}
416
417impl<O: TextOut> TextOut for Concat<O> {
418    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
419        self.0.write_str(s)
420    }
421
422    fn flush(&mut self) -> Result<(), Fail> {
423        self.0.flush()
424    }
425
426    fn has_committed(&self) -> bool {
427        self.0.has_committed()
428    }
429}
430
431/// Replaces every occurrence of a fixed literal, across fragment
432/// boundaries.
433///
434/// The text is treated as one string however it is chunked, with the same
435/// left-to-right, non-overlapping matches as `str::replace`. To do that
436/// without holding the text, at most `literal.len() - 1` bytes are carried
437/// from one fragment to the next: the longest suffix of what has been seen
438/// that could still begin a match. `flush` writes the carry out, because
439/// nothing can complete it once the caller has declared the text at a
440/// boundary; the renderer's single flush at the end of a document is what
441/// makes that safe. An empty literal matches nothing and the text passes
442/// through unchanged, the one well-defined meaning it can have in a stream.
443pub struct ReplaceText<O: TextOut> {
444    out: O,
445    from: Box<str>,
446    to: Box<str>,
447    carry: String,
448}
449
450impl<O: TextOut> ReplaceText<O> {
451    pub fn new(out: O, from: impl Into<Box<str>>, to: impl Into<Box<str>>) -> Self {
452        ReplaceText {
453            out,
454            from: from.into(),
455            to: to.into(),
456            carry: String::new(),
457        }
458    }
459
460    pub fn into_inner(self) -> O {
461        self.out
462    }
463
464    /// The longest proper prefix of the literal that `rest` ends with, in
465    /// bytes; zero when there is none. Only char boundaries of the literal
466    /// are tried: a match ending inside a character is impossible between
467    /// two valid strings, and skipping them keeps every slice below safe.
468    fn pending_len(&self, rest: &str) -> usize {
469        let max = rest.len().min(self.from.len() - 1);
470        (1..=max)
471            .rev()
472            .find(|&k| self.from.is_char_boundary(k) && rest.ends_with(&self.from[..k]))
473            .unwrap_or(0)
474    }
475
476    fn scan(&mut self, text: &str) -> Result<(), Fail> {
477        let mut rest = text;
478        while let Some(i) = rest.find(&*self.from) {
479            if i > 0 {
480                self.out.write_str(&rest[..i])?;
481            }
482            if !self.to.is_empty() {
483                self.out.write_str(&self.to)?;
484            }
485            rest = &rest[i + self.from.len()..];
486        }
487        let keep = self.pending_len(rest);
488        // `pending_len` returns the length of a suffix of `rest` that equals a
489        // prefix of the literal beginning with a character's first byte, so
490        // the split lands on a char boundary; the checked form keeps the
491        // guarantee without a panic path.
492        match rest.split_at_checked(rest.len() - keep) {
493            Some((emit, pending)) => {
494                if !emit.is_empty() {
495                    self.out.write_str(emit)?;
496                }
497                self.carry.clear();
498                self.carry.push_str(pending);
499            }
500            None => {
501                self.out.write_str(rest)?;
502                self.carry.clear();
503            }
504        }
505        Ok(())
506    }
507}
508
509impl<O: TextOut> TextOut for ReplaceText<O> {
510    fn write_str(&mut self, s: &str) -> Result<(), Fail> {
511        if self.from.is_empty() {
512            return self.out.write_str(s);
513        }
514        if self.carry.is_empty() {
515            self.scan(s)
516        } else {
517            let mut text = std::mem::take(&mut self.carry);
518            text.push_str(s);
519            self.scan(&text)
520        }
521    }
522
523    fn flush(&mut self) -> Result<(), Fail> {
524        if !self.carry.is_empty() {
525            let pending = std::mem::take(&mut self.carry);
526            self.out.write_str(&pending)?;
527        }
528        self.out.flush()
529    }
530
531    /// The carry has not gone anywhere; only the output beneath knows.
532    fn has_committed(&self) -> bool {
533        self.out.has_committed()
534    }
535}
536
537#[cfg(test)]
538mod tests {
539    use super::*;
540    use tabnas_transduce::Code;
541
542    /// A writer that records each `write` as one chunk, so coalescing is
543    /// observable, and fails after a set number of bytes when asked. It
544    /// takes a buffer whole or refuses it whole: the all-or-nothing writer.
545    #[derive(Default, Debug)]
546    struct Chunks {
547        chunks: Vec<Vec<u8>>,
548        flushes: usize,
549        fail_after: Option<usize>,
550    }
551
552    impl io::Write for Chunks {
553        fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
554            let so_far: usize = self.chunks.iter().map(Vec::len).sum();
555            if self.fail_after.is_some_and(|n| so_far + buf.len() > n) {
556                return Err(io::Error::other("disk full"));
557            }
558            self.chunks.push(buf.to_vec());
559            Ok(buf.len())
560        }
561
562        fn flush(&mut self) -> io::Result<()> {
563            self.flushes += 1;
564            Ok(())
565        }
566    }
567
568    fn joined(chunks: &[Vec<u8>]) -> String {
569        String::from_utf8(chunks.concat()).unwrap()
570    }
571
572    #[test]
573    fn fragments_coalesce_up_to_the_budget_and_the_buffer_never_exceeds_it() {
574        let mut out = WriteOut::new(Chunks::default()).with_budget(8);
575        out.write_str("abc").unwrap();
576        out.write_str("def").unwrap();
577        out.write_str("gh").unwrap();
578        // Exactly the budget: still held.
579        assert!(out.writer.chunks.is_empty());
580        out.write_str("i").unwrap();
581        assert_eq!(out.writer.chunks, vec![b"abcdefgh".to_vec()]);
582        assert_eq!(out.committed(), 8);
583        assert_eq!(out.accepted(), 9);
584        // A fragment at least as large as the budget bypasses the buffer,
585        // after what was buffered before it.
586        out.write_str("0123456789").unwrap();
587        assert_eq!(
588            out.writer.chunks,
589            vec![b"abcdefgh".to_vec(), b"i".to_vec(), b"0123456789".to_vec()]
590        );
591        out.write_str("z").unwrap();
592        out.flush().unwrap();
593        let w = out.into_inner();
594        assert_eq!(joined(&w.chunks), "abcdefghi0123456789z");
595        assert_eq!(w.flushes, 1);
596    }
597
598    #[test]
599    fn a_zero_budget_writes_every_fragment_as_it_arrives() {
600        let mut out = WriteOut::new(Chunks::default()).with_budget(0);
601        out.write_str("a").unwrap();
602        out.write_str("bc").unwrap();
603        assert_eq!(out.writer.chunks, vec![b"a".to_vec(), b"bc".to_vec()]);
604    }
605
606    #[test]
607    fn the_output_limit_fails_before_the_fragment_that_would_exceed_it() {
608        let limits = Limits {
609            max_output_bytes: Some(10),
610            ..Limits::default()
611        };
612        let metrics = Metrics::new();
613        let mut out = WriteOut::new(Chunks::default())
614            .with_budget(4)
615            .with_limits(&limits)
616            .with_metrics(Arc::clone(&metrics));
617        out.write_str("hello").unwrap();
618        out.write_str("worl").unwrap();
619        let err = out.write_str("d!").unwrap_err();
620        assert_eq!(err.code, Code::ResourceLimitExceeded);
621        let limit = err.limit.as_ref().unwrap();
622        assert_eq!(limit.name, "max_output_bytes");
623        assert_eq!(limit.value, 10);
624        // "hello" crossed the budget when "worl" arrived, so it was written
625        // before the failure and the failure says so.
626        assert!(err.committed_output);
627        assert_eq!(out.accepted(), 9);
628        // "worl" is as large as the budget, so it went to the writer too;
629        // the writer holds what `committed()` says and nothing more.
630        assert_eq!(out.committed(), 9);
631        let w = out.into_inner();
632        assert_eq!(joined(&w.chunks), "helloworl");
633        assert_eq!(Metrics::get(&metrics.output_bytes), 9);
634    }
635
636    #[test]
637    fn into_inner_after_a_failure_hands_back_exactly_the_committed_bytes() {
638        let limits = Limits {
639            max_output_bytes: Some(5),
640            ..Limits::default()
641        };
642        let mut out = WriteOut::new(Chunks::default())
643            .with_budget(100)
644            .with_limits(&limits);
645        out.write_str("abc").unwrap();
646        let err = out.write_str("xyz").unwrap_err();
647        assert_eq!(err.code, Code::ResourceLimitExceeded);
648        assert!(!err.committed_output);
649        assert_eq!(out.committed(), 0);
650        let w = out.into_inner();
651        assert!(w.chunks.is_empty(), "no committed output means none: {w:?}");
652        assert_eq!(w.flushes, 0);
653    }
654
655    #[test]
656    fn a_caller_that_wants_the_partial_output_flushes_before_into_inner() {
657        let mut out = WriteOut::new(Chunks::default()).with_budget(100);
658        out.write_str("abc").unwrap();
659        out.flush().unwrap();
660        let w = out.into_inner();
661        assert_eq!(joined(&w.chunks), "abc");
662        assert_eq!(w.flushes, 1);
663    }
664
665    #[test]
666    fn a_limit_failure_with_nothing_written_is_not_committed() {
667        let limits = Limits {
668            max_output_bytes: Some(3),
669            ..Limits::default()
670        };
671        let mut out = WriteOut::new(Chunks::default()).with_limits(&limits);
672        let err = out.write_str("abcd").unwrap_err();
673        assert_eq!(err.code, Code::ResourceLimitExceeded);
674        assert!(!err.committed_output);
675        assert_eq!(out.into_inner().chunks, Vec::<Vec<u8>>::new());
676    }
677
678    #[test]
679    fn an_io_error_is_output_failed_and_says_whether_bytes_were_committed() {
680        let writer = Chunks {
681            fail_after: Some(4),
682            ..Chunks::default()
683        };
684        let mut out = WriteOut::new(writer).with_budget(3);
685        out.write_str("abc").unwrap();
686        assert_eq!(out.committed(), 3);
687        out.write_str("de").unwrap();
688        assert_eq!(out.committed(), 3);
689        let err = out.flush().unwrap_err();
690        assert_eq!(err.code, Code::OutputFailed);
691        assert!(err.committed_output);
692        assert!(err.message.contains("disk full"));
693
694        let writer = Chunks {
695            fail_after: Some(0),
696            ..Chunks::default()
697        };
698        let mut out = WriteOut::new(writer).with_budget(0);
699        let err = out.write_str("x").unwrap_err();
700        assert_eq!(err.code, Code::OutputFailed);
701        assert!(!err.committed_output);
702    }
703
704    /// A writer with room for `room` bytes that takes what fits of each
705    /// write, as `io::Write` permits, and fails once it is full: the short
706    /// write before "no space left on device".
707    #[derive(Debug)]
708    struct Cramped {
709        room: usize,
710        taken: Vec<u8>,
711    }
712
713    impl Cramped {
714        fn with_room(room: usize) -> Self {
715            Cramped {
716                room,
717                taken: Vec::new(),
718            }
719        }
720    }
721
722    impl io::Write for Cramped {
723        fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
724            let left = self.room - self.taken.len();
725            if left == 0 {
726                return Err(io::Error::other("disk full"));
727            }
728            let n = buf.len().min(left);
729            self.taken.extend_from_slice(&buf[..n]);
730            Ok(n)
731        }
732
733        fn flush(&mut self) -> io::Result<()> {
734            Ok(())
735        }
736    }
737
738    #[test]
739    fn a_short_write_before_the_failure_counts_the_bytes_the_writer_took() {
740        // Buffered, then flushed: the writer takes three bytes of the six
741        // and fails on the rest.
742        let metrics = Metrics::new();
743        let mut out = WriteOut::new(Cramped::with_room(3))
744            .with_budget(100)
745            .with_metrics(Arc::clone(&metrics));
746        out.write_str("abc").unwrap();
747        out.write_str("def").unwrap();
748        assert_eq!(out.committed(), 0, "still buffered");
749        let err = out.flush().unwrap_err();
750        assert_eq!(err.code, Code::OutputFailed);
751        assert!(err.message.contains("disk full"));
752        assert!(
753            err.committed_output,
754            "three bytes reached the writer before it failed"
755        );
756        assert!(out.has_committed());
757        assert_eq!(out.committed(), 3);
758        assert_eq!(out.accepted(), 6);
759        assert_eq!(Metrics::get(&metrics.output_bytes), 3);
760        let w = out.into_inner();
761        assert_eq!(
762            w.taken, b"abc",
763            "the writer holds exactly committed() bytes"
764        );
765
766        // Written directly: a fragment as large as the budget takes the
767        // same path and is counted the same way.
768        let mut out = WriteOut::new(Cramped::with_room(2)).with_budget(0);
769        let err = out.write_str("abcdef").unwrap_err();
770        assert_eq!(err.code, Code::OutputFailed);
771        assert!(err.committed_output);
772        assert_eq!(out.committed(), 2);
773        assert_eq!(out.accepted(), 0, "the fragment was not accepted");
774        assert_eq!(out.into_inner().taken, b"ab");
775
776        // No room at all: nothing was taken, and the failure says so.
777        let mut out = WriteOut::new(Cramped::with_room(0)).with_budget(0);
778        let err = out.write_str("abc").unwrap_err();
779        assert_eq!(err.code, Code::OutputFailed);
780        assert!(!err.committed_output);
781        assert!(!out.has_committed());
782        assert_eq!(out.committed(), 0);
783        assert!(out.into_inner().taken.is_empty());
784    }
785
786    /// A writer that accepts nothing and reports no error, which
787    /// `io::Write` allows and `write_all` treats as the end of the writer.
788    struct Zero;
789
790    impl io::Write for Zero {
791        fn write(&mut self, _buf: &[u8]) -> io::Result<usize> {
792            Ok(0)
793        }
794
795        fn flush(&mut self) -> io::Result<()> {
796            Ok(())
797        }
798    }
799
800    #[test]
801    fn a_writer_that_takes_nothing_is_write_zero_not_a_spin() {
802        let mut out = WriteOut::new(Zero).with_budget(0);
803        let err = out.write_str("abc").unwrap_err();
804        assert_eq!(err.code, Code::OutputFailed);
805        assert!(err.message.contains("failed to write whole buffer"));
806        assert!(!err.committed_output);
807        assert_eq!(out.committed(), 0);
808    }
809
810    /// A writer that is interrupted before every write it accepts.
811    #[derive(Default)]
812    struct Interrupting {
813        taken: Vec<u8>,
814        interruptions: usize,
815        ready: bool,
816    }
817
818    impl io::Write for Interrupting {
819        fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
820            if !self.ready {
821                self.ready = true;
822                self.interruptions += 1;
823                return Err(io::Error::from(io::ErrorKind::Interrupted));
824            }
825            self.ready = false;
826            self.taken.extend_from_slice(buf);
827            Ok(buf.len())
828        }
829
830        fn flush(&mut self) -> io::Result<()> {
831            Ok(())
832        }
833    }
834
835    #[test]
836    fn an_interrupted_write_is_retried_and_counted_once() {
837        let mut out = WriteOut::new(Interrupting::default()).with_budget(4);
838        out.write_str("abcd").unwrap();
839        out.write_str("efgh").unwrap();
840        out.flush().unwrap();
841        assert_eq!(out.committed(), 8);
842        let w = out.into_inner();
843        assert_eq!(w.taken, b"abcdefgh");
844        assert_eq!(w.interruptions, 2, "one retry per buffer, none counted");
845    }
846
847    /// The environment variable that marks the process running under the
848    /// file-size limit in the test below.
849    #[cfg(unix)]
850    const SHORT_WRITE_CHILD: &str = "TABNAS_RENDER_SHORT_WRITE_CHILD";
851
852    /// The half of the test that runs under the limit: a real file, a
853    /// buffer larger than the limit, one flush. It reports on standard
854    /// output with a `short-write:` line, which the parent reads.
855    #[cfg(unix)]
856    fn short_write_child() {
857        println!("short-write: start");
858        let path = std::env::temp_dir().join(format!(
859            "tabnas-render-short-write-{}.txt",
860            std::process::id()
861        ));
862        let file = std::fs::File::create(&path).expect("create the output file");
863        let metrics = Metrics::new();
864        let mut out = WriteOut::new(file).with_metrics(Arc::clone(&metrics));
865        let text = "0123456789abcdef".repeat(1000);
866        out.write_str(&text).unwrap();
867        assert_eq!(
868            out.committed(),
869            0,
870            "under the default budget it is buffered"
871        );
872        let result = out.flush();
873        let committed = out.committed();
874        let has_committed = out.has_committed();
875        let output_bytes = Metrics::get(&metrics.output_bytes);
876        let file = out.into_inner();
877        let on_disk = file.metadata().map(|m| m.len());
878        drop(file);
879        let _ = std::fs::remove_file(&path);
880        let on_disk = on_disk.expect("the file's length");
881        match result {
882            Ok(()) => println!(
883                "short-write: skipped, no file-size limit was in force ({on_disk} bytes on disk)"
884            ),
885            Err(err) if committed == 0 && on_disk == 0 => println!(
886                "short-write: skipped, the limit refused the write whole: {}",
887                err.message
888            ),
889            Err(err) => {
890                println!(
891                    "short-write: ran committed={committed} on_disk={on_disk} \
892                     output_bytes={output_bytes} code={:?} committed_output={}",
893                    err.code, err.committed_output
894                );
895                assert_eq!(err.code, Code::OutputFailed);
896                assert!(
897                    err.committed_output,
898                    "the kernel took part of the buffer, so output is partial"
899                );
900                assert!(has_committed);
901                assert_eq!(
902                    committed, on_disk,
903                    "committed() must be what the file holds"
904                );
905                assert_eq!(output_bytes, committed);
906                assert!(committed < text.len() as u64, "the limit cut the write");
907            }
908        }
909    }
910
911    /// The short write the reviewer reproduced on a full file system, on a
912    /// real file: the kernel accepts the bytes up to the limit and refuses
913    /// the rest, and the file on disk holds exactly `committed()` bytes.
914    /// `RLIMIT_FSIZE` is the limit that needs no privileges; it is set per
915    /// process by the shell's `ulimit -f`, so the run happens in a child
916    /// process (this binary, this test, with `SHORT_WRITE_CHILD` set), and
917    /// the shell ignores `SIGXFSZ` first so that the write past the limit
918    /// fails with `EFBIG` instead of ending the child. Where no shell can
919    /// impose the limit, the test says so and passes.
920    #[cfg(unix)]
921    #[test]
922    fn a_file_under_a_size_limit_holds_exactly_the_committed_bytes() {
923        use std::process::Command;
924
925        if std::env::var_os(SHORT_WRITE_CHILD).is_some() {
926            short_write_child();
927            return;
928        }
929        let exe = std::env::current_exe().expect("the test binary");
930        let module = module_path!()
931            .split_once("::")
932            .map_or(module_path!(), |(_, m)| m);
933        let name = format!("{module}::a_file_under_a_size_limit_holds_exactly_the_committed_bytes");
934        let output = match Command::new("sh")
935            .arg("-c")
936            .arg(r#"trap "" XFSZ && ulimit -f 8 && exec "$0" "$@""#)
937            .arg(&exe)
938            .args(["--exact", &name, "--nocapture"])
939            .env(SHORT_WRITE_CHILD, "1")
940            .output()
941        {
942            Ok(output) => output,
943            Err(e) => {
944                eprintln!("skipped: no `sh` to impose a file-size limit with ({e})");
945                return;
946            }
947        };
948        let stdout = String::from_utf8_lossy(&output.stdout);
949        let stderr = String::from_utf8_lossy(&output.stderr);
950        let report = stdout.lines().rfind(|l| l.starts_with("short-write: "));
951        match report {
952            Some(line) if line.starts_with("short-write: ran") => {
953                // The child's own assertions ran; its report is worth
954                // seeing under `--nocapture`.
955                println!("{line}");
956                assert!(
957                    output.status.success(),
958                    "the run under the limit failed: {line}\n{stdout}\n{stderr}"
959                );
960            }
961            Some(line) if line.starts_with("short-write: skipped") => eprintln!("{line}"),
962            Some(line) => panic!(
963                "the child stopped after `{line}` with status {}:\n{stdout}\n{stderr}",
964                output.status
965            ),
966            None if stdout.contains("running ") => {
967                panic!("the child ran the harness but not the test:\n{stdout}\n{stderr}")
968            }
969            None => eprintln!(
970                "skipped: the shell could not impose a file-size limit (status {}): {}",
971                output.status,
972                stderr.trim()
973            ),
974        }
975    }
976
977    #[test]
978    fn metrics_count_bytes_handed_to_the_writer() {
979        let metrics = Metrics::new();
980        let mut out = WriteOut::new(Vec::new())
981            .with_budget(100)
982            .with_metrics(Arc::clone(&metrics));
983        out.write_str("twelve bytes").unwrap();
984        assert_eq!(Metrics::get(&metrics.output_bytes), 0);
985        out.flush().unwrap();
986        assert_eq!(Metrics::get(&metrics.output_bytes), 12);
987        assert_eq!(out.into_inner(), b"twelve bytes");
988    }
989
990    #[test]
991    fn has_committed_is_answered_by_the_destination_not_the_buffer() {
992        let mut out = WriteOut::new(Chunks::default()).with_budget(100);
993        assert!(!out.has_committed());
994        out.write_str("abc").unwrap();
995        assert!(!out.has_committed(), "buffered is not committed");
996        out.flush().unwrap();
997        assert!(out.has_committed());
998
999        let mut s = StringOut::new();
1000        assert!(!s.has_committed());
1001        s.write_str("x").unwrap();
1002        assert!(s.has_committed(), "the string is the destination");
1003
1004        // The combinators forward the question; a replacer's carry has
1005        // gone nowhere yet.
1006        let inner = WriteOut::new(Vec::new()).with_budget(100);
1007        let mut r = ReplaceText::new(Join::new(inner, ","), "ab", "");
1008        r.write_str("xa").unwrap();
1009        assert!(!r.has_committed());
1010        r.flush().unwrap();
1011        assert!(r.has_committed());
1012        let by_ref: &mut ReplaceText<_> = &mut r;
1013        assert!(TextOut::has_committed(&by_ref));
1014        let boxed: Box<dyn TextOut> = Box::new(r);
1015        assert!(boxed.has_committed());
1016    }
1017
1018    #[test]
1019    fn string_out_keeps_the_text() {
1020        let mut s = StringOut::new();
1021        s.write_str("a").unwrap();
1022        s.write_str("b").unwrap();
1023        s.flush().unwrap();
1024        assert_eq!(s.as_str(), "ab");
1025        assert_eq!(s.into_string(), "ab");
1026    }
1027
1028    #[test]
1029    fn join_separates_items_not_fragments() {
1030        let mut j = Join::new(StringOut::new(), ", ");
1031        j.item_start().unwrap();
1032        j.write_str("a").unwrap();
1033        j.write_str("b").unwrap();
1034        j.item_end().unwrap();
1035        j.item_start().unwrap();
1036        j.write_str("c").unwrap();
1037        j.item_end().unwrap();
1038        assert_eq!(j.items(), 2);
1039        assert_eq!(j.into_inner().as_str(), "ab, c");
1040    }
1041
1042    #[test]
1043    fn join_counts_empty_items() {
1044        let mut j = Join::new(StringOut::new(), ",");
1045        for _ in 0..3 {
1046            j.item_start().unwrap();
1047            j.item_end().unwrap();
1048        }
1049        j.item_start().unwrap();
1050        j.write_str("x").unwrap();
1051        j.item_end().unwrap();
1052        j.item_start().unwrap();
1053        j.item_end().unwrap();
1054        assert_eq!(j.into_inner().as_str(), ",,,x,");
1055    }
1056
1057    #[test]
1058    fn join_treats_a_fragment_outside_an_item_as_an_item() {
1059        let mut j = Join::new(StringOut::new(), "|");
1060        j.write_str("a").unwrap();
1061        j.write_str("").unwrap();
1062        j.write_str("b").unwrap();
1063        j.flush().unwrap();
1064        assert_eq!(j.into_inner().as_str(), "a||b");
1065    }
1066
1067    #[test]
1068    fn join_with_no_items_writes_nothing() {
1069        let mut j = Join::new(StringOut::new(), ",");
1070        j.flush().unwrap();
1071        assert_eq!(j.into_inner().as_str(), "");
1072    }
1073
1074    #[test]
1075    fn join_rejects_unbalanced_item_markers() {
1076        let mut j = Join::new(StringOut::new(), ",");
1077        assert_eq!(j.item_end().unwrap_err().code, Code::ProtocolOrderError);
1078        j.item_start().unwrap();
1079        assert_eq!(j.item_start().unwrap_err().code, Code::ProtocolOrderError);
1080    }
1081
1082    #[test]
1083    fn concat_appends_items_and_fragments_with_nothing_between_them() {
1084        let mut c = Concat::new(StringOut::new());
1085        c.item_start().unwrap();
1086        c.write_str("a").unwrap();
1087        c.write_str("b").unwrap();
1088        c.item_end().unwrap();
1089        c.item_start().unwrap();
1090        c.item_end().unwrap();
1091        c.write_str("c").unwrap();
1092        c.flush().unwrap();
1093        assert_eq!(c.items(), 3);
1094        assert!(c.has_committed());
1095        assert_eq!(c.into_inner().as_str(), "abc");
1096    }
1097
1098    #[test]
1099    fn concat_keeps_joins_item_discipline() {
1100        let mut c = Concat::new(WriteOut::new(Vec::new()));
1101        assert_eq!(c.item_end().unwrap_err().code, Code::ProtocolOrderError);
1102        c.item_start().unwrap();
1103        assert_eq!(c.item_start().unwrap_err().code, Code::ProtocolOrderError);
1104        c.write_str("x").unwrap();
1105        assert!(!c.has_committed(), "buffered beneath, not yet written");
1106        c.item_end().unwrap();
1107        c.flush().unwrap();
1108        assert_eq!(c.into_inner().into_inner(), b"x");
1109    }
1110
1111    /// Feed `text` to a replacer split at `at`, then flushed, and give the
1112    /// result back.
1113    fn replaced_split(text: &str, at: usize, from: &str, to: &str) -> String {
1114        let mut r = ReplaceText::new(StringOut::new(), from, to);
1115        let (a, b) = text.split_at(at);
1116        r.write_str(a).unwrap();
1117        r.write_str(b).unwrap();
1118        r.flush().unwrap();
1119        r.into_inner().into_string()
1120    }
1121
1122    #[test]
1123    fn replace_matches_str_replace_when_split_at_every_boundary() {
1124        let cases = [
1125            ("abcabc", "abc", "X"),
1126            ("xxabcxxabcxx", "abc", ""),
1127            ("aaaa", "aa", "b"),
1128            ("aaaaa", "aa", "b"),
1129            ("ababab", "aba", "_"),
1130            ("no match here", "zzz", "Y"),
1131            ("abab", "abab", "1"),
1132            ("ab", "abc", "1"),
1133            ("héllo wörld héllo", "héllo", "hi"),
1134            ("日本語日本", "日本", "*"),
1135            ("a\r\nb\r\n", "\r\n", "\n"),
1136        ];
1137        for (text, from, to) in cases {
1138            let want = text.replace(from, to);
1139            for at in (0..=text.len()).filter(|&i| text.is_char_boundary(i)) {
1140                assert_eq!(
1141                    replaced_split(text, at, from, to),
1142                    want,
1143                    "{text:?} split at {at} replacing {from:?}"
1144                );
1145            }
1146        }
1147    }
1148
1149    #[test]
1150    fn replace_across_many_one_character_fragments() {
1151        let text = "the cat sat on the mat with the hat";
1152        let mut r = ReplaceText::new(StringOut::new(), "the", "a");
1153        for c in text.chars() {
1154            let mut buf = [0u8; 4];
1155            r.write_str(c.encode_utf8(&mut buf)).unwrap();
1156        }
1157        r.flush().unwrap();
1158        assert_eq!(r.into_inner().as_str(), text.replace("the", "a"));
1159    }
1160
1161    #[test]
1162    fn replace_never_carries_more_than_the_literal_less_one_byte() {
1163        let mut r = ReplaceText::new(StringOut::new(), "abcd", "");
1164        r.write_str("xxabc").unwrap();
1165        assert_eq!(r.carry, "abc");
1166        assert_eq!(r.out.as_str(), "xx");
1167        r.write_str("ab").unwrap();
1168        assert_eq!(r.carry, "ab");
1169        assert_eq!(r.out.as_str(), "xxabc");
1170        r.write_str("cdab").unwrap();
1171        assert_eq!(r.carry, "ab");
1172        assert_eq!(r.out.as_str(), "xxabc");
1173        r.flush().unwrap();
1174        assert_eq!(r.carry, "");
1175        assert_eq!(r.into_inner().as_str(), "xxabcab");
1176    }
1177
1178    #[test]
1179    fn replace_with_an_empty_literal_passes_text_through() {
1180        let mut r = ReplaceText::new(StringOut::new(), "", "X");
1181        r.write_str("abc").unwrap();
1182        r.flush().unwrap();
1183        assert_eq!(r.into_inner().as_str(), "abc");
1184    }
1185
1186    #[test]
1187    fn replace_flushes_the_carry_at_flush_so_a_later_match_cannot_span_it() {
1188        let mut r = ReplaceText::new(StringOut::new(), "ab", "X");
1189        r.write_str("a").unwrap();
1190        r.flush().unwrap();
1191        r.write_str("b").unwrap();
1192        r.flush().unwrap();
1193        assert_eq!(r.into_inner().as_str(), "ab");
1194    }
1195
1196    #[test]
1197    fn combinators_stack_over_a_writer() {
1198        let inner = WriteOut::new(Vec::new()).with_budget(3);
1199        let mut j = Join::new(ReplaceText::new(inner, "-", "+"), ";");
1200        j.write_str("a-b").unwrap();
1201        j.write_str("c-").unwrap();
1202        j.flush().unwrap();
1203        let bytes = j.into_inner().into_inner().into_inner();
1204        assert_eq!(String::from_utf8(bytes).unwrap(), "a+b;c+");
1205    }
1206}