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::with_capacity(8);
116 let mut current_components: SmallVec<[(Cow<'a, str>, Span); 4]> = SmallVec::new();
117 let mut current_element_start: Option<usize> = None;
118 let mut in_element = false;
119 let mut segment_end = tag_span.end;
120
121 loop {
122 let tok = match self.tokenizer.next() {
123 Some(Ok(t)) => t,
124 Some(Err(e)) => return Some(Err(e)),
125 None => {
126 if in_element {
128 finish_element(
129 &mut elements,
130 &mut current_components,
131 &mut current_element_start,
132 );
133 if let Some(last) = elements.last() {
134 segment_end = last.span.end;
135 }
136 }
137 break;
138 }
139 };
140
141 match tok {
142 Token::SegmentTag {
143 value: next_tag,
144 span,
145 } => {
146 self.peeked = Some(Token::SegmentTag {
148 value: next_tag,
149 span,
150 });
151 if in_element {
152 finish_element(
153 &mut elements,
154 &mut current_components,
155 &mut current_element_start,
156 );
157 if let Some(last) = elements.last() {
158 segment_end = last.span.end;
159 }
160 }
161 break;
162 }
163 Token::SegmentTerminator { span } => {
164 if in_element {
165 finish_element(
166 &mut elements,
167 &mut current_components,
168 &mut current_element_start,
169 );
170 }
171 segment_end = span.end;
172 break;
173 }
174 Token::DataElement { value, span } => {
175 if in_element {
176 finish_element(
177 &mut elements,
178 &mut current_components,
179 &mut current_element_start,
180 );
181 }
182 let resolved = match resolve_release(value, self.release_char, span.start) {
183 Ok(v) => v,
184 Err(error) => return Some(Err(error)),
185 };
186 current_components.push((resolved, span));
187 current_element_start = Some(span.start);
188 in_element = true;
189 }
190 Token::ComponentElement { value, span } => {
191 if !in_element {
192 in_element = true;
194 current_element_start = Some(span.start);
195 }
196 let resolved = match resolve_release(value, self.release_char, span.start) {
197 Ok(v) => v,
198 Err(error) => return Some(Err(error)),
199 };
200 current_components.push((resolved, span));
201 }
202 }
203 }
204
205 Some(Ok(Segment {
206 tag,
207 span: Span::new(tag_span.start, segment_end),
208 tag_span,
209 elements,
210 }))
211 }
212}
213
214pub fn from_reader<R: Read>(reader: R) -> Result<Vec<OwnedSegment>, EdifactError> {
220 from_reader_stream(reader).collect()
221}
222
223pub fn from_bufread<R: BufRead>(reader: R) -> Result<Vec<OwnedSegment>, EdifactError> {
225 from_bufread_stream(reader).collect()
226}
227
228#[derive(Debug, Clone, Copy)]
244pub struct ReaderConfig {
245 pub max_segment_bytes: usize,
255 pub max_segments: Option<usize>,
264 pub max_input_bytes: Option<u64>,
276 pub max_messages: Option<usize>,
285}
286
287impl Default for ReaderConfig {
288 fn default() -> Self {
289 Self {
290 max_segment_bytes: 65_536,
291 max_segments: None,
292 max_input_bytes: None,
293 max_messages: None,
294 }
295 }
296}
297
298impl ReaderConfig {
299 #[must_use]
301 pub fn max_segment_bytes(mut self, limit: usize) -> Self {
302 self.max_segment_bytes = limit;
303 self
304 }
305
306 #[must_use]
308 pub fn max_segments(mut self, limit: usize) -> Self {
309 self.max_segments = Some(limit);
310 self
311 }
312
313 #[must_use]
315 pub fn max_input_bytes(mut self, limit: u64) -> Self {
316 self.max_input_bytes = Some(limit);
317 self
318 }
319
320 #[must_use]
322 pub fn max_messages(mut self, limit: usize) -> Self {
323 self.max_messages = Some(limit);
324 self
325 }
326}
327
328#[derive(Debug, Clone, Copy, PartialEq, Eq)]
330enum StreamState {
331 Init,
333 Running,
335 Done,
337}
338
339pub struct OwnedSegmentStream<R: BufRead> {
357 reader: R,
358 ssa: crate::tokenizer::ServiceStringAdvice,
359 state: StreamState,
360 stream_offset: u64,
361 config: ReaderConfig,
362 segments_yielded: usize,
364 messages_yielded: usize,
366 in_message: bool,
368 bytes_consumed: u64,
370}
371
372impl<R: BufRead> OwnedSegmentStream<R> {
373 fn new(reader: R) -> Self {
374 Self::with_config(reader, ReaderConfig::default())
375 }
376
377 fn with_config(reader: R, config: ReaderConfig) -> Self {
378 Self {
379 reader,
380 ssa: crate::tokenizer::ServiceStringAdvice::default(),
381 state: StreamState::Init,
382 stream_offset: 0,
383 config,
384 segments_yielded: 0,
385 messages_yielded: 0,
386 in_message: false,
387 bytes_consumed: 0,
388 }
389 }
390}
391
392enum FastSegment {
396 Parsed(OwnedSegment, usize),
398 Skip(usize),
400 NeedMore,
402 Eof,
404 Err(EdifactError),
406}
407
408fn find_unescaped_term(buf: &[u8], term: u8, release: u8) -> Option<usize> {
420 let mut i = 0;
421 while i < buf.len() {
422 let rel = memchr2(release, term, &buf[i..])?;
424 let pos = i + rel;
425 if buf[pos] == release {
426 i = pos + 2;
428 } else {
429 return Some(pos);
431 }
432 }
433 None
434}
435
436fn try_fast_segment<R: BufRead>(
441 reader: &mut R,
442 ssa: crate::tokenizer::ServiceStringAdvice,
443 seg_start: usize,
444 max_segment_bytes: usize,
445) -> FastSegment {
446 let buf = match reader.fill_buf() {
447 Ok(b) => b,
448 Err(e) => return FastSegment::Err(e.into()),
449 };
450
451 if buf.is_empty() {
452 return FastSegment::Eof;
453 }
454
455 let Some(pos) = find_unescaped_term(buf, ssa.segment_term, ssa.release_char) else {
456 return FastSegment::NeedMore;
457 };
458
459 if pos > max_segment_bytes {
462 return FastSegment::Err(EdifactError::SegmentTooLong {
463 offset: seg_start,
464 limit: max_segment_bytes,
465 });
466 }
467
468 let seg_bytes = &buf[..pos];
470
471 if seg_bytes
473 .iter()
474 .all(|&b| matches!(b, b' ' | b'\t' | b'\r' | b'\n'))
475 {
476 return FastSegment::Skip(pos + 1);
477 }
478
479 let tok = Tokenizer::with_limit(&buf[..pos + 1], ssa, max_segment_bytes);
486 let mut parser_iter = Parser::new(tok);
487 match parser_iter.next() {
488 None => FastSegment::Skip(pos + 1),
489 Some(Err(e)) => FastSegment::Err(e),
490 Some(Ok(s)) => FastSegment::Parsed(OwnedSegment::from(s).offset(seg_start), pos + 1),
491 }
492 }
494
495impl<R: BufRead> Iterator for OwnedSegmentStream<R> {
498 type Item = Result<OwnedSegment, EdifactError>;
499
500 fn next(&mut self) -> Option<Self::Item> {
501 if self.state == StreamState::Done {
502 return None;
503 }
504
505 if let Some(max) = self.config.max_segments {
507 if self.segments_yielded >= max {
508 self.state = StreamState::Done;
509 return None;
510 }
511 }
512
513 if let Some(max) = self.config.max_input_bytes {
515 if self.bytes_consumed >= max {
516 self.state = StreamState::Done;
517 return None;
518 }
519 }
520
521 if let Some(max) = self.config.max_messages {
523 if self.messages_yielded >= max {
524 self.state = StreamState::Done;
525 return None;
526 }
527 }
528
529 loop {
530 if self.state == StreamState::Running {
532 let seg_start = self.stream_offset;
533 match try_fast_segment(
534 &mut self.reader,
535 self.ssa,
536 seg_start.min(usize::MAX as u64) as usize,
540 self.config.max_segment_bytes,
541 ) {
542 FastSegment::Parsed(seg, n) => {
543 let n = n as u64;
544 self.reader.consume(n as usize);
545 self.stream_offset += n;
546 self.bytes_consumed = self.stream_offset;
547 self.segments_yielded += 1;
548 if seg.tag == "UNT" {
553 if self.in_message {
554 self.messages_yielded += 1;
555 }
556 self.in_message = false;
557 } else if seg.tag == "UNH" {
558 self.in_message = true;
559 }
560 if let Some(max) = self.config.max_input_bytes {
564 if self.bytes_consumed >= max {
565 self.state = StreamState::Done;
566 }
567 }
568 return Some(Ok(seg));
569 }
570 FastSegment::Skip(n) => {
571 let n = n as u64;
572 self.reader.consume(n as usize);
573 self.stream_offset += n;
574 self.bytes_consumed = self.stream_offset;
575 continue;
576 }
577 FastSegment::Eof => return None,
578 FastSegment::Err(e) => {
579 self.state = StreamState::Done;
580 return Some(Err(e));
581 }
582 FastSegment::NeedMore => {
583 }
585 }
586 }
587
588 let mut scanned = self.state != StreamState::Init;
590 let mut slow_offset: usize = self.stream_offset.min(usize::MAX as u64) as usize;
595 let mut raw = match read_next_raw_segment(
596 &mut self.reader,
597 &mut self.ssa,
598 &mut scanned,
599 &mut slow_offset,
600 self.config.max_segment_bytes,
601 ) {
602 Ok(Some(r)) => r,
603 Ok(None) => return None,
604 Err(e) => {
605 self.state = StreamState::Done;
606 return Some(Err(e));
607 }
608 };
609 self.stream_offset = slow_offset as u64;
610 if scanned {
611 self.state = StreamState::Running;
612 }
613 self.bytes_consumed = self.stream_offset;
614
615 raw.bytes.push(self.ssa.segment_term);
616 let tok = Tokenizer::with_limit(
620 raw.bytes.as_slice(),
621 self.ssa,
622 self.config.max_segment_bytes,
623 );
624 let mut parser_iter = Parser::new(tok);
625 match parser_iter.next() {
626 Some(Ok(s)) => {
627 self.segments_yielded += 1;
628 let seg = OwnedSegment::from(s).offset(raw.start_offset);
629 if seg.tag == "UNT" {
630 if self.in_message {
631 self.messages_yielded += 1;
632 }
633 self.in_message = false;
634 } else if seg.tag == "UNH" {
635 self.in_message = true;
636 }
637 return Some(Ok(seg));
638 }
639 Some(Err(e)) => {
640 self.state = StreamState::Done;
641 return Some(Err(e));
642 }
643 None => {} }
645 }
646 }
647}
648
649pub fn from_bufread_stream<R: BufRead>(reader: R) -> OwnedSegmentStream<R> {
651 OwnedSegmentStream::new(reader)
652}
653
654pub fn from_bufread_stream_with_config<R: BufRead>(
656 reader: R,
657 config: ReaderConfig,
658) -> OwnedSegmentStream<R> {
659 OwnedSegmentStream::with_config(reader, config)
660}
661
662pub fn from_reader_stream<R: Read>(reader: R) -> OwnedSegmentStream<BufReader<R>> {
664 from_bufread_stream(BufReader::new(reader))
665}
666
667pub fn from_reader_with_config<R: Read>(
680 reader: R,
681 config: ReaderConfig,
682) -> OwnedSegmentStream<BufReader<R>> {
683 from_bufread_stream_with_config(BufReader::new(reader), config)
684}
685
686fn read_next_raw_segment<R: BufRead>(
687 reader: &mut R,
688 ssa: &mut crate::tokenizer::ServiceStringAdvice,
689 scanned_header: &mut bool,
690 stream_offset: &mut usize,
691 max_segment_bytes: usize,
692) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
693 loop {
694 let Some((first_offset, first)) = read_next_non_ws_byte(reader, stream_offset)? else {
695 return Ok(None);
696 };
697
698 if !*scanned_header && first == b'U' {
699 let second = read_required_byte(reader, stream_offset)?;
700 let third = read_required_byte(reader, stream_offset)?;
701 if second == b'N' && third == b'A' {
702 let mut una = [0u8; 9];
703 una[0] = b'U';
704 una[1] = b'N';
705 una[2] = b'A';
706 for slot in una.iter_mut().skip(3) {
707 *slot = read_required_byte(reader, stream_offset)?;
708 }
709 *ssa = crate::tokenizer::ServiceStringAdvice {
710 component_sep: una[3],
711 element_sep: una[4],
712 decimal_mark: una[5],
713 release_char: una[6],
714 repetition_sep: una[7],
715 segment_term: una[8],
716 };
717 if !ssa.is_valid() {
718 return Err(EdifactError::InvalidUna);
719 }
720 *scanned_header = true;
721 continue;
722 }
723
724 *scanned_header = true;
725 return read_remainder_of_segment(
726 reader,
727 ssa,
728 crate::tokenizer::RawSegment {
729 bytes: vec![first, second, third],
730 start_offset: first_offset,
731 },
732 stream_offset,
733 max_segment_bytes,
734 );
735 }
736
737 *scanned_header = true;
738 return read_remainder_of_segment(
739 reader,
740 ssa,
741 crate::tokenizer::RawSegment {
742 bytes: vec![first],
743 start_offset: first_offset,
744 },
745 stream_offset,
746 max_segment_bytes,
747 );
748 }
749}
750
751fn read_remainder_of_segment<R: BufRead>(
752 reader: &mut R,
753 ssa: &crate::tokenizer::ServiceStringAdvice,
754 mut out: crate::tokenizer::RawSegment,
755 stream_offset: &mut usize,
756 max_segment_bytes: usize,
757) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
758 let mut escaped = false;
759 loop {
760 if out.bytes.len() >= max_segment_bytes {
761 return Err(EdifactError::SegmentTooLong {
762 offset: out.start_offset,
763 limit: max_segment_bytes,
764 });
765 }
766 let Some(byte) = read_next_byte(reader, stream_offset)? else {
767 return if out.bytes.is_empty() {
768 Ok(None)
769 } else if escaped {
770 Err(EdifactError::InvalidReleaseSequence {
771 offset: out.start_offset + out.bytes.len().saturating_sub(1),
772 })
773 } else {
774 Err(EdifactError::UnexpectedEof {
775 offset: out.start_offset + out.bytes.len(),
776 })
777 };
778 };
779
780 if !escaped && byte == ssa.segment_term {
781 return Ok(Some(out));
782 }
783
784 if !escaped && byte == ssa.release_char {
785 escaped = true;
786 out.bytes.push(byte);
787 continue;
788 }
789
790 escaped = false;
791 out.bytes.push(byte);
792 }
793}
794
795fn read_next_byte<R: BufRead>(
796 reader: &mut R,
797 stream_offset: &mut usize,
798) -> Result<Option<u8>, EdifactError> {
799 let buf = reader.fill_buf()?;
800 if buf.is_empty() {
801 return Ok(None);
802 }
803
804 let byte = buf[0];
805 reader.consume(1);
806 let next_offset = stream_offset.saturating_add(1);
811 *stream_offset = next_offset;
812 Ok(Some(byte))
813}
814
815fn read_required_byte<R: BufRead>(
816 reader: &mut R,
817 stream_offset: &mut usize,
818) -> Result<u8, EdifactError> {
819 read_next_byte(reader, stream_offset)?.ok_or(EdifactError::UnexpectedEof {
820 offset: *stream_offset,
821 })
822}
823
824fn read_next_non_ws_byte<R: BufRead>(
825 reader: &mut R,
826 stream_offset: &mut usize,
827) -> Result<Option<(usize, u8)>, EdifactError> {
828 loop {
829 let current_offset = *stream_offset;
830 let Some(byte) = read_next_byte(reader, stream_offset)? else {
831 return Ok(None);
832 };
833 if !matches!(byte, b' ' | b'\t' | b'\r' | b'\n') {
834 return Ok(Some((current_offset, byte)));
835 }
836 }
837}
838
839#[cfg(test)]
840mod tests {
841 use super::*;
842 use crate::tokenizer::ServiceStringAdvice;
843
844 fn parse_all(input: &[u8]) -> Vec<Segment<'_>> {
845 let ssa = ServiceStringAdvice::from_bytes_unchecked(input);
846 let tok = Tokenizer::new(input, ssa);
847 Parser::new(tok)
848 .collect::<Result<Vec<_>, _>>()
849 .expect("parse failed")
850 }
851
852 #[test]
853 fn parses_unb_unz() {
854 let input = b"UNB+UNOA:1+SENDER+RECEIVER+200101:0900+1'UNZ+0+1'";
855 let segs = parse_all(input);
856 assert_eq!(segs.len(), 2);
857 assert_eq!(segs[0].tag, "UNB");
858 assert_eq!(segs[1].tag, "UNZ");
859 assert_eq!(segs[0].tag_span, Span::new(0, 3));
860 assert_eq!(segs[0].span, Span::new(0, 41));
861 }
862
863 #[test]
864 fn element_access() {
865 let input = b"BGM+220+ORDER123+9'";
866 let segs = parse_all(input);
867 assert_eq!(segs[0].element_str(0), Some("220"));
868 assert_eq!(segs[0].element_str(1), Some("ORDER123"));
869 }
870
871 #[test]
872 fn component_access() {
873 let input = b"DTM+137:20200101:102'";
874 let segs = parse_all(input);
875 let dtm = &segs[0];
876 assert_eq!(dtm.get_element(0).unwrap().get_component(0), Some("137"));
877 assert_eq!(
878 dtm.get_element(0).unwrap().get_component(1),
879 Some("20200101")
880 );
881 assert_eq!(dtm.get_element(0).unwrap().get_component(2), Some("102"));
882 }
883
884 #[test]
885 fn release_char_resolved() {
886 let input = b"FTX+AAA++test?+value'";
887 let segs = parse_all(input);
888 assert_eq!(segs[0].element_str(2), Some("test+value"));
889 assert_eq!(
890 segs[0].get_element(2).unwrap().component_span(0),
891 Some(Span::new(9, 20))
892 );
893 }
894
895 #[test]
896 fn reader_path_preserves_custom_una_delimiters() {
897 let input = b"UNA:;.? 'BGM;220;test?;value'";
898 let segments = super::from_bufread(std::io::BufReader::new(std::io::Cursor::new(input)))
899 .expect("reader parse should succeed");
900 let bgm = segments
901 .iter()
902 .find(|segment| segment.tag == "BGM")
903 .expect("BGM segment should be present");
904 assert_eq!(bgm.elements[0].components[0].0, "220");
905 assert_eq!(bgm.elements[1].components[0].0, "test;value");
906 }
907
908 #[test]
909 fn arbitrary_bytes_no_panic() {
910 let garbage: &[u8] = b"\xff\x00\x01\x02ABC+++'''???";
912 let _ = crate::from_bytes(garbage).collect::<Vec<_>>();
913 }
914
915 #[test]
916 fn from_reader_handles_chunk_boundaries() {
917 let input = b"UNA:+.? 'BGM+220+test?+value'UNT+2+1'";
918 let reader = std::io::BufReader::with_capacity(5, std::io::Cursor::new(input));
919 let parsed = from_bufread(reader).expect("reader parsing should succeed");
920 assert_eq!(parsed.len(), 2);
921 assert_eq!(parsed[0].tag, "BGM");
922 assert_eq!(parsed[0].elements[1].components[0].0, "test+value");
923 assert_eq!(parsed[1].tag, "UNT");
924 }
925
926 #[test]
927 fn from_reader_without_una_uses_default_delimiters() {
928 let input = b"BGM+220+X'UNT+2+1'";
929 let parsed =
930 from_reader(std::io::Cursor::new(input)).expect("reader parsing should succeed");
931 assert_eq!(parsed.len(), 2);
932 assert_eq!(parsed[0].tag, "BGM");
933 assert_eq!(parsed[0].elements[0].components[0].0, "220");
934 assert_eq!(parsed[1].span, Span::new(10, 18));
935 }
936
937 #[test]
938 fn dangling_release_sequence_is_error() {
939 let input = b"FTX+AAA++dangling?";
940 let err = crate::from_bytes(input)
941 .collect::<Result<Vec<_>, _>>()
942 .expect_err("expected dangling release to fail");
943
944 assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
945 }
946
947 #[test]
948 fn from_reader_reports_dangling_release_sequence() {
949 let input = b"FTX+AAA++dangling?";
950 let err = from_reader(std::io::Cursor::new(input))
951 .expect_err("expected dangling release from reader path");
952 assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
953 }
954
955 #[test]
956 fn from_reader_rejects_invalid_una() {
957 let input = b"UNA::.? 'BGM:220'";
958 let err = from_reader(std::io::Cursor::new(input))
959 .expect_err("invalid UNA should fail reader parsing");
960 assert!(matches!(err, EdifactError::InvalidUna));
961 }
962}