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