1use crate::{
4 error::EdifactError,
5 model::{Element, OwnedSegment, Segment, Span},
6 tokenizer::{Token, Tokenizer},
7};
8use memchr::memchr2;
9use smallvec::SmallVec;
10use std::borrow::Cow;
11use std::io::{BufRead, BufReader, Read};
12
13fn finish_element<'a>(
14 elements: &mut Vec<Element<'a>>,
15 current_components: &mut SmallVec<[(Cow<'a, str>, Span); 4]>,
16 current_element_start: &mut Option<usize>,
17) {
18 if let (Some(start), Some((_, last_span))) =
19 (current_element_start.take(), current_components.last())
20 {
21 let last_end = last_span.end;
22 elements.push(Element {
23 span: Span::new(start, last_end),
24 components: std::mem::take(current_components),
25 });
26 }
27}
28
29fn resolve_release(
30 val: &str,
31 release_char: char,
32 start_offset: usize,
33) -> Result<Cow<'_, str>, EdifactError> {
34 if !val.contains(release_char) {
35 return Ok(Cow::Borrowed(val));
36 }
37 resolve_release_owned(val, release_char, start_offset).map(Cow::Owned)
38}
39
40fn resolve_release_owned(
41 val: &str,
42 release_char: char,
43 start_offset: usize,
44) -> Result<String, EdifactError> {
45 let cap = val.len() - val.len() / 4;
51 let mut out = String::with_capacity(cap);
52 let mut chars = val.chars();
53 while let Some(c) = chars.next() {
54 if c == release_char {
55 if let Some(escaped) = chars.next() {
56 out.push(escaped);
57 } else {
58 return Err(EdifactError::InvalidReleaseSequence {
59 offset: start_offset + val.len().saturating_sub(1),
60 });
61 }
62 } else {
63 out.push(c);
64 }
65 }
66 Ok(out)
67}
68
69pub struct Parser<'a> {
74 tokenizer: Tokenizer<'a>,
75 peeked: Option<Token<'a>>,
77 release_char: char,
79}
80
81impl<'a> Parser<'a> {
82 pub fn new(tokenizer: Tokenizer<'a>) -> Self {
84 let release_char = tokenizer.service_string_advice().release_char as char;
85 Self {
86 tokenizer,
87 peeked: None,
88 release_char,
89 }
90 }
91}
92
93impl<'a> Iterator for Parser<'a> {
94 type Item = Result<Segment<'a>, EdifactError>;
95
96 fn next(&mut self) -> Option<Self::Item> {
97 let (tag, tag_span) = loop {
99 let tok = match self.peeked.take() {
100 Some(t) => Ok(t),
101 None => self.tokenizer.next()?,
102 };
103 match tok {
104 Ok(Token::SegmentTag { value, span }) => break (value, span),
105 Ok(Token::SegmentTerminator { .. }) => continue, Ok(Token::DataElement { span, .. }) | Ok(Token::ComponentElement { span, .. }) => {
107 return Some(Err(EdifactError::UnexpectedDataToken {
108 offset: span.start,
109 }));
110 }
111 Err(e) => return Some(Err(e)),
112 }
113 };
114
115 let mut elements: Vec<Element<'a>> = Vec::new();
120 let mut current_components: SmallVec<[(Cow<'a, str>, Span); 4]> = SmallVec::new();
121 let mut current_element_start: Option<usize> = None;
122 let mut in_element = false;
123 let mut segment_end = tag_span.end;
124
125 loop {
126 let tok = match self.tokenizer.next() {
127 Some(Ok(t)) => t,
128 Some(Err(e)) => return Some(Err(e)),
129 None => {
130 if in_element {
132 finish_element(
133 &mut elements,
134 &mut current_components,
135 &mut current_element_start,
136 );
137 if let Some(last) = elements.last() {
138 segment_end = last.span.end;
139 }
140 }
141 break;
142 }
143 };
144
145 match tok {
146 Token::SegmentTag {
147 value: next_tag,
148 span,
149 } => {
150 self.peeked = Some(Token::SegmentTag {
152 value: next_tag,
153 span,
154 });
155 if in_element {
156 finish_element(
157 &mut elements,
158 &mut current_components,
159 &mut current_element_start,
160 );
161 if let Some(last) = elements.last() {
162 segment_end = last.span.end;
163 }
164 }
165 break;
166 }
167 Token::SegmentTerminator { span } => {
168 if in_element {
169 finish_element(
170 &mut elements,
171 &mut current_components,
172 &mut current_element_start,
173 );
174 }
175 segment_end = span.end;
176 break;
177 }
178 Token::DataElement { value, span } => {
179 if in_element {
180 finish_element(
181 &mut elements,
182 &mut current_components,
183 &mut current_element_start,
184 );
185 }
186 let resolved = match resolve_release(value, self.release_char, span.start) {
187 Ok(v) => v,
188 Err(error) => return Some(Err(error)),
189 };
190 current_components.push((resolved, span));
191 current_element_start = Some(span.start);
192 in_element = true;
193 }
194 Token::ComponentElement { value, span } => {
195 if !in_element {
196 in_element = true;
198 current_element_start = Some(span.start);
199 }
200 let resolved = match resolve_release(value, self.release_char, span.start) {
201 Ok(v) => v,
202 Err(error) => return Some(Err(error)),
203 };
204 current_components.push((resolved, span));
205 }
206 }
207 }
208
209 Some(Ok(Segment {
210 tag,
211 span: Span::new(tag_span.start, segment_end),
212 tag_span,
213 elements,
214 }))
215 }
216}
217
218pub fn from_reader<R: Read>(reader: R) -> Result<Vec<OwnedSegment>, EdifactError> {
224 from_reader_stream(reader).collect()
225}
226
227pub fn from_bufread<R: BufRead>(reader: R) -> Result<Vec<OwnedSegment>, EdifactError> {
229 from_bufread_stream(reader).collect()
230}
231
232#[derive(Debug, Clone, Copy)]
248pub struct ReaderConfig {
249 pub max_segment_bytes: usize,
259 pub max_segments: Option<usize>,
268 pub max_input_bytes: Option<u64>,
280 pub max_messages: Option<usize>,
289}
290
291impl Default for ReaderConfig {
292 fn default() -> Self {
293 Self {
294 max_segment_bytes: 65_536,
295 max_segments: None,
296 max_input_bytes: None,
297 max_messages: None,
298 }
299 }
300}
301
302impl ReaderConfig {
303 #[must_use]
305 pub fn max_segment_bytes(mut self, limit: usize) -> Self {
306 self.max_segment_bytes = limit;
307 self
308 }
309
310 #[must_use]
312 pub fn max_segments(mut self, limit: usize) -> Self {
313 self.max_segments = Some(limit);
314 self
315 }
316
317 #[must_use]
319 pub fn max_input_bytes(mut self, limit: u64) -> Self {
320 self.max_input_bytes = Some(limit);
321 self
322 }
323
324 #[must_use]
326 pub fn max_messages(mut self, limit: usize) -> Self {
327 self.max_messages = Some(limit);
328 self
329 }
330}
331
332#[derive(Debug, Clone, Copy, PartialEq, Eq)]
334enum StreamState {
335 Init,
337 Running,
339 Done,
341}
342
343pub struct OwnedSegmentStream<R: BufRead> {
361 reader: R,
362 ssa: crate::tokenizer::ServiceStringAdvice,
363 state: StreamState,
364 stream_offset: u64,
365 config: ReaderConfig,
366 segments_yielded: usize,
368 messages_yielded: usize,
370 in_message: bool,
372 bytes_consumed: u64,
374}
375
376impl<R: BufRead> OwnedSegmentStream<R> {
377 fn new(reader: R) -> Self {
378 Self::with_config(reader, ReaderConfig::default())
379 }
380
381 fn with_config(reader: R, config: ReaderConfig) -> Self {
382 Self {
383 reader,
384 ssa: crate::tokenizer::ServiceStringAdvice::default(),
385 state: StreamState::Init,
386 stream_offset: 0,
387 config,
388 segments_yielded: 0,
389 messages_yielded: 0,
390 in_message: false,
391 bytes_consumed: 0,
392 }
393 }
394}
395
396enum FastSegment {
400 Parsed(OwnedSegment, usize),
402 Skip(usize),
404 NeedMore,
406 Eof,
408 Err(EdifactError),
410}
411
412fn find_unescaped_term(buf: &[u8], term: u8, release: u8) -> Option<usize> {
424 let mut i = 0;
425 while i < buf.len() {
426 let rel = memchr2(release, term, &buf[i..])?;
428 let pos = i + rel;
429 if buf[pos] == release {
430 i = pos + 2;
432 } else {
433 return Some(pos);
435 }
436 }
437 None
438}
439
440fn try_fast_segment<R: BufRead>(
445 reader: &mut R,
446 ssa: crate::tokenizer::ServiceStringAdvice,
447 seg_start: usize,
448 max_segment_bytes: usize,
449) -> FastSegment {
450 let buf = match reader.fill_buf() {
451 Ok(b) => b,
452 Err(e) => return FastSegment::Err(e.into()),
453 };
454
455 if buf.is_empty() {
456 return FastSegment::Eof;
457 }
458
459 let Some(pos) = find_unescaped_term(buf, ssa.segment_term, ssa.release_char) else {
460 return FastSegment::NeedMore;
461 };
462
463 if pos > max_segment_bytes {
466 return FastSegment::Err(EdifactError::SegmentTooLong {
467 offset: seg_start,
468 limit: max_segment_bytes,
469 });
470 }
471
472 let seg_bytes = &buf[..pos];
474
475 if seg_bytes
477 .iter()
478 .all(|&b| matches!(b, b' ' | b'\t' | b'\r' | b'\n'))
479 {
480 return FastSegment::Skip(pos + 1);
481 }
482
483 let tok = Tokenizer::with_limit(&buf[..pos + 1], ssa, max_segment_bytes);
490 let mut parser_iter = Parser::new(tok);
491 match parser_iter.next() {
492 None => FastSegment::Skip(pos + 1),
493 Some(Err(e)) => FastSegment::Err(e),
494 Some(Ok(s)) => FastSegment::Parsed(OwnedSegment::from(s).offset(seg_start), pos + 1),
495 }
496 }
498
499impl<R: BufRead> Iterator for OwnedSegmentStream<R> {
502 type Item = Result<OwnedSegment, EdifactError>;
503
504 fn next(&mut self) -> Option<Self::Item> {
505 if self.state == StreamState::Done {
506 return None;
507 }
508
509 if let Some(max) = self.config.max_segments {
511 if self.segments_yielded >= max {
512 self.state = StreamState::Done;
513 return None;
514 }
515 }
516
517 if let Some(max) = self.config.max_input_bytes {
519 if self.bytes_consumed >= max {
520 self.state = StreamState::Done;
521 return None;
522 }
523 }
524
525 if let Some(max) = self.config.max_messages {
527 if self.messages_yielded >= max {
528 self.state = StreamState::Done;
529 return None;
530 }
531 }
532
533 loop {
534 if self.state == StreamState::Running {
536 let seg_start = self.stream_offset;
537 match try_fast_segment(
538 &mut self.reader,
539 self.ssa,
540 seg_start.min(usize::MAX as u64) as usize,
544 self.config.max_segment_bytes,
545 ) {
546 FastSegment::Parsed(seg, n) => {
547 let n = n as u64;
548 self.reader.consume(n as usize);
549 self.stream_offset += n;
550 self.bytes_consumed = self.stream_offset;
551 self.segments_yielded += 1;
552 if seg.tag == "UNT" {
557 if self.in_message {
558 self.messages_yielded += 1;
559 }
560 self.in_message = false;
561 } else if seg.tag == "UNH" {
562 self.in_message = true;
563 }
564 if let Some(max) = self.config.max_input_bytes {
568 if self.bytes_consumed >= max {
569 self.state = StreamState::Done;
570 }
571 }
572 return Some(Ok(seg));
573 }
574 FastSegment::Skip(n) => {
575 let n = n as u64;
576 self.reader.consume(n as usize);
577 self.stream_offset += n;
578 self.bytes_consumed = self.stream_offset;
579 continue;
580 }
581 FastSegment::Eof => return None,
582 FastSegment::Err(e) => {
583 self.state = StreamState::Done;
584 return Some(Err(e));
585 }
586 FastSegment::NeedMore => {
587 }
589 }
590 }
591
592 let mut scanned = self.state != StreamState::Init;
594 let mut slow_offset: usize = self.stream_offset.min(usize::MAX as u64) as usize;
599 let mut raw = match read_next_raw_segment(
600 &mut self.reader,
601 &mut self.ssa,
602 &mut scanned,
603 &mut slow_offset,
604 self.config.max_segment_bytes,
605 ) {
606 Ok(Some(r)) => r,
607 Ok(None) => return None,
608 Err(e) => {
609 self.state = StreamState::Done;
610 return Some(Err(e));
611 }
612 };
613 self.stream_offset = slow_offset as u64;
614 if scanned {
615 self.state = StreamState::Running;
616 }
617 self.bytes_consumed = self.stream_offset;
618
619 raw.bytes.push(self.ssa.segment_term);
620 let tok = Tokenizer::with_limit(
624 raw.bytes.as_slice(),
625 self.ssa,
626 self.config.max_segment_bytes,
627 );
628 let mut parser_iter = Parser::new(tok);
629 match parser_iter.next() {
630 Some(Ok(s)) => {
631 self.segments_yielded += 1;
632 let seg = OwnedSegment::from(s).offset(raw.start_offset);
633 if seg.tag == "UNT" {
634 if self.in_message {
635 self.messages_yielded += 1;
636 }
637 self.in_message = false;
638 } else if seg.tag == "UNH" {
639 self.in_message = true;
640 }
641 return Some(Ok(seg));
642 }
643 Some(Err(e)) => {
644 self.state = StreamState::Done;
645 return Some(Err(e));
646 }
647 None => {} }
649 }
650 }
651}
652
653pub fn from_bufread_stream<R: BufRead>(reader: R) -> OwnedSegmentStream<R> {
655 OwnedSegmentStream::new(reader)
656}
657
658pub fn from_bufread_stream_with_config<R: BufRead>(
660 reader: R,
661 config: ReaderConfig,
662) -> OwnedSegmentStream<R> {
663 OwnedSegmentStream::with_config(reader, config)
664}
665
666pub fn from_reader_stream<R: Read>(reader: R) -> OwnedSegmentStream<BufReader<R>> {
668 from_bufread_stream(BufReader::new(reader))
669}
670
671pub fn from_reader_with_config<R: Read>(
684 reader: R,
685 config: ReaderConfig,
686) -> OwnedSegmentStream<BufReader<R>> {
687 from_bufread_stream_with_config(BufReader::new(reader), config)
688}
689
690fn read_next_raw_segment<R: BufRead>(
691 reader: &mut R,
692 ssa: &mut crate::tokenizer::ServiceStringAdvice,
693 scanned_header: &mut bool,
694 stream_offset: &mut usize,
695 max_segment_bytes: usize,
696) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
697 loop {
698 let Some((first_offset, first)) = read_next_non_ws_byte(reader, stream_offset)? else {
699 return Ok(None);
700 };
701
702 if !*scanned_header && first == b'U' {
703 let second = read_required_byte(reader, stream_offset)?;
704 let third = read_required_byte(reader, stream_offset)?;
705 if second == b'N' && third == b'A' {
706 let mut una = [0u8; 9];
707 una[0] = b'U';
708 una[1] = b'N';
709 una[2] = b'A';
710 for slot in una.iter_mut().skip(3) {
711 *slot = read_required_byte(reader, stream_offset)?;
712 }
713 *ssa = crate::tokenizer::ServiceStringAdvice {
714 component_sep: una[3],
715 element_sep: una[4],
716 decimal_mark: una[5],
717 release_char: una[6],
718 repetition_sep: una[7],
719 segment_term: una[8],
720 };
721 if !ssa.is_valid() {
722 return Err(EdifactError::InvalidUna);
723 }
724 *scanned_header = true;
725 continue;
726 }
727
728 *scanned_header = true;
729 return read_remainder_of_segment(
730 reader,
731 ssa,
732 crate::tokenizer::RawSegment {
733 bytes: vec![first, second, third],
734 start_offset: first_offset,
735 },
736 stream_offset,
737 max_segment_bytes,
738 );
739 }
740
741 *scanned_header = true;
742 return read_remainder_of_segment(
743 reader,
744 ssa,
745 crate::tokenizer::RawSegment {
746 bytes: vec![first],
747 start_offset: first_offset,
748 },
749 stream_offset,
750 max_segment_bytes,
751 );
752 }
753}
754
755fn read_remainder_of_segment<R: BufRead>(
756 reader: &mut R,
757 ssa: &crate::tokenizer::ServiceStringAdvice,
758 mut out: crate::tokenizer::RawSegment,
759 stream_offset: &mut usize,
760 max_segment_bytes: usize,
761) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
762 let mut escaped = false;
763 loop {
764 if out.bytes.len() > max_segment_bytes {
770 return Err(EdifactError::SegmentTooLong {
771 offset: out.start_offset,
772 limit: max_segment_bytes,
773 });
774 }
775 let Some(byte) = read_next_byte(reader, stream_offset)? else {
776 return if out.bytes.is_empty() {
777 Ok(None)
778 } else if escaped {
779 Err(EdifactError::InvalidReleaseSequence {
780 offset: out.start_offset + out.bytes.len().saturating_sub(1),
781 })
782 } else {
783 Err(EdifactError::UnexpectedEof {
784 offset: out.start_offset + out.bytes.len(),
785 })
786 };
787 };
788
789 if !escaped && byte == ssa.segment_term {
790 return Ok(Some(out));
791 }
792
793 if !escaped && byte == ssa.release_char {
794 escaped = true;
795 out.bytes.push(byte);
796 continue;
797 }
798
799 escaped = false;
800 out.bytes.push(byte);
801 }
802}
803
804fn read_next_byte<R: BufRead>(
805 reader: &mut R,
806 stream_offset: &mut usize,
807) -> Result<Option<u8>, EdifactError> {
808 let buf = reader.fill_buf()?;
809 if buf.is_empty() {
810 return Ok(None);
811 }
812
813 let byte = buf[0];
814 reader.consume(1);
815 let next_offset = stream_offset.saturating_add(1);
820 *stream_offset = next_offset;
821 Ok(Some(byte))
822}
823
824fn read_required_byte<R: BufRead>(
825 reader: &mut R,
826 stream_offset: &mut usize,
827) -> Result<u8, EdifactError> {
828 read_next_byte(reader, stream_offset)?.ok_or(EdifactError::UnexpectedEof {
829 offset: *stream_offset,
830 })
831}
832
833fn read_next_non_ws_byte<R: BufRead>(
834 reader: &mut R,
835 stream_offset: &mut usize,
836) -> Result<Option<(usize, u8)>, EdifactError> {
837 loop {
838 let current_offset = *stream_offset;
839 let Some(byte) = read_next_byte(reader, stream_offset)? else {
840 return Ok(None);
841 };
842 if !matches!(byte, b' ' | b'\t' | b'\r' | b'\n') {
843 return Ok(Some((current_offset, byte)));
844 }
845 }
846}
847
848#[cfg(test)]
849mod tests {
850 use super::*;
851 use crate::tokenizer::ServiceStringAdvice;
852
853 fn parse_all(input: &[u8]) -> Vec<Segment<'_>> {
854 let ssa = ServiceStringAdvice::from_bytes_unchecked(input);
855 let tok = Tokenizer::new(input, ssa);
856 Parser::new(tok)
857 .collect::<Result<Vec<_>, _>>()
858 .expect("parse failed")
859 }
860
861 #[test]
862 fn parses_unb_unz() {
863 let input = b"UNB+UNOA:1+SENDER+RECEIVER+200101:0900+1'UNZ+0+1'";
864 let segs = parse_all(input);
865 assert_eq!(segs.len(), 2);
866 assert_eq!(segs[0].tag, "UNB");
867 assert_eq!(segs[1].tag, "UNZ");
868 assert_eq!(segs[0].tag_span, Span::new(0, 3));
869 assert_eq!(segs[0].span, Span::new(0, 41));
870 }
871
872 #[test]
873 fn element_access() {
874 let input = b"BGM+220+ORDER123+9'";
875 let segs = parse_all(input);
876 assert_eq!(segs[0].element_str(0), Some("220"));
877 assert_eq!(segs[0].element_str(1), Some("ORDER123"));
878 }
879
880 #[test]
881 fn component_access() {
882 let input = b"DTM+137:20200101:102'";
883 let segs = parse_all(input);
884 let dtm = &segs[0];
885 assert_eq!(dtm.get_element(0).unwrap().get_component(0), Some("137"));
886 assert_eq!(
887 dtm.get_element(0).unwrap().get_component(1),
888 Some("20200101")
889 );
890 assert_eq!(dtm.get_element(0).unwrap().get_component(2), Some("102"));
891 }
892
893 #[test]
894 fn release_char_resolved() {
895 let input = b"FTX+AAA++test?+value'";
896 let segs = parse_all(input);
897 assert_eq!(segs[0].element_str(2), Some("test+value"));
898 assert_eq!(
899 segs[0].get_element(2).unwrap().component_span(0),
900 Some(Span::new(9, 20))
901 );
902 }
903
904 #[test]
905 fn reader_path_preserves_custom_una_delimiters() {
906 let input = b"UNA:;.? 'BGM;220;test?;value'";
907 let segments = super::from_bufread(std::io::BufReader::new(std::io::Cursor::new(input)))
908 .expect("reader parse should succeed");
909 let bgm = segments
910 .iter()
911 .find(|segment| segment.tag == "BGM")
912 .expect("BGM segment should be present");
913 assert_eq!(bgm.elements[0].components[0].0, "220");
914 assert_eq!(bgm.elements[1].components[0].0, "test;value");
915 }
916
917 #[test]
918 fn arbitrary_bytes_no_panic() {
919 let garbage: &[u8] = b"\xff\x00\x01\x02ABC+++'''???";
921 let _ = crate::from_bytes(garbage).collect::<Vec<_>>();
922 }
923
924 #[test]
925 fn from_reader_handles_chunk_boundaries() {
926 let input = b"UNA:+.? 'BGM+220+test?+value'UNT+2+1'";
927 let reader = std::io::BufReader::with_capacity(5, std::io::Cursor::new(input));
928 let parsed = from_bufread(reader).expect("reader parsing should succeed");
929 assert_eq!(parsed.len(), 2);
930 assert_eq!(parsed[0].tag, "BGM");
931 assert_eq!(parsed[0].elements[1].components[0].0, "test+value");
932 assert_eq!(parsed[1].tag, "UNT");
933 }
934
935 #[test]
936 fn from_reader_without_una_uses_default_delimiters() {
937 let input = b"BGM+220+X'UNT+2+1'";
938 let parsed =
939 from_reader(std::io::Cursor::new(input)).expect("reader parsing should succeed");
940 assert_eq!(parsed.len(), 2);
941 assert_eq!(parsed[0].tag, "BGM");
942 assert_eq!(parsed[0].elements[0].components[0].0, "220");
943 assert_eq!(parsed[1].span, Span::new(10, 18));
944 }
945
946 #[test]
947 fn dangling_release_sequence_is_error() {
948 let input = b"FTX+AAA++dangling?";
949 let err = crate::from_bytes(input)
950 .collect::<Result<Vec<_>, _>>()
951 .expect_err("expected dangling release to fail");
952
953 assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
954 }
955
956 #[test]
957 fn from_reader_reports_dangling_release_sequence() {
958 let input = b"FTX+AAA++dangling?";
959 let err = from_reader(std::io::Cursor::new(input))
960 .expect_err("expected dangling release from reader path");
961 assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
962 }
963
964 #[test]
965 fn from_reader_rejects_invalid_una() {
966 let input = b"UNA::.? 'BGM:220'";
967 let err = from_reader(std::io::Cursor::new(input))
968 .expect_err("invalid UNA should fail reader parsing");
969 assert!(matches!(err, EdifactError::InvalidUna));
970 }
971}