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                    // EOF — flush whatever we have.
173                    if let Some(end) = pending.flush(&mut elements) {
174                        segment_end = end;
175                    }
176                    break;
177                }
178            };
179
180            match tok {
181                Token::SegmentTag {
182                    value: next_tag,
183                    span,
184                } => {
185                    // We consumed the first token of the *next* segment; save it.
186                    self.peeked = Some(Token::SegmentTag {
187                        value: next_tag,
188                        span,
189                    });
190                    if let Some(end) = pending.flush(&mut elements) {
191                        segment_end = end;
192                    }
193                    break;
194                }
195                Token::SegmentTerminator { span } => {
196                    pending.flush(&mut elements);
197                    segment_end = span.end;
198                    break;
199                }
200                Token::DataElement { value, span } => {
201                    pending.flush(&mut elements);
202                    let resolved = match resolve_release(value, self.release_char, span.start) {
203                        Ok(v) => v,
204                        Err(error) => return Some(Err(error)),
205                    };
206                    pending.open(span.start);
207                    pending.components.push((resolved, span));
208                }
209                Token::ComponentElement { value, span } => {
210                    // A component before any element opens the first one.
211                    let resolved = match resolve_release(value, self.release_char, span.start) {
212                        Ok(v) => v,
213                        Err(error) => return Some(Err(error)),
214                    };
215                    pending.open(span.start);
216                    pending.components.push((resolved, span));
217                }
218                Token::RepeatElement { value, span } => {
219                    let resolved = match resolve_release(value, self.release_char, span.start) {
220                        Ok(v) => v,
221                        Err(error) => return Some(Err(error)),
222                    };
223                    // A repetition separator before any element is malformed:
224                    // there is nothing to repeat.
225                    if !pending.is_open() {
226                        return Some(Err(EdifactError::UnexpectedDataToken {
227                            offset: span.start,
228                        }));
229                    }
230                    pending.end_repetition();
231                    pending.components.push((resolved, span));
232                }
233            }
234        }
235
236        Some(Ok(Segment {
237            tag,
238            span: Span::new(tag_span.start, segment_end),
239            tag_span,
240            elements,
241        }))
242    }
243}
244
245/// Parse EDIFACT from an arbitrary reader.
246///
247/// This path is optimized for bounded-memory ingest and returns owned segments,
248/// allowing the parser to advance across chunk boundaries without requiring a
249/// fully-buffered input slice.
250pub fn from_reader<R: Read>(reader: R) -> Result<Vec<OwnedSegment>, EdifactError> {
251    from_reader_stream(reader).collect()
252}
253
254/// Parse EDIFACT from a buffered reader.
255pub fn from_bufread<R: BufRead>(reader: R) -> Result<Vec<OwnedSegment>, EdifactError> {
256    from_bufread_stream(reader).collect()
257}
258
259/// Configuration for reader-based EDIFACT parsers.
260///
261/// Pass to [`from_reader_with_config`] or [`from_bufread_stream_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_stream_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(OwnedSegment::from(s).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            let mut scanned = self.state != StreamState::Init;
697            // `read_next_raw_segment` tracks offset as `usize` for segment
698            // start positions; sync back to the `u64` field afterward.
699            // Saturate rather than wrap on 32-bit targets (same rationale as
700            // the fast-path cast above).
701            let mut slow_offset: usize = self.stream_offset.min(usize::MAX as u64) as usize;
702            let mut raw = match read_next_raw_segment(
703                &mut self.reader,
704                &mut self.ssa,
705                &mut scanned,
706                &mut slow_offset,
707                self.config.max_segment_bytes,
708                &mut self.delimiters_settled,
709                self.config.service_string_advice.is_some(),
710            ) {
711                Ok(Some(r)) => r,
712                Ok(None) => return None,
713                Err(e) => {
714                    self.state = StreamState::Done;
715                    return Some(Err(e));
716                }
717            };
718            self.stream_offset = slow_offset as u64;
719            if scanned {
720                self.state = StreamState::Running;
721            }
722            self.bytes_consumed = self.stream_offset;
723
724            raw.bytes.push(self.ssa.segment_term);
725            // `for_segment`, and with the configured `max_segment_bytes`: this
726            // buffer holds one bare segment, so neither the UNA skip nor the
727            // 64 KiB default of `Tokenizer::new` applies.
728            let tok = Tokenizer::for_segment(
729                raw.bytes.as_slice(),
730                self.ssa,
731                self.config.max_segment_bytes,
732            );
733            let mut parser_iter = Parser::new(tok);
734            match parser_iter.next() {
735                Some(Ok(s)) => {
736                    let seg = OwnedSegment::from(s).offset(raw.start_offset);
737                    if let Some(error) = self.check_limits(&seg.tag) {
738                        self.state = StreamState::Done;
739                        return Some(Err(error));
740                    }
741                    if !self.delimiters_settled && seg.tag == "UNB" {
742                        self.adopt_syntax_version(&seg);
743                    }
744                    self.account(&seg.tag);
745                    return Some(Ok(seg));
746                }
747                Some(Err(e)) => {
748                    self.state = StreamState::Done;
749                    return Some(Err(e));
750                }
751                None => {} // Empty segment — loop back.
752            }
753        }
754    }
755}
756
757/// Parse EDIFACT from a buffered reader as a streaming iterator.
758pub fn from_bufread_stream<R: BufRead>(reader: R) -> OwnedSegmentStream<R> {
759    OwnedSegmentStream::new(reader)
760}
761
762/// Parse EDIFACT from a buffered reader as a streaming iterator with custom config.
763pub fn from_bufread_stream_with_config<R: BufRead>(
764    reader: R,
765    config: ReaderConfig,
766) -> OwnedSegmentStream<R> {
767    OwnedSegmentStream::with_config(reader, config)
768}
769
770/// Parse EDIFACT from an arbitrary reader as a streaming iterator.
771pub fn from_reader_stream<R: Read>(reader: R) -> OwnedSegmentStream<BufReader<R>> {
772    from_bufread_stream(BufReader::new(reader))
773}
774
775/// Parse EDIFACT from an arbitrary reader as a streaming iterator with custom config.
776///
777/// # Example
778/// ```
779/// use edifact_rs::{ReaderConfig, from_reader_with_config};
780///
781/// let cfg = ReaderConfig::default().max_segment_bytes(4_096);
782/// let segs: Vec<_> = from_reader_with_config(b"BGM+220+1+9'".as_ref(), cfg)
783///     .collect::<Result<_, _>>()
784///     .unwrap();
785/// assert_eq!(segs[0].tag, "BGM");
786/// ```
787pub fn from_reader_with_config<R: Read>(
788    reader: R,
789    config: ReaderConfig,
790) -> OwnedSegmentStream<BufReader<R>> {
791    from_bufread_stream_with_config(BufReader::new(reader), config)
792}
793
794#[allow(clippy::too_many_arguments)]
795fn read_next_raw_segment<R: BufRead>(
796    reader: &mut R,
797    ssa: &mut crate::tokenizer::ServiceStringAdvice,
798    scanned_header: &mut bool,
799    stream_offset: &mut usize,
800    max_segment_bytes: usize,
801    delimiters_settled: &mut bool,
802    advice_overridden: bool,
803) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
804    loop {
805        let Some((first_offset, first)) = read_next_non_ws_byte(reader, stream_offset)? else {
806            return Ok(None);
807        };
808
809        if !*scanned_header && first == b'U' {
810            let second = read_required_byte(reader, stream_offset)?;
811            let third = read_required_byte(reader, stream_offset)?;
812            if second == b'N' && third == b'A' {
813                let mut una = [0u8; 9];
814                una[0] = b'U';
815                una[1] = b'N';
816                una[2] = b'A';
817                for slot in una.iter_mut().skip(3) {
818                    *slot = read_required_byte(reader, stream_offset)?;
819                }
820                let declared = crate::tokenizer::ServiceStringAdvice {
821                    component_sep: una[3],
822                    element_sep: una[4],
823                    decimal_mark: una[5],
824                    release_char: una[6],
825                    repetition_sep: una[7],
826                    segment_term: una[8],
827                };
828                if !declared.is_valid() {
829                    return Err(EdifactError::InvalidUna);
830                }
831                // A malformed UNA is still an error when the caller supplied its
832                // own advice — the input is broken either way — but the caller's
833                // choice, not the header's, is what parsing then uses.
834                if !advice_overridden {
835                    *ssa = declared;
836                }
837                *delimiters_settled = true;
838                *scanned_header = true;
839                continue;
840            }
841
842            *scanned_header = true;
843            return read_remainder_of_segment(
844                reader,
845                ssa,
846                crate::tokenizer::RawSegment {
847                    bytes: vec![first, second, third],
848                    start_offset: first_offset,
849                },
850                stream_offset,
851                max_segment_bytes,
852            );
853        }
854
855        *scanned_header = true;
856        return read_remainder_of_segment(
857            reader,
858            ssa,
859            crate::tokenizer::RawSegment {
860                bytes: vec![first],
861                start_offset: first_offset,
862            },
863            stream_offset,
864            max_segment_bytes,
865        );
866    }
867}
868
869fn read_remainder_of_segment<R: BufRead>(
870    reader: &mut R,
871    ssa: &crate::tokenizer::ServiceStringAdvice,
872    mut out: crate::tokenizer::RawSegment,
873    stream_offset: &mut usize,
874    max_segment_bytes: usize,
875) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
876    let mut escaped = false;
877    loop {
878        // Strictly greater-than, matching the fast path (`pos > max_segment_bytes`)
879        // and the tokenizer guard.  Using `>=` here made a segment of exactly
880        // `max_segment_bytes` bytes parse on the fast path but fail on the slow
881        // path, so identical input succeeded or failed depending only on whether
882        // it happened to straddle a read-buffer boundary.
883        if out.bytes.len() > max_segment_bytes {
884            return Err(EdifactError::SegmentTooLong {
885                offset: out.start_offset,
886                limit: max_segment_bytes,
887            });
888        }
889        let Some(byte) = read_next_byte(reader, stream_offset)? else {
890            return if out.bytes.is_empty() {
891                Ok(None)
892            } else if escaped {
893                Err(EdifactError::InvalidReleaseSequence {
894                    offset: out.start_offset + out.bytes.len().saturating_sub(1),
895                })
896            } else {
897                Err(EdifactError::UnexpectedEof {
898                    offset: out.start_offset + out.bytes.len(),
899                })
900            };
901        };
902
903        if !escaped && byte == ssa.segment_term {
904            return Ok(Some(out));
905        }
906
907        if !escaped && byte == ssa.release_char {
908            escaped = true;
909            out.bytes.push(byte);
910            continue;
911        }
912
913        escaped = false;
914        out.bytes.push(byte);
915    }
916}
917
918fn read_next_byte<R: BufRead>(
919    reader: &mut R,
920    stream_offset: &mut usize,
921) -> Result<Option<u8>, EdifactError> {
922    let buf = reader.fill_buf()?;
923    if buf.is_empty() {
924        return Ok(None);
925    }
926
927    let byte = buf[0];
928    reader.consume(1);
929    // Saturating add: on 32-bit targets `stream_offset` is a `usize` clamped from a
930    // `u64` field.  Plain `+= 1` would wrap to 0 once the counter reaches `usize::MAX`
931    // and corrupt subsequent span diagnostics / `bytes_consumed` accounting.
932    // Using an explicit local avoids relying on `&mut` auto-deref evaluation order.
933    let next_offset = stream_offset.saturating_add(1);
934    *stream_offset = next_offset;
935    Ok(Some(byte))
936}
937
938fn read_required_byte<R: BufRead>(
939    reader: &mut R,
940    stream_offset: &mut usize,
941) -> Result<u8, EdifactError> {
942    read_next_byte(reader, stream_offset)?.ok_or(EdifactError::UnexpectedEof {
943        offset: *stream_offset,
944    })
945}
946
947fn read_next_non_ws_byte<R: BufRead>(
948    reader: &mut R,
949    stream_offset: &mut usize,
950) -> Result<Option<(usize, u8)>, EdifactError> {
951    loop {
952        let current_offset = *stream_offset;
953        let Some(byte) = read_next_byte(reader, stream_offset)? else {
954            return Ok(None);
955        };
956        if !matches!(byte, b' ' | b'\t' | b'\r' | b'\n') {
957            return Ok(Some((current_offset, byte)));
958        }
959    }
960}
961
962#[cfg(test)]
963mod tests {
964    use super::*;
965    use crate::tokenizer::ServiceStringAdvice;
966
967    fn parse_all(input: &[u8]) -> Vec<Segment<'_>> {
968        let ssa = ServiceStringAdvice::from_bytes_unchecked(input);
969        let tok = Tokenizer::new(input, ssa);
970        Parser::new(tok)
971            .collect::<Result<Vec<_>, _>>()
972            .expect("parse failed")
973    }
974
975    #[test]
976    fn parses_unb_unz() {
977        let input = b"UNB+UNOA:1+SENDER+RECEIVER+200101:0900+1'UNZ+0+1'";
978        let segs = parse_all(input);
979        assert_eq!(segs.len(), 2);
980        assert_eq!(segs[0].tag, "UNB");
981        assert_eq!(segs[1].tag, "UNZ");
982        assert_eq!(segs[0].tag_span, Span::new(0, 3));
983        assert_eq!(segs[0].span, Span::new(0, 41));
984    }
985
986    #[test]
987    fn element_access() {
988        let input = b"BGM+220+ORDER123+9'";
989        let segs = parse_all(input);
990        assert_eq!(segs[0].element_str(0), Some("220"));
991        assert_eq!(segs[0].element_str(1), Some("ORDER123"));
992    }
993
994    #[test]
995    fn component_access() {
996        let input = b"DTM+137:20200101:102'";
997        let segs = parse_all(input);
998        let dtm = &segs[0];
999        assert_eq!(dtm.get_element(0).unwrap().get_component(0), Some("137"));
1000        assert_eq!(
1001            dtm.get_element(0).unwrap().get_component(1),
1002            Some("20200101")
1003        );
1004        assert_eq!(dtm.get_element(0).unwrap().get_component(2), Some("102"));
1005    }
1006
1007    #[test]
1008    fn release_char_resolved() {
1009        let input = b"FTX+AAA++test?+value'";
1010        let segs = parse_all(input);
1011        assert_eq!(segs[0].element_str(2), Some("test+value"));
1012        assert_eq!(
1013            segs[0].get_element(2).unwrap().component_span(0),
1014            Some(Span::new(9, 20))
1015        );
1016    }
1017
1018    #[test]
1019    fn reader_path_preserves_custom_una_delimiters() {
1020        let input = b"UNA:;.? 'BGM;220;test?;value'";
1021        let segments = super::from_bufread(std::io::BufReader::new(std::io::Cursor::new(input)))
1022            .expect("reader parse should succeed");
1023        let bgm = segments
1024            .iter()
1025            .find(|segment| segment.tag == "BGM")
1026            .expect("BGM segment should be present");
1027        assert_eq!(bgm.elements[0].components[0].0, "220");
1028        assert_eq!(bgm.elements[1].components[0].0, "test;value");
1029    }
1030
1031    #[test]
1032    fn arbitrary_bytes_no_panic() {
1033        // This is the stable no-panic property — arbitrary input must not panic
1034        let garbage: &[u8] = b"\xff\x00\x01\x02ABC+++'''???";
1035        let _ = crate::from_bytes(garbage).collect::<Vec<_>>();
1036    }
1037
1038    #[test]
1039    fn from_reader_handles_chunk_boundaries() {
1040        let input = b"UNA:+.? 'BGM+220+test?+value'UNT+2+1'";
1041        let reader = std::io::BufReader::with_capacity(5, std::io::Cursor::new(input));
1042        let parsed = from_bufread(reader).expect("reader parsing should succeed");
1043        assert_eq!(parsed.len(), 2);
1044        assert_eq!(parsed[0].tag, "BGM");
1045        assert_eq!(parsed[0].elements[1].components[0].0, "test+value");
1046        assert_eq!(parsed[1].tag, "UNT");
1047    }
1048
1049    #[test]
1050    fn from_reader_without_una_uses_default_delimiters() {
1051        let input = b"BGM+220+X'UNT+2+1'";
1052        let parsed =
1053            from_reader(std::io::Cursor::new(input)).expect("reader parsing should succeed");
1054        assert_eq!(parsed.len(), 2);
1055        assert_eq!(parsed[0].tag, "BGM");
1056        assert_eq!(parsed[0].elements[0].components[0].0, "220");
1057        assert_eq!(parsed[1].span, Span::new(10, 18));
1058    }
1059
1060    #[test]
1061    fn dangling_release_sequence_is_error() {
1062        let input = b"FTX+AAA++dangling?";
1063        let err = crate::from_bytes(input)
1064            .collect::<Result<Vec<_>, _>>()
1065            .expect_err("expected dangling release to fail");
1066
1067        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
1068    }
1069
1070    #[test]
1071    fn from_reader_reports_dangling_release_sequence() {
1072        let input = b"FTX+AAA++dangling?";
1073        let err = from_reader(std::io::Cursor::new(input))
1074            .expect_err("expected dangling release from reader path");
1075        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
1076    }
1077
1078    #[test]
1079    fn a_segment_tagged_una_parses_the_same_on_both_paths() {
1080        // `UNA` is three ASCII uppercase letters, so it is a syntactically legal
1081        // segment tag.  The reader re-tokenizes each segment from its own slice,
1082        // where the whole-interchange "skip nine bytes of service string advice"
1083        // heuristic is wrong: it ate the tag and everything after it, and the
1084        // identical bytes that parsed cleanly through `from_bytes` came back as
1085        // `InvalidSegmentTag` through a reader.
1086        let input = b"BGM+220'UNA+XXXXXX'BGM+221'";
1087
1088        let from_slice: Vec<_> = crate::from_bytes(input)
1089            .collect::<Result<Vec<_>, _>>()
1090            .expect("slice path");
1091        let from_reader =
1092            from_reader(std::io::Cursor::new(&input[..])).expect("reader path must agree");
1093
1094        assert_eq!(
1095            from_slice.iter().map(|s| s.tag).collect::<Vec<_>>(),
1096            from_reader
1097                .iter()
1098                .map(|s| s.tag.as_str())
1099                .collect::<Vec<_>>(),
1100        );
1101        assert_eq!(from_reader[1].element_str(0), Some("XXXXXX"));
1102    }
1103
1104    #[test]
1105    fn the_reader_adopts_the_repetition_separator_from_the_unb_syntax_version() {
1106        // The slice path resolves this before tokenizing; the reader can only
1107        // learn it once the UNB has been parsed.  Both must agree.
1108        let input = b"UNB+UNOC:4+S+R+260101:0900+IC1'RFF+ON:1*ON:2'UNZ+0+IC1'";
1109
1110        let from_slice: Vec<_> = crate::from_bytes(input)
1111            .collect::<Result<Vec<_>, _>>()
1112            .expect("slice path");
1113        let from_reader = from_reader(std::io::Cursor::new(&input[..])).expect("reader path");
1114
1115        assert_eq!(from_slice[1].get_element(0).unwrap().repeat_count(), 2);
1116        assert_eq!(from_reader[1].elements[0].repeat_count(), 2);
1117    }
1118
1119    #[test]
1120    fn an_explicit_service_string_advice_parses_a_fragment_with_no_header() {
1121        // A message lifted out of an interchange carries neither UNA nor UNB, so
1122        // nothing in it records the delimiters — the caller has to supply them.
1123        let ssa = ServiceStringAdvice::from_bytes(b"UNA:;.? ~").expect("UNA");
1124        let config = ReaderConfig::default().with_service_string_advice(ssa);
1125
1126        let from_slice: Vec<_> = crate::from_bytes_with_config(b"BGM;220;PO-4711~", config)
1127            .collect::<Result<Vec<_>, _>>()
1128            .expect("slice path");
1129        assert_eq!(from_slice[0].element_str(1), Some("PO-4711"));
1130
1131        let from_reader: Vec<_> =
1132            from_reader_with_config(std::io::Cursor::new(b"BGM;220;PO-4711~"), config)
1133                .collect::<Result<Vec<_>, _>>()
1134                .expect("reader path");
1135        assert_eq!(from_reader[0].element_str(1), Some("PO-4711"));
1136    }
1137
1138    #[test]
1139    fn an_explicit_service_string_advice_outranks_the_una_in_the_input() {
1140        let ssa = ServiceStringAdvice::default();
1141        let config = ReaderConfig::default().with_service_string_advice(ssa);
1142        // The UNA declares `;`, the caller insists on `+`.
1143        let input = b"UNA:;.? 'BGM+220'";
1144
1145        for segments in [
1146            crate::from_bytes_with_config(input, config)
1147                .map(|r| r.map(crate::OwnedSegment::from))
1148                .collect::<Result<Vec<_>, _>>()
1149                .expect("slice path"),
1150            from_reader_with_config(std::io::Cursor::new(&input[..]), config)
1151                .collect::<Result<Vec<_>, _>>()
1152                .expect("reader path"),
1153        ] {
1154            assert_eq!(segments[0].element_str(0).unwrap(), "220");
1155        }
1156    }
1157
1158    #[test]
1159    fn from_reader_rejects_invalid_una() {
1160        let input = b"UNA::.? 'BGM:220'";
1161        let err = from_reader(std::io::Cursor::new(input))
1162            .expect_err("invalid UNA should fail reader parsing");
1163        assert!(matches!(err, EdifactError::InvalidUna));
1164    }
1165}