1use std::io;
16use std::sync::Arc;
17
18use tabnas_alchemy::shared::{Fail, JoinOut, Limits, Metrics};
19
20pub use tabnas_alchemy::shared::text::TextOut;
23
24pub const DEFAULT_BUDGET: usize = 32 * 1024;
28
29pub 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 accepted: u64,
50 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 pub fn with_budget(mut self, budget: usize) -> Self {
73 self.budget = budget;
74 self
75 }
76
77 pub fn with_limits(mut self, limits: &Limits) -> Self {
80 self.max_output_bytes = limits.max_output_bytes;
81 self
82 }
83
84 pub fn with_metrics(mut self, metrics: Arc<Metrics>) -> Self {
86 self.metrics = Some(metrics);
87 self
88 }
89
90 pub fn accepted(&self) -> u64 {
92 self.accepted
93 }
94
95 pub fn committed(&self) -> u64 {
99 self.committed
100 }
101
102 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 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 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 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#[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 fn has_committed(&self) -> bool {
245 !self.0.is_empty()
246 }
247}
248
249pub 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 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 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 pub fn items(&self) -> u64 {
304 self.items
305 }
306
307 pub fn into_inner(self) -> O {
308 self.out
309 }
310}
311
312impl<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 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
344pub 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 pub fn item_start(&mut self) -> Result<(), Fail> {
361 self.0.item_start()
362 }
363
364 pub fn item_end(&mut self) -> Result<(), Fail> {
366 self.0.item_end()
367 }
368
369 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
393pub 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 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 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 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 #[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 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 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 assert!(err.committed_output);
589 assert_eq!(out.accepted(), 9);
590 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 #[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 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 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 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 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 #[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 #[cfg(unix)]
812 const SHORT_WRITE_CHILD: &str = "TABNAS_RENDER_SHORT_WRITE_CHILD";
813
814 #[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 #[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 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 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 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}