Skip to main content

edifact_rs/
parser.rs

1//! Streaming EDIFACT parser — wraps a [`Tokenizer`] and assembles [`Segment`]s.
2
3use crate::{
4    error::EdifactError,
5    model::{Components, 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
13/// In-progress data element: the components of the repetition currently being
14/// read, plus every repetition already closed by a repetition separator.
15struct PendingElement<'a> {
16    start: Option<usize>,
17    components: Components<'a>,
18    repeats: Vec<Components<'a>>,
19}
20
21impl<'a> PendingElement<'a> {
22    fn new() -> Self {
23        Self {
24            start: None,
25            components: SmallVec::new(),
26            repeats: Vec::new(),
27        }
28    }
29
30    #[inline]
31    fn is_open(&self) -> bool {
32        self.start.is_some()
33    }
34
35    /// Open a new element (or, when already open, a further repetition of it).
36    #[inline]
37    fn open(&mut self, start: usize) {
38        if self.start.is_none() {
39            self.start = Some(start);
40        }
41    }
42
43    /// Close the current repetition and begin the next one.
44    #[inline]
45    fn end_repetition(&mut self) {
46        self.repeats.push(std::mem::take(&mut self.components));
47    }
48
49    /// Close the element and push it onto `elements`, returning its end offset.
50    fn flush(&mut self, elements: &mut Vec<Element<'a>>) -> Option<usize> {
51        let start = self.start.take()?;
52        let mut repeats = std::mem::take(&mut self.repeats);
53        repeats.push(std::mem::take(&mut self.components));
54        // The element ends at the last component of the last *non-empty*
55        // repetition; a trailing repetition separator with no value after it
56        // contributes no span.
57        let end = repeats
58            .iter()
59            .rev()
60            .find_map(|rep| rep.last().map(|(_, span)| span.end))?;
61        let components = repeats.remove(0);
62        elements.push(Element {
63            span: Span::new(start, end),
64            components,
65            repeats,
66        });
67        Some(end)
68    }
69}
70
71fn resolve_release(
72    val: &str,
73    release_char: char,
74    start_offset: usize,
75) -> Result<Cow<'_, str>, EdifactError> {
76    if !val.contains(release_char) {
77        return Ok(Cow::Borrowed(val));
78    }
79    resolve_release_owned(val, release_char, start_offset).map(Cow::Owned)
80}
81
82fn resolve_release_owned(
83    val: &str,
84    release_char: char,
85    start_offset: usize,
86) -> Result<String, EdifactError> {
87    // The unescaped output is always shorter than the input (every release-char
88    // pair shrinks by one), so `val.len()` is a safe upper bound.  However,
89    // allocating exactly `val.len()` wastes capacity for inputs with many
90    // release sequences.  A conservative 75 % heuristic saves memory on
91    // release-heavy values while avoiding reallocation for typical inputs.
92    let cap = val.len() - val.len() / 4;
93    let mut out = String::with_capacity(cap);
94    let mut chars = val.chars();
95    while let Some(c) = chars.next() {
96        if c == release_char {
97            if let Some(escaped) = chars.next() {
98                out.push(escaped);
99            } else {
100                return Err(EdifactError::InvalidReleaseSequence {
101                    offset: start_offset + val.len().saturating_sub(1),
102                });
103            }
104        } else {
105            out.push(c);
106        }
107    }
108    Ok(out)
109}
110
111/// Streaming parser over a [`Tokenizer`].
112///
113/// Implements `Iterator<Item = Result<Segment<'a>, EdifactError>>`.
114/// Each `next()` call produces one fully-assembled segment.
115pub struct Parser<'a> {
116    tokenizer: Tokenizer<'a>,
117    /// Buffered token from a previous `next()` call (tag peeked ahead).
118    peeked: Option<Token<'a>>,
119    /// Release character from the active service string advice.
120    release_char: char,
121}
122
123impl<'a> Parser<'a> {
124    /// Construct a parser from a tokenizer.
125    pub fn new(tokenizer: Tokenizer<'a>) -> Self {
126        let release_char = tokenizer.service_string_advice().release_char as char;
127        Self {
128            tokenizer,
129            peeked: None,
130            release_char,
131        }
132    }
133}
134
135impl<'a> Iterator for Parser<'a> {
136    type Item = Result<Segment<'a>, EdifactError>;
137
138    fn next(&mut self) -> Option<Self::Item> {
139        // Obtain the segment tag (may have been peeked from a previous iteration)
140        let (tag, tag_span) = loop {
141            let tok = match self.peeked.take() {
142                Some(t) => Ok(t),
143                None => self.tokenizer.next()?,
144            };
145            match tok {
146                Ok(Token::SegmentTag { value, span }) => break (value, span),
147                Ok(Token::SegmentTerminator { .. }) => continue, // stray terminator — tolerated (blank line)
148                Ok(Token::DataElement { span, .. })
149                | Ok(Token::ComponentElement { span, .. })
150                | Ok(Token::RepeatElement { span, .. }) => {
151                    return Some(Err(EdifactError::UnexpectedDataToken {
152                        offset: span.start,
153                    }));
154                }
155                Err(e) => return Some(Err(e)),
156            }
157        };
158
159        // Deliberately unallocated: `Element` is a large struct (inline component
160        // storage), so eagerly reserving 8 slots cost well over a kilobyte per
161        // segment while typical EDIFACT segments carry 2–5 elements.  Amortised
162        // growth costs at most two reallocations for the common case.
163        let mut elements: Vec<Element<'a>> = Vec::new();
164        let mut pending = PendingElement::new();
165        let mut segment_end = tag_span.end;
166
167        loop {
168            let tok = match self.tokenizer.next() {
169                Some(Ok(t)) => t,
170                Some(Err(e)) => return Some(Err(e)),
171                None => {
172                    // Input ran out mid-segment: no terminator arrived.
173                    //
174                    // Yielding what was read would let a truncated file parse as
175                    // a complete one — a silent data-integrity failure in a
176                    // format whose trailers declare counts.
177                    if let Some(end) = pending.flush(&mut elements) {
178                        segment_end = end;
179                    }
180                    return Some(Err(EdifactError::UnexpectedEof {
181                        offset: segment_end,
182                    }));
183                }
184            };
185
186            match tok {
187                Token::SegmentTag {
188                    value: next_tag,
189                    span,
190                } => {
191                    // Not reachable through the current tokenizer, which only
192                    // leaves `InSegment` on a terminator (handled below) or on
193                    // an error it reports first.  Kept so the match stays total
194                    // and a future tokenizer change cannot silently splice two
195                    // segments together: a tag where a terminator belongs means
196                    // the current segment was never closed.  The token is saved
197                    // so a caller that continues sees the next segment intact.
198                    self.peeked = Some(Token::SegmentTag {
199                        value: next_tag,
200                        span,
201                    });
202                    if let Some(end) = pending.flush(&mut elements) {
203                        segment_end = end;
204                    }
205                    return Some(Err(EdifactError::UnexpectedEof {
206                        offset: segment_end,
207                    }));
208                }
209                Token::SegmentTerminator { span } => {
210                    pending.flush(&mut elements);
211                    segment_end = span.end;
212                    break;
213                }
214                Token::DataElement { value, span } => {
215                    pending.flush(&mut elements);
216                    let resolved = match resolve_release(value, self.release_char, span.start) {
217                        Ok(v) => v,
218                        Err(error) => return Some(Err(error)),
219                    };
220                    pending.open(span.start);
221                    pending.components.push((resolved, span));
222                }
223                Token::ComponentElement { value, span } => {
224                    // A component before any element opens the first one.
225                    let resolved = match resolve_release(value, self.release_char, span.start) {
226                        Ok(v) => v,
227                        Err(error) => return Some(Err(error)),
228                    };
229                    pending.open(span.start);
230                    pending.components.push((resolved, span));
231                }
232                Token::RepeatElement { value, span } => {
233                    let resolved = match resolve_release(value, self.release_char, span.start) {
234                        Ok(v) => v,
235                        Err(error) => return Some(Err(error)),
236                    };
237                    // A repetition separator before any element is malformed:
238                    // there is nothing to repeat.
239                    if !pending.is_open() {
240                        return Some(Err(EdifactError::UnexpectedDataToken {
241                            offset: span.start,
242                        }));
243                    }
244                    pending.end_repetition();
245                    pending.components.push((resolved, span));
246                }
247            }
248        }
249
250        Some(Ok(Segment {
251            tag: Cow::Borrowed(tag),
252            span: Span::new(tag_span.start, segment_end),
253            tag_span,
254            elements,
255        }))
256    }
257}
258
259/// Configuration for reader-based EDIFACT parsers.
260///
261/// Pass to [`from_reader_with_config`] or [`from_bufread_with_config`] to
262/// override default limits.
263///
264/// # Example
265/// ```
266/// use edifact_rs::{ReaderConfig, from_reader_with_config};
267///
268/// let cfg = ReaderConfig::default().max_segment_bytes(4_096);
269/// let segments: Vec<_> = from_reader_with_config(b"BGM+220+1+9'".as_ref(), cfg)
270///     .collect::<Result<_, _>>()
271///     .unwrap();
272/// assert_eq!(segments[0].tag, "BGM");
273/// ```
274#[derive(Debug, Clone, Copy)]
275pub struct ReaderConfig {
276    /// Maximum allowed segment byte length (excluding the segment terminator).
277    ///
278    /// If a segment accumulates more bytes than this limit without a terminator
279    /// the parser returns [`EdifactError::SegmentTooLong`].  This prevents
280    /// unbounded allocation when processing malformed or adversarially crafted
281    /// input streams.
282    ///
283    /// Default: 65 536 bytes (64 KiB).  Real-world EDIFACT segments are almost
284    /// always below 4 KiB; consider using a tighter limit for untrusted inputs.
285    pub max_segment_bytes: usize,
286    /// Maximum number of segments the input may contain.
287    ///
288    /// A segment beyond this budget yields [`EdifactError::LimitExceeded`].
289    /// Input that ends exactly at the limit is not a violation.
290    ///
291    /// Default: `None` (unlimited).
292    pub max_segments: Option<usize>,
293    /// Maximum number of input bytes the interchange may occupy.
294    ///
295    /// A segment whose end offset passes this budget yields
296    /// [`EdifactError::LimitExceeded`] *instead of* the segment, so no data
297    /// beyond the cap is ever handed to the caller.  Combine with
298    /// `max_segment_bytes` for defence in depth against oversized input.
299    ///
300    /// Default: `None` (unlimited).
301    pub max_input_bytes: Option<u64>,
302    /// Maximum number of EDIFACT messages (`UNH`/`UNT` pairs) the input may contain.
303    ///
304    /// A segment belonging to a message beyond this budget yields
305    /// [`EdifactError::LimitExceeded`].
306    ///
307    /// Default: `None` (unlimited).
308    pub max_messages: Option<usize>,
309    /// Service characters to use instead of reading them from the input.
310    ///
311    /// `None` — the default — discovers them the way ISO 9735-1 says a receiver
312    /// should: from a leading `UNA` if there is one, otherwise the §5.1 defaults
313    /// with the repetition separator resolved from `UNB` S001 DE 0002.
314    ///
315    /// Set this when the input is a **fragment** that carries neither: a single
316    /// message lifted out of an interchange has no `UNA` and no `UNB`, so nothing
317    /// in it records the delimiters its interchange declared, and parsing it with
318    /// the defaults silently mis-splits every value.
319    ///
320    /// Default: `None`.
321    pub service_string_advice: Option<crate::tokenizer::ServiceStringAdvice>,
322}
323
324impl Default for ReaderConfig {
325    fn default() -> Self {
326        Self {
327            max_segment_bytes: 65_536,
328            max_segments: None,
329            max_input_bytes: None,
330            max_messages: None,
331            service_string_advice: None,
332        }
333    }
334}
335
336impl ReaderConfig {
337    /// Set the maximum segment byte length and return `self`.
338    #[must_use]
339    pub fn max_segment_bytes(mut self, limit: usize) -> Self {
340        self.max_segment_bytes = limit;
341        self
342    }
343
344    /// Set the maximum number of segments to yield and return `self`.
345    #[must_use]
346    pub fn max_segments(mut self, limit: usize) -> Self {
347        self.max_segments = Some(limit);
348        self
349    }
350
351    /// Set the maximum total input bytes to consume and return `self`.
352    #[must_use]
353    pub fn max_input_bytes(mut self, limit: u64) -> Self {
354        self.max_input_bytes = Some(limit);
355        self
356    }
357
358    /// Set the maximum number of EDIFACT messages (UNH/UNT pairs) to yield and return `self`.
359    #[must_use]
360    pub fn max_messages(mut self, limit: usize) -> Self {
361        self.max_messages = Some(limit);
362        self
363    }
364
365    /// Parse with these service characters instead of discovering them.
366    ///
367    /// See [`service_string_advice`][Self::service_string_advice].
368    ///
369    /// # Example
370    ///
371    /// ```
372    /// use edifact_rs::{ReaderConfig, ServiceStringAdvice, from_bytes_with_config};
373    ///
374    /// // A message lifted out of an interchange whose UNA declared `;` and `~`.
375    /// let ssa = ServiceStringAdvice::from_bytes(b"UNA:;.? ~")?;
376    /// let config = ReaderConfig::default().with_service_string_advice(ssa);
377    ///
378    /// let segments: Vec<_> = from_bytes_with_config(b"BGM;220;PO-4711~", config)
379    ///     .collect::<Result<Vec<_>, _>>()?;
380    /// assert_eq!(segments[0].element_str(1), Some("PO-4711"));
381    /// # Ok::<(), edifact_rs::EdifactError>(())
382    /// ```
383    #[must_use]
384    pub fn with_service_string_advice(
385        mut self,
386        ssa: crate::tokenizer::ServiceStringAdvice,
387    ) -> Self {
388        self.service_string_advice = Some(ssa);
389        self
390    }
391}
392
393/// Streaming state for [`OwnedSegmentStream`].
394#[derive(Debug, Clone, Copy, PartialEq, Eq)]
395enum StreamState {
396    /// UNA header not yet scanned; must inspect first bytes.
397    Init,
398    /// UNA has been scanned (or was absent); streaming segments.
399    Running,
400    /// A terminal error was encountered; no more items.
401    Done,
402}
403
404/// Streaming iterator over owned segments from a buffered reader.
405///
406/// # Performance
407///
408/// A **fast path** uses [`BufRead::fill_buf`] + `memchr` to locate the segment
409/// terminator within the OS-level read buffer (typically 8 KB) without any
410/// intermediate heap allocation.  The segment bytes are parsed directly from
411/// the buffer slice and converted to an [`OwnedSegment`] in a single pass.
412///
413/// For segments that span read-buffer boundaries the implementation falls back
414/// to the byte-accumulation slow path, which allocates a temporary `Vec<u8>`
415/// and re-tokenizes — the same behaviour as in older versions of the library.
416/// In practice this fallback is rare because the default `BufReader` buffer
417/// (8 KB) is far larger than a typical EDIFACT segment (<300 bytes).
418///
419/// Configure limits via [`ReaderConfig`] and [`from_reader_with_config`] /
420/// [`from_bufread_with_config`].
421pub struct OwnedSegmentStream<R: BufRead> {
422    reader: R,
423    ssa: crate::tokenizer::ServiceStringAdvice,
424    state: StreamState,
425    stream_offset: u64,
426    config: ReaderConfig,
427    /// Number of segments successfully yielded so far.
428    segments_yielded: usize,
429    /// Number of complete EDIFACT messages (UNH/UNT pairs) seen so far.
430    messages_yielded: usize,
431    /// Whether the last segment tag was UNH (inside a message).
432    in_message: bool,
433    /// Total bytes consumed from the reader (UNA header + segment data).
434    bytes_consumed: u64,
435    /// Whether the delimiters are already settled — because the caller supplied
436    /// them, or because a `UNA` stated them.  While this is `false` the `UNB`'s
437    /// own syntax version still gets to choose the repetition separator.
438    delimiters_settled: bool,
439}
440
441impl<R: BufRead> OwnedSegmentStream<R> {
442    fn new(reader: R) -> Self {
443        Self::with_config(reader, ReaderConfig::default())
444    }
445
446    fn with_config(reader: R, config: ReaderConfig) -> Self {
447        let (ssa, delimiters_settled) = match config.service_string_advice {
448            Some(ssa) => (ssa, true),
449            None => (crate::tokenizer::ServiceStringAdvice::default(), false),
450        };
451        Self {
452            reader,
453            ssa,
454            state: StreamState::Init,
455            stream_offset: 0,
456            config,
457            segments_yielded: 0,
458            messages_yielded: 0,
459            in_message: false,
460            bytes_consumed: 0,
461            delimiters_settled,
462        }
463    }
464
465    /// Adopt the repetition separator implied by a `UNB`'s syntax version.
466    ///
467    /// Only reached when nothing has already settled the delimiters: a `UNA`
468    /// states all six explicitly, and a caller-supplied advice is by definition
469    /// the last word.  Version 4 is the version that has a repetition separator
470    /// at all, and its default is `*` (ISO 9735-1 §5.1).
471    ///
472    /// The `UNB` itself has no repeating data elements, so re-reading it under
473    /// the updated advice would not change a thing.
474    fn adopt_syntax_version(&mut self, unb: &OwnedSegment) {
475        self.delimiters_settled = true;
476        if unb.component_str(0, 1) == Some("4") {
477            self.ssa.repetition_sep = b'*';
478        }
479    }
480
481    /// Check the whole-input budgets against a segment that has just been read.
482    ///
483    /// Checking *after* the read rather than before is what makes an input that
484    /// ends exactly at a limit finish cleanly: a budget is only exceeded when
485    /// there was genuinely more data than the caller allowed.
486    fn check_limits(&self, tag: &str) -> Option<EdifactError> {
487        if let Some(max) = self.config.max_segments {
488            if self.segments_yielded >= max {
489                return Some(EdifactError::LimitExceeded {
490                    limit: "max_segments",
491                    max: max as u64,
492                });
493            }
494        }
495        // Only a `UNH` opens a message, so only a `UNH` can push the count past
496        // the budget.  Testing every segment would trip on the interchange
497        // trailer, which belongs to no message.
498        if let Some(max) = self.config.max_messages {
499            if tag == "UNH" && self.messages_yielded >= max {
500                return Some(EdifactError::LimitExceeded {
501                    limit: "max_messages",
502                    max: max as u64,
503                });
504            }
505        }
506        if let Some(max) = self.config.max_input_bytes {
507            if self.bytes_consumed > max {
508                return Some(EdifactError::LimitExceeded {
509                    limit: "max_input_bytes",
510                    max,
511                });
512            }
513        }
514        None
515    }
516
517    /// Record a yielded segment against the segment and message counters.
518    ///
519    /// Only a `UNT` that closes a `UNH` already seen advances the message
520    /// counter; a bare `UNT` is malformed and must not inflate it.
521    fn account(&mut self, tag: &str) {
522        self.segments_yielded += 1;
523        if tag == "UNT" {
524            if self.in_message {
525                self.messages_yielded += 1;
526            }
527            self.in_message = false;
528        } else if tag == "UNH" {
529            self.in_message = true;
530        }
531    }
532}
533
534// ── fast-path helpers ─────────────────────────────────────────────────────────
535
536/// Outcome of a single-buffer segment extraction attempt.
537enum FastSegment {
538    /// Segment parsed; second field = bytes to consume (content + terminator).
539    Parsed(OwnedSegment, usize),
540    /// Only whitespace or an isolated terminator; bytes to skip and continue.
541    Skip(usize),
542    /// Terminator not present in the current buffer; caller must use slow path.
543    NeedMore,
544    /// Buffer is empty — no more input.
545    Eof,
546    /// Parse error.
547    Err(EdifactError),
548}
549
550/// Return the byte offset of the first **unescaped** occurrence of `term` in `buf`.
551///
552/// A byte is *escaped* when it is immediately preceded by `release` (e.g.
553/// `?'` escapes `'`).  Two consecutive release chars cancel each other, so
554/// `??'` contains an *unescaped* `'`.
555///
556/// # Complexity
557///
558/// O(n) in the length of `buf`.  `memchr2` is used to fast-scan past bytes
559/// that are neither `release` nor `term`, so SIMD acceleration applies on
560/// platforms where `memchr` provides it.
561fn find_unescaped_term(buf: &[u8], term: u8, release: u8) -> Option<usize> {
562    let mut i = 0;
563    while i < buf.len() {
564        // Fast-skip to the next byte that might be a release char or terminator.
565        let rel = memchr2(release, term, &buf[i..])?;
566        let pos = i + rel;
567        if buf[pos] == release {
568            // Release char: the next byte is escaped — skip both.
569            i = pos + 2;
570        } else {
571            // Unescaped terminator found.
572            return Some(pos);
573        }
574    }
575    None
576}
577
578/// Try to parse one segment directly from the `BufRead` buffer.
579///
580/// This function borrows `reader` only for the duration of the call.  After it
581/// returns the caller is free to call `reader.consume(n)`.
582fn try_fast_segment<R: BufRead>(
583    reader: &mut R,
584    ssa: crate::tokenizer::ServiceStringAdvice,
585    seg_start: usize,
586    max_segment_bytes: usize,
587) -> FastSegment {
588    let buf = match reader.fill_buf() {
589        Ok(b) => b,
590        Err(e) => return FastSegment::Err(e.into()),
591    };
592
593    if buf.is_empty() {
594        return FastSegment::Eof;
595    }
596
597    let Some(pos) = find_unescaped_term(buf, ssa.segment_term, ssa.release_char) else {
598        return FastSegment::NeedMore;
599    };
600
601    // Enforce the segment-size guard *before* any allocation.
602    // `pos` is the index of the terminator byte, so the segment body is `buf[..pos]`.
603    if pos > max_segment_bytes {
604        return FastSegment::Err(EdifactError::SegmentTooLong {
605            offset: seg_start,
606            limit: max_segment_bytes,
607        });
608    }
609
610    // `buf[..pos]` is the segment content without the terminator.
611    let seg_bytes = &buf[..pos];
612
613    // Skip isolated terminators / pure-whitespace slots between segments.
614    if seg_bytes
615        .iter()
616        .all(|&b| matches!(b, b' ' | b'\t' | b'\r' | b'\n'))
617    {
618        return FastSegment::Skip(pos + 1);
619    }
620
621    // Parse directly from the buffer slice — zero intermediate allocation.
622    // Include the terminator byte so the parser sees a `SegmentTerminator`
623    // token and records a span that is consistent with the `from_bytes` path.
624    // `for_segment` (not `with_limit`) because this slice is one bare segment:
625    // the whole-interchange constructors would mistake a segment tagged `UNA`
626    // for a service string advice and skip nine bytes of it.  The configured
627    // `max_segment_bytes` is passed through so the tokenizer never applies a
628    // tighter cap than the caller allowed.
629    let tok = Tokenizer::for_segment(&buf[..pos + 1], ssa, max_segment_bytes);
630    let mut parser_iter = Parser::new(tok);
631    match parser_iter.next() {
632        None => FastSegment::Skip(pos + 1),
633        Some(Err(e)) => FastSegment::Err(e),
634        Some(Ok(s)) => FastSegment::Parsed(s.into_owned().offset(seg_start), pos + 1),
635    }
636    // `buf` borrow released here — `reader.consume()` is safe to call in the caller.
637}
638
639// ── Iterator impl ─────────────────────────────────────────────────────────────
640
641impl<R: BufRead> Iterator for OwnedSegmentStream<R> {
642    type Item = Result<OwnedSegment, EdifactError>;
643
644    fn next(&mut self) -> Option<Self::Item> {
645        if self.state == StreamState::Done {
646            return None;
647        }
648
649        loop {
650            // ── Fast path (after UNA has been consumed) ───────────────────
651            if self.state == StreamState::Running {
652                let seg_start = self.stream_offset;
653                match try_fast_segment(
654                    &mut self.reader,
655                    self.ssa,
656                    // Saturate rather than wrap on 32-bit targets; Span offsets
657                    // are `usize` so streams > 4 GiB on 32-bit produce clamped
658                    // (but monotonic) diagnostic positions.
659                    seg_start.min(usize::MAX as u64) as usize,
660                    self.config.max_segment_bytes,
661                ) {
662                    FastSegment::Parsed(seg, n) => {
663                        let n = n as u64;
664                        self.reader.consume(n as usize);
665                        self.stream_offset += n;
666                        self.bytes_consumed = self.stream_offset;
667                        if let Some(error) = self.check_limits(&seg.tag) {
668                            self.state = StreamState::Done;
669                            return Some(Err(error));
670                        }
671                        if !self.delimiters_settled && seg.tag == "UNB" {
672                            self.adopt_syntax_version(&seg);
673                        }
674                        self.account(&seg.tag);
675                        return Some(Ok(seg));
676                    }
677                    FastSegment::Skip(n) => {
678                        let n = n as u64;
679                        self.reader.consume(n as usize);
680                        self.stream_offset += n;
681                        self.bytes_consumed = self.stream_offset;
682                        continue;
683                    }
684                    FastSegment::Eof => return None,
685                    FastSegment::Err(e) => {
686                        self.state = StreamState::Done;
687                        return Some(Err(e));
688                    }
689                    FastSegment::NeedMore => {
690                        // Segment spans buffer boundary — fall through to slow path.
691                    }
692                }
693            }
694
695            // ── Slow path: byte accumulation (also handles UNA header) ────
696            if self.state == StreamState::Init {
697                // Strip a byte order mark before anything tries to read a tag
698                // out of it.  The slice path does the same (see
699                // `tokenizer::prologue_len`), and the two must agree or the
700                // identical file parses only from one of them.
701                match self.reader.fill_buf() {
702                    Ok(buf) if buf.starts_with(&crate::tokenizer::UTF8_BOM) => {
703                        self.reader.consume(3);
704                        self.stream_offset += 3;
705                        self.bytes_consumed = self.stream_offset;
706                    }
707                    Ok(_) => {}
708                    Err(error) => {
709                        self.state = StreamState::Done;
710                        return Some(Err(error.into()));
711                    }
712                }
713            }
714            let mut scanned = self.state != StreamState::Init;
715            // `read_next_raw_segment` tracks offset as `usize` for segment
716            // start positions; sync back to the `u64` field afterward.
717            // Saturate rather than wrap on 32-bit targets (same rationale as
718            // the fast-path cast above).
719            let mut slow_offset: usize = self.stream_offset.min(usize::MAX as u64) as usize;
720            let mut raw = match read_next_raw_segment(
721                &mut self.reader,
722                &mut self.ssa,
723                &mut scanned,
724                &mut slow_offset,
725                self.config.max_segment_bytes,
726                &mut self.delimiters_settled,
727                self.config.service_string_advice.is_some(),
728            ) {
729                Ok(Some(r)) => r,
730                Ok(None) => return None,
731                Err(e) => {
732                    self.state = StreamState::Done;
733                    return Some(Err(e));
734                }
735            };
736            self.stream_offset = slow_offset as u64;
737            if scanned {
738                self.state = StreamState::Running;
739            }
740            self.bytes_consumed = self.stream_offset;
741
742            raw.bytes.push(self.ssa.segment_term);
743            // `for_segment`, and with the configured `max_segment_bytes`: this
744            // buffer holds one bare segment, so neither the UNA skip nor the
745            // 64 KiB default of `Tokenizer::new` applies.
746            let tok = Tokenizer::for_segment(
747                raw.bytes.as_slice(),
748                self.ssa,
749                self.config.max_segment_bytes,
750            );
751            let mut parser_iter = Parser::new(tok);
752            match parser_iter.next() {
753                Some(Ok(s)) => {
754                    let seg = s.into_owned().offset(raw.start_offset);
755                    if let Some(error) = self.check_limits(&seg.tag) {
756                        self.state = StreamState::Done;
757                        return Some(Err(error));
758                    }
759                    if !self.delimiters_settled && seg.tag == "UNB" {
760                        self.adopt_syntax_version(&seg);
761                    }
762                    self.account(&seg.tag);
763                    return Some(Ok(seg));
764                }
765                Some(Err(e)) => {
766                    self.state = StreamState::Done;
767                    return Some(Err(e));
768                }
769                None => {} // Empty segment — loop back.
770            }
771        }
772    }
773}
774
775/// Parse an already-buffered reader into a lazy iterator of [`OwnedSegment`]s.
776///
777/// The same as [`from_reader`][crate::from_reader] without the internal
778/// [`BufReader`] wrapper, for callers who already hold a [`BufRead`] and would
779/// otherwise pay for a second layer of buffering.
780///
781/// # Errors
782///
783/// Each `next()` yields a parse or I/O failure as `Some(Err(_))`.
784pub fn from_bufread<R: BufRead>(reader: R) -> OwnedSegmentStream<R> {
785    OwnedSegmentStream::new(reader)
786}
787
788/// [`from_bufread`] with explicit [`ReaderConfig`] limits.
789pub fn from_bufread_with_config<R: BufRead>(
790    reader: R,
791    config: ReaderConfig,
792) -> OwnedSegmentStream<R> {
793    OwnedSegmentStream::with_config(reader, config)
794}
795
796/// Parse EDIFACT from an arbitrary reader as a streaming iterator.
797pub fn from_reader_stream<R: Read>(reader: R) -> OwnedSegmentStream<BufReader<R>> {
798    from_bufread(BufReader::new(reader))
799}
800
801/// Parse EDIFACT from an arbitrary reader as a streaming iterator with custom config.
802///
803/// # Example
804/// ```
805/// use edifact_rs::{ReaderConfig, from_reader_with_config};
806///
807/// let cfg = ReaderConfig::default().max_segment_bytes(4_096);
808/// let segs: Vec<_> = from_reader_with_config(b"BGM+220+1+9'".as_ref(), cfg)
809///     .collect::<Result<_, _>>()
810///     .unwrap();
811/// assert_eq!(segs[0].tag, "BGM");
812/// ```
813pub fn from_reader_with_config<R: Read>(
814    reader: R,
815    config: ReaderConfig,
816) -> OwnedSegmentStream<BufReader<R>> {
817    from_bufread_with_config(BufReader::new(reader), config)
818}
819
820#[allow(clippy::too_many_arguments)]
821fn read_next_raw_segment<R: BufRead>(
822    reader: &mut R,
823    ssa: &mut crate::tokenizer::ServiceStringAdvice,
824    scanned_header: &mut bool,
825    stream_offset: &mut usize,
826    max_segment_bytes: usize,
827    delimiters_settled: &mut bool,
828    advice_overridden: bool,
829) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
830    loop {
831        let Some((first_offset, first)) = read_next_non_ws_byte(reader, stream_offset)? else {
832            return Ok(None);
833        };
834
835        if !*scanned_header && first == b'U' {
836            let second = read_required_byte(reader, stream_offset)?;
837            let third = read_required_byte(reader, stream_offset)?;
838            if second == b'N' && third == b'A' {
839                let mut una = [0u8; 9];
840                una[0] = b'U';
841                una[1] = b'N';
842                una[2] = b'A';
843                for slot in una.iter_mut().skip(3) {
844                    *slot = read_required_byte(reader, stream_offset)?;
845                }
846                let declared = crate::tokenizer::ServiceStringAdvice {
847                    component_sep: una[3],
848                    element_sep: una[4],
849                    decimal_mark: una[5],
850                    release_char: una[6],
851                    repetition_sep: una[7],
852                    segment_term: una[8],
853                };
854                if !declared.is_valid() {
855                    return Err(EdifactError::InvalidUna);
856                }
857                // A malformed UNA is still an error when the caller supplied its
858                // own advice — the input is broken either way — but the caller's
859                // choice, not the header's, is what parsing then uses.
860                if !advice_overridden {
861                    *ssa = declared;
862                }
863                *delimiters_settled = true;
864                *scanned_header = true;
865                continue;
866            }
867
868            *scanned_header = true;
869            return read_remainder_of_segment(
870                reader,
871                ssa,
872                crate::tokenizer::RawSegment {
873                    bytes: vec![first, second, third],
874                    start_offset: first_offset,
875                },
876                stream_offset,
877                max_segment_bytes,
878            );
879        }
880
881        *scanned_header = true;
882        return read_remainder_of_segment(
883            reader,
884            ssa,
885            crate::tokenizer::RawSegment {
886                bytes: vec![first],
887                start_offset: first_offset,
888            },
889            stream_offset,
890            max_segment_bytes,
891        );
892    }
893}
894
895fn read_remainder_of_segment<R: BufRead>(
896    reader: &mut R,
897    ssa: &crate::tokenizer::ServiceStringAdvice,
898    mut out: crate::tokenizer::RawSegment,
899    stream_offset: &mut usize,
900    max_segment_bytes: usize,
901) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
902    let mut escaped = false;
903    loop {
904        // Strictly greater-than, matching the fast path (`pos > max_segment_bytes`)
905        // and the tokenizer guard.  Using `>=` here made a segment of exactly
906        // `max_segment_bytes` bytes parse on the fast path but fail on the slow
907        // path, so identical input succeeded or failed depending only on whether
908        // it happened to straddle a read-buffer boundary.
909        if out.bytes.len() > max_segment_bytes {
910            return Err(EdifactError::SegmentTooLong {
911                offset: out.start_offset,
912                limit: max_segment_bytes,
913            });
914        }
915        let Some(byte) = read_next_byte(reader, stream_offset)? else {
916            return if out.bytes.is_empty() {
917                Ok(None)
918            } else if escaped {
919                Err(EdifactError::InvalidReleaseSequence {
920                    offset: out.start_offset + out.bytes.len().saturating_sub(1),
921                })
922            } else {
923                Err(EdifactError::UnexpectedEof {
924                    offset: out.start_offset + out.bytes.len(),
925                })
926            };
927        };
928
929        if !escaped && byte == ssa.segment_term {
930            return Ok(Some(out));
931        }
932
933        if !escaped && byte == ssa.release_char {
934            escaped = true;
935            out.bytes.push(byte);
936            continue;
937        }
938
939        escaped = false;
940        out.bytes.push(byte);
941    }
942}
943
944fn read_next_byte<R: BufRead>(
945    reader: &mut R,
946    stream_offset: &mut usize,
947) -> Result<Option<u8>, EdifactError> {
948    let buf = reader.fill_buf()?;
949    if buf.is_empty() {
950        return Ok(None);
951    }
952
953    let byte = buf[0];
954    reader.consume(1);
955    // Saturating add: on 32-bit targets `stream_offset` is a `usize` clamped from a
956    // `u64` field.  Plain `+= 1` would wrap to 0 once the counter reaches `usize::MAX`
957    // and corrupt subsequent span diagnostics / `bytes_consumed` accounting.
958    // Using an explicit local avoids relying on `&mut` auto-deref evaluation order.
959    let next_offset = stream_offset.saturating_add(1);
960    *stream_offset = next_offset;
961    Ok(Some(byte))
962}
963
964fn read_required_byte<R: BufRead>(
965    reader: &mut R,
966    stream_offset: &mut usize,
967) -> Result<u8, EdifactError> {
968    read_next_byte(reader, stream_offset)?.ok_or(EdifactError::UnexpectedEof {
969        offset: *stream_offset,
970    })
971}
972
973fn read_next_non_ws_byte<R: BufRead>(
974    reader: &mut R,
975    stream_offset: &mut usize,
976) -> Result<Option<(usize, u8)>, EdifactError> {
977    loop {
978        let current_offset = *stream_offset;
979        let Some(byte) = read_next_byte(reader, stream_offset)? else {
980            return Ok(None);
981        };
982        if !matches!(byte, b' ' | b'\t' | b'\r' | b'\n') {
983            return Ok(Some((current_offset, byte)));
984        }
985    }
986}
987
988#[cfg(test)]
989mod tests {
990    use super::*;
991    use crate::tokenizer::ServiceStringAdvice;
992
993    fn parse_all(input: &[u8]) -> Vec<Segment<'_>> {
994        let ssa = ServiceStringAdvice::from_bytes_unchecked(input);
995        let tok = Tokenizer::new(input, ssa);
996        Parser::new(tok)
997            .collect::<Result<Vec<_>, _>>()
998            .expect("parse failed")
999    }
1000
1001    #[test]
1002    fn parses_unb_unz() {
1003        let input = b"UNB+UNOA:1+SENDER+RECEIVER+200101:0900+1'UNZ+0+1'";
1004        let segs = parse_all(input);
1005        assert_eq!(segs.len(), 2);
1006        assert_eq!(segs[0].tag, "UNB");
1007        assert_eq!(segs[1].tag, "UNZ");
1008        assert_eq!(segs[0].tag_span, Span::new(0, 3));
1009        assert_eq!(segs[0].span, Span::new(0, 41));
1010    }
1011
1012    #[test]
1013    fn element_access() {
1014        let input = b"BGM+220+ORDER123+9'";
1015        let segs = parse_all(input);
1016        assert_eq!(segs[0].element_str(0), Some("220"));
1017        assert_eq!(segs[0].element_str(1), Some("ORDER123"));
1018    }
1019
1020    #[test]
1021    fn component_access() {
1022        let input = b"DTM+137:20200101:102'";
1023        let segs = parse_all(input);
1024        let dtm = &segs[0];
1025        assert_eq!(dtm.get_element(0).unwrap().get_component(0), Some("137"));
1026        assert_eq!(
1027            dtm.get_element(0).unwrap().get_component(1),
1028            Some("20200101")
1029        );
1030        assert_eq!(dtm.get_element(0).unwrap().get_component(2), Some("102"));
1031    }
1032
1033    #[test]
1034    fn release_char_resolved() {
1035        let input = b"FTX+AAA++test?+value'";
1036        let segs = parse_all(input);
1037        assert_eq!(segs[0].element_str(2), Some("test+value"));
1038        assert_eq!(
1039            segs[0].get_element(2).unwrap().component_span(0),
1040            Some(Span::new(9, 20))
1041        );
1042    }
1043
1044    #[test]
1045    fn reader_path_preserves_custom_una_delimiters() {
1046        let input = b"UNA:;.? 'BGM;220;test?;value'";
1047        let segments: Vec<_> =
1048            super::from_bufread(std::io::BufReader::new(std::io::Cursor::new(input)))
1049                .collect::<Result<_, _>>()
1050                .expect("reader parse should succeed");
1051        let bgm = segments
1052            .iter()
1053            .find(|segment| segment.tag == "BGM")
1054            .expect("BGM segment should be present");
1055        assert_eq!(bgm.elements[0].components[0].0, "220");
1056        assert_eq!(bgm.elements[1].components[0].0, "test;value");
1057    }
1058
1059    #[test]
1060    fn arbitrary_bytes_no_panic() {
1061        // This is the stable no-panic property — arbitrary input must not panic
1062        let garbage: &[u8] = b"\xff\x00\x01\x02ABC+++'''???";
1063        let _ = crate::from_bytes(garbage).collect::<Vec<_>>();
1064    }
1065
1066    #[test]
1067    fn from_reader_handles_chunk_boundaries() {
1068        let input = b"UNA:+.? 'BGM+220+test?+value'UNT+2+1'";
1069        let reader = std::io::BufReader::with_capacity(5, std::io::Cursor::new(input));
1070        let parsed = from_bufread(reader)
1071            .collect::<Result<Vec<_>, _>>()
1072            .expect("reader parsing should succeed");
1073        assert_eq!(parsed.len(), 2);
1074        assert_eq!(parsed[0].tag, "BGM");
1075        assert_eq!(parsed[0].elements[1].components[0].0, "test+value");
1076        assert_eq!(parsed[1].tag, "UNT");
1077    }
1078
1079    #[test]
1080    fn from_reader_without_una_uses_default_delimiters() {
1081        let input = b"BGM+220+X'UNT+2+1'";
1082        let parsed = crate::from_reader(std::io::Cursor::new(input))
1083            .collect::<Result<Vec<_>, _>>()
1084            .expect("reader parsing should succeed");
1085        assert_eq!(parsed.len(), 2);
1086        assert_eq!(parsed[0].tag, "BGM");
1087        assert_eq!(parsed[0].elements[0].components[0].0, "220");
1088        assert_eq!(parsed[1].span, Span::new(10, 18));
1089    }
1090
1091    #[test]
1092    fn dangling_release_sequence_is_error() {
1093        let input = b"FTX+AAA++dangling?";
1094        let err = crate::from_bytes(input)
1095            .collect::<Result<Vec<_>, _>>()
1096            .expect_err("expected dangling release to fail");
1097
1098        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
1099    }
1100
1101    #[test]
1102    fn from_reader_reports_dangling_release_sequence() {
1103        let input = b"FTX+AAA++dangling?";
1104        let err = crate::from_reader(std::io::Cursor::new(input))
1105            .collect::<Result<Vec<_>, _>>()
1106            .expect_err("expected dangling release from reader path");
1107        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
1108    }
1109
1110    #[test]
1111    fn a_segment_tagged_una_parses_the_same_on_both_paths() {
1112        // `UNA` is three ASCII uppercase letters, so it is a syntactically legal
1113        // segment tag.  The reader re-tokenizes each segment from its own slice,
1114        // where the whole-interchange "skip nine bytes of service string advice"
1115        // heuristic is wrong: it ate the tag and everything after it, and the
1116        // identical bytes that parsed cleanly through `from_bytes` came back as
1117        // `InvalidSegmentTag` through a reader.
1118        let input = b"BGM+220'UNA+XXXXXX'BGM+221'";
1119
1120        let from_slice: Vec<_> = crate::from_bytes(input)
1121            .collect::<Result<Vec<_>, _>>()
1122            .expect("slice path");
1123        let from_reader = crate::from_reader(std::io::Cursor::new(&input[..]))
1124            .collect::<Result<Vec<_>, _>>()
1125            .expect("reader path must agree");
1126
1127        assert_eq!(
1128            from_slice.iter().map(|s| s.tag()).collect::<Vec<_>>(),
1129            from_reader.iter().map(|s| s.tag()).collect::<Vec<_>>(),
1130        );
1131        assert_eq!(from_reader[1].element_str(0), Some("XXXXXX"));
1132    }
1133
1134    #[test]
1135    fn the_reader_adopts_the_repetition_separator_from_the_unb_syntax_version() {
1136        // The slice path resolves this before tokenizing; the reader can only
1137        // learn it once the UNB has been parsed.  Both must agree.
1138        let input = b"UNB+UNOC:4+S+R+260101:0900+IC1'RFF+ON:1*ON:2'UNZ+0+IC1'";
1139
1140        let from_slice: Vec<_> = crate::from_bytes(input)
1141            .collect::<Result<Vec<_>, _>>()
1142            .expect("slice path");
1143        let from_reader = crate::from_reader(std::io::Cursor::new(&input[..]))
1144            .collect::<Result<Vec<_>, _>>()
1145            .expect("reader path");
1146
1147        assert_eq!(from_slice[1].get_element(0).unwrap().repeat_count(), 2);
1148        assert_eq!(from_reader[1].elements[0].repeat_count(), 2);
1149    }
1150
1151    #[test]
1152    fn an_explicit_service_string_advice_parses_a_fragment_with_no_header() {
1153        // A message lifted out of an interchange carries neither UNA nor UNB, so
1154        // nothing in it records the delimiters — the caller has to supply them.
1155        let ssa = ServiceStringAdvice::from_bytes(b"UNA:;.? ~").expect("UNA");
1156        let config = ReaderConfig::default().with_service_string_advice(ssa);
1157
1158        let from_slice: Vec<_> = crate::from_bytes_with_config(b"BGM;220;PO-4711~", config)
1159            .collect::<Result<Vec<_>, _>>()
1160            .expect("slice path");
1161        assert_eq!(from_slice[0].element_str(1), Some("PO-4711"));
1162
1163        let from_reader: Vec<_> =
1164            from_reader_with_config(std::io::Cursor::new(b"BGM;220;PO-4711~"), config)
1165                .collect::<Result<Vec<_>, _>>()
1166                .expect("reader path");
1167        assert_eq!(from_reader[0].element_str(1), Some("PO-4711"));
1168    }
1169
1170    #[test]
1171    fn an_explicit_service_string_advice_outranks_the_una_in_the_input() {
1172        let ssa = ServiceStringAdvice::default();
1173        let config = ReaderConfig::default().with_service_string_advice(ssa);
1174        // The UNA declares `;`, the caller insists on `+`.
1175        let input = b"UNA:;.? 'BGM+220'";
1176
1177        for segments in [
1178            crate::from_bytes_with_config(input, config)
1179                .map(|r| r.map(crate::Segment::into_owned))
1180                .collect::<Result<Vec<_>, _>>()
1181                .expect("slice path"),
1182            from_reader_with_config(std::io::Cursor::new(&input[..]), config)
1183                .collect::<Result<Vec<_>, _>>()
1184                .expect("reader path"),
1185        ] {
1186            assert_eq!(segments[0].element_str(0).unwrap(), "220");
1187        }
1188    }
1189
1190    /// The slice and reader paths must agree byte for byte: the same file has
1191    /// to parse — or fail — identically however it is handed in.
1192    #[test]
1193    fn both_parsing_paths_agree_on_every_edge_case() {
1194        fn bom(rest: &[u8]) -> Vec<u8> {
1195            let mut v = crate::tokenizer::UTF8_BOM.to_vec();
1196            v.extend_from_slice(rest);
1197            v
1198        }
1199
1200        let cases: Vec<Vec<u8>> = vec![
1201            b"BGM+220".to_vec(),                           // truncated: no terminator
1202            b"BGM+220'".to_vec(),                          // the ordinary case
1203            b"BGM+220'\r\n".to_vec(),                      // trailing whitespace
1204            b"\n  UNA:+.? 'BGM+220'".to_vec(),             // UNA behind whitespace
1205            bom(b"UNA:+.? 'BGM+220'"),                     // UNA behind a BOM
1206            bom(b"UNB+UNOA:1+S+R+200101:0900+1'UNZ+0+1'"), // UNB behind a BOM
1207            b"UNB+UNOC:4+S+R+260101:0900+I'RFF+ON:1*ON:2'UNZ+0+I'".to_vec(),
1208            b"FTX+AAA++dangling?".to_vec(), // dangling release
1209            b"UNA::.? 'BGM:220'".to_vec(),  // malformed UNA
1210        ];
1211
1212        for input in cases {
1213            let readable = String::from_utf8_lossy(&input).into_owned();
1214            let sliced: Result<Vec<OwnedSegment>, _> = crate::from_bytes(&input)
1215                .map(|r| r.map(|s| s.into_owned()))
1216                .collect();
1217            let streamed: Result<Vec<OwnedSegment>, _> =
1218                crate::from_reader(std::io::Cursor::new(input.clone())).collect();
1219            assert_eq!(
1220                format!("{sliced:?}"),
1221                format!("{streamed:?}"),
1222                "slice and reader paths disagree on {readable:?}",
1223            );
1224        }
1225    }
1226
1227    #[test]
1228    fn a_segment_with_no_terminator_is_rejected_not_accepted() {
1229        // Accepting it would let a file cut off mid-transfer parse as a
1230        // complete interchange whose UNZ count happens to agree.
1231        for input in [
1232            &b"BGM+220"[..],
1233            &b"UNB+UNOA:1+S+R+200101:0900+1'UNZ+0+1"[..],
1234        ] {
1235            let err = crate::from_bytes(input)
1236                .collect::<Result<Vec<_>, _>>()
1237                .expect_err("an unterminated segment must not parse");
1238            assert!(
1239                matches!(err, EdifactError::UnexpectedEof { .. }),
1240                "expected UnexpectedEof for {:?}, got {err:?}",
1241                std::str::from_utf8(input).unwrap(),
1242            );
1243        }
1244    }
1245
1246    #[test]
1247    fn a_byte_order_mark_does_not_hide_the_first_segment() {
1248        let mut input = crate::tokenizer::UTF8_BOM.to_vec();
1249        input.extend_from_slice(b"UNB+UNOA:1+S+R+200101:0900+1'UNZ+0+1'");
1250        let segments: Vec<_> = crate::from_bytes(&input)
1251            .collect::<Result<Vec<_>, _>>()
1252            .expect("a BOM is prologue, not a segment tag");
1253        assert_eq!(
1254            segments.iter().map(Segment::tag).collect::<Vec<_>>(),
1255            ["UNB", "UNZ"],
1256        );
1257    }
1258
1259    #[test]
1260    fn from_reader_rejects_invalid_una() {
1261        let input = b"UNA::.? 'BGM:220'";
1262        let err = crate::from_reader(std::io::Cursor::new(input))
1263            .collect::<Result<Vec<_>, _>>()
1264            .expect_err("invalid UNA should fail reader parsing");
1265        assert!(matches!(err, EdifactError::InvalidUna));
1266    }
1267}