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}
310
311impl Default for ReaderConfig {
312    fn default() -> Self {
313        Self {
314            max_segment_bytes: 65_536,
315            max_segments: None,
316            max_input_bytes: None,
317            max_messages: None,
318        }
319    }
320}
321
322impl ReaderConfig {
323    /// Set the maximum segment byte length and return `self`.
324    #[must_use]
325    pub fn max_segment_bytes(mut self, limit: usize) -> Self {
326        self.max_segment_bytes = limit;
327        self
328    }
329
330    /// Set the maximum number of segments to yield and return `self`.
331    #[must_use]
332    pub fn max_segments(mut self, limit: usize) -> Self {
333        self.max_segments = Some(limit);
334        self
335    }
336
337    /// Set the maximum total input bytes to consume and return `self`.
338    #[must_use]
339    pub fn max_input_bytes(mut self, limit: u64) -> Self {
340        self.max_input_bytes = Some(limit);
341        self
342    }
343
344    /// Set the maximum number of EDIFACT messages (UNH/UNT pairs) to yield and return `self`.
345    #[must_use]
346    pub fn max_messages(mut self, limit: usize) -> Self {
347        self.max_messages = Some(limit);
348        self
349    }
350}
351
352/// Streaming state for [`OwnedSegmentStream`].
353#[derive(Debug, Clone, Copy, PartialEq, Eq)]
354enum StreamState {
355    /// UNA header not yet scanned; must inspect first bytes.
356    Init,
357    /// UNA has been scanned (or was absent); streaming segments.
358    Running,
359    /// A terminal error was encountered; no more items.
360    Done,
361}
362
363/// Streaming iterator over owned segments from a buffered reader.
364///
365/// # Performance
366///
367/// A **fast path** uses [`BufRead::fill_buf`] + `memchr` to locate the segment
368/// terminator within the OS-level read buffer (typically 8 KB) without any
369/// intermediate heap allocation.  The segment bytes are parsed directly from
370/// the buffer slice and converted to an [`OwnedSegment`] in a single pass.
371///
372/// For segments that span read-buffer boundaries the implementation falls back
373/// to the byte-accumulation slow path, which allocates a temporary `Vec<u8>`
374/// and re-tokenizes — the same behaviour as in older versions of the library.
375/// In practice this fallback is rare because the default `BufReader` buffer
376/// (8 KB) is far larger than a typical EDIFACT segment (<300 bytes).
377///
378/// Configure limits via [`ReaderConfig`] and [`from_reader_with_config`] /
379/// [`from_bufread_stream_with_config`].
380pub struct OwnedSegmentStream<R: BufRead> {
381    reader: R,
382    ssa: crate::tokenizer::ServiceStringAdvice,
383    state: StreamState,
384    stream_offset: u64,
385    config: ReaderConfig,
386    /// Number of segments successfully yielded so far.
387    segments_yielded: usize,
388    /// Number of complete EDIFACT messages (UNH/UNT pairs) seen so far.
389    messages_yielded: usize,
390    /// Whether the last segment tag was UNH (inside a message).
391    in_message: bool,
392    /// Total bytes consumed from the reader (UNA header + segment data).
393    bytes_consumed: u64,
394}
395
396impl<R: BufRead> OwnedSegmentStream<R> {
397    fn new(reader: R) -> Self {
398        Self::with_config(reader, ReaderConfig::default())
399    }
400
401    fn with_config(reader: R, config: ReaderConfig) -> Self {
402        Self {
403            reader,
404            ssa: crate::tokenizer::ServiceStringAdvice::default(),
405            state: StreamState::Init,
406            stream_offset: 0,
407            config,
408            segments_yielded: 0,
409            messages_yielded: 0,
410            in_message: false,
411            bytes_consumed: 0,
412        }
413    }
414
415    /// Check the whole-input budgets against a segment that has just been read.
416    ///
417    /// Checking *after* the read rather than before is what makes an input that
418    /// ends exactly at a limit finish cleanly: a budget is only exceeded when
419    /// there was genuinely more data than the caller allowed.
420    fn check_limits(&self, tag: &str) -> Option<EdifactError> {
421        if let Some(max) = self.config.max_segments {
422            if self.segments_yielded >= max {
423                return Some(EdifactError::LimitExceeded {
424                    limit: "max_segments",
425                    max: max as u64,
426                });
427            }
428        }
429        // Only a `UNH` opens a message, so only a `UNH` can push the count past
430        // the budget.  Testing every segment would trip on the interchange
431        // trailer, which belongs to no message.
432        if let Some(max) = self.config.max_messages {
433            if tag == "UNH" && self.messages_yielded >= max {
434                return Some(EdifactError::LimitExceeded {
435                    limit: "max_messages",
436                    max: max as u64,
437                });
438            }
439        }
440        if let Some(max) = self.config.max_input_bytes {
441            if self.bytes_consumed > max {
442                return Some(EdifactError::LimitExceeded {
443                    limit: "max_input_bytes",
444                    max,
445                });
446            }
447        }
448        None
449    }
450
451    /// Record a yielded segment against the segment and message counters.
452    ///
453    /// Only a `UNT` that closes a `UNH` already seen advances the message
454    /// counter; a bare `UNT` is malformed and must not inflate it.
455    fn account(&mut self, tag: &str) {
456        self.segments_yielded += 1;
457        if tag == "UNT" {
458            if self.in_message {
459                self.messages_yielded += 1;
460            }
461            self.in_message = false;
462        } else if tag == "UNH" {
463            self.in_message = true;
464        }
465    }
466}
467
468// ── fast-path helpers ─────────────────────────────────────────────────────────
469
470/// Outcome of a single-buffer segment extraction attempt.
471enum FastSegment {
472    /// Segment parsed; second field = bytes to consume (content + terminator).
473    Parsed(OwnedSegment, usize),
474    /// Only whitespace or an isolated terminator; bytes to skip and continue.
475    Skip(usize),
476    /// Terminator not present in the current buffer; caller must use slow path.
477    NeedMore,
478    /// Buffer is empty — no more input.
479    Eof,
480    /// Parse error.
481    Err(EdifactError),
482}
483
484/// Return the byte offset of the first **unescaped** occurrence of `term` in `buf`.
485///
486/// A byte is *escaped* when it is immediately preceded by `release` (e.g.
487/// `?'` escapes `'`).  Two consecutive release chars cancel each other, so
488/// `??'` contains an *unescaped* `'`.
489///
490/// # Complexity
491///
492/// O(n) in the length of `buf`.  `memchr2` is used to fast-scan past bytes
493/// that are neither `release` nor `term`, so SIMD acceleration applies on
494/// platforms where `memchr` provides it.
495fn find_unescaped_term(buf: &[u8], term: u8, release: u8) -> Option<usize> {
496    let mut i = 0;
497    while i < buf.len() {
498        // Fast-skip to the next byte that might be a release char or terminator.
499        let rel = memchr2(release, term, &buf[i..])?;
500        let pos = i + rel;
501        if buf[pos] == release {
502            // Release char: the next byte is escaped — skip both.
503            i = pos + 2;
504        } else {
505            // Unescaped terminator found.
506            return Some(pos);
507        }
508    }
509    None
510}
511
512/// Try to parse one segment directly from the `BufRead` buffer.
513///
514/// This function borrows `reader` only for the duration of the call.  After it
515/// returns the caller is free to call `reader.consume(n)`.
516fn try_fast_segment<R: BufRead>(
517    reader: &mut R,
518    ssa: crate::tokenizer::ServiceStringAdvice,
519    seg_start: usize,
520    max_segment_bytes: usize,
521) -> FastSegment {
522    let buf = match reader.fill_buf() {
523        Ok(b) => b,
524        Err(e) => return FastSegment::Err(e.into()),
525    };
526
527    if buf.is_empty() {
528        return FastSegment::Eof;
529    }
530
531    let Some(pos) = find_unescaped_term(buf, ssa.segment_term, ssa.release_char) else {
532        return FastSegment::NeedMore;
533    };
534
535    // Enforce the segment-size guard *before* any allocation.
536    // `pos` is the index of the terminator byte, so the segment body is `buf[..pos]`.
537    if pos > max_segment_bytes {
538        return FastSegment::Err(EdifactError::SegmentTooLong {
539            offset: seg_start,
540            limit: max_segment_bytes,
541        });
542    }
543
544    // `buf[..pos]` is the segment content without the terminator.
545    let seg_bytes = &buf[..pos];
546
547    // Skip isolated terminators / pure-whitespace slots between segments.
548    if seg_bytes
549        .iter()
550        .all(|&b| matches!(b, b' ' | b'\t' | b'\r' | b'\n'))
551    {
552        return FastSegment::Skip(pos + 1);
553    }
554
555    // Parse directly from the buffer slice — zero intermediate allocation.
556    // Include the terminator byte so the parser sees a `SegmentTerminator`
557    // token and records a span that is consistent with the `from_bytes` path.
558    // Use `with_limit(max_segment_bytes)` so the tokenizer respects the caller's
559    // configured limit; `Tokenizer::new` would impose a hard 64 KiB cap that
560    // could reject segments already allowed by a larger `max_segment_bytes`.
561    let tok = Tokenizer::with_limit(&buf[..pos + 1], ssa, max_segment_bytes);
562    let mut parser_iter = Parser::new(tok);
563    match parser_iter.next() {
564        None => FastSegment::Skip(pos + 1),
565        Some(Err(e)) => FastSegment::Err(e),
566        Some(Ok(s)) => FastSegment::Parsed(OwnedSegment::from(s).offset(seg_start), pos + 1),
567    }
568    // `buf` borrow released here — `reader.consume()` is safe to call in the caller.
569}
570
571// ── Iterator impl ─────────────────────────────────────────────────────────────
572
573impl<R: BufRead> Iterator for OwnedSegmentStream<R> {
574    type Item = Result<OwnedSegment, EdifactError>;
575
576    fn next(&mut self) -> Option<Self::Item> {
577        if self.state == StreamState::Done {
578            return None;
579        }
580
581        loop {
582            // ── Fast path (after UNA has been consumed) ───────────────────
583            if self.state == StreamState::Running {
584                let seg_start = self.stream_offset;
585                match try_fast_segment(
586                    &mut self.reader,
587                    self.ssa,
588                    // Saturate rather than wrap on 32-bit targets; Span offsets
589                    // are `usize` so streams > 4 GiB on 32-bit produce clamped
590                    // (but monotonic) diagnostic positions.
591                    seg_start.min(usize::MAX as u64) as usize,
592                    self.config.max_segment_bytes,
593                ) {
594                    FastSegment::Parsed(seg, n) => {
595                        let n = n as u64;
596                        self.reader.consume(n as usize);
597                        self.stream_offset += n;
598                        self.bytes_consumed = self.stream_offset;
599                        if let Some(error) = self.check_limits(&seg.tag) {
600                            self.state = StreamState::Done;
601                            return Some(Err(error));
602                        }
603                        self.account(&seg.tag);
604                        return Some(Ok(seg));
605                    }
606                    FastSegment::Skip(n) => {
607                        let n = n as u64;
608                        self.reader.consume(n as usize);
609                        self.stream_offset += n;
610                        self.bytes_consumed = self.stream_offset;
611                        continue;
612                    }
613                    FastSegment::Eof => return None,
614                    FastSegment::Err(e) => {
615                        self.state = StreamState::Done;
616                        return Some(Err(e));
617                    }
618                    FastSegment::NeedMore => {
619                        // Segment spans buffer boundary — fall through to slow path.
620                    }
621                }
622            }
623
624            // ── Slow path: byte accumulation (also handles UNA header) ────
625            let mut scanned = self.state != StreamState::Init;
626            // `read_next_raw_segment` tracks offset as `usize` for segment
627            // start positions; sync back to the `u64` field afterward.
628            // Saturate rather than wrap on 32-bit targets (same rationale as
629            // the fast-path cast above).
630            let mut slow_offset: usize = self.stream_offset.min(usize::MAX as u64) as usize;
631            let mut raw = match read_next_raw_segment(
632                &mut self.reader,
633                &mut self.ssa,
634                &mut scanned,
635                &mut slow_offset,
636                self.config.max_segment_bytes,
637            ) {
638                Ok(Some(r)) => r,
639                Ok(None) => return None,
640                Err(e) => {
641                    self.state = StreamState::Done;
642                    return Some(Err(e));
643                }
644            };
645            self.stream_offset = slow_offset as u64;
646            if scanned {
647                self.state = StreamState::Running;
648            }
649            self.bytes_consumed = self.stream_offset;
650
651            raw.bytes.push(self.ssa.segment_term);
652            // Use `with_limit` so the configured max_segment_bytes is honoured on the
653            // slow path as well; `Tokenizer::new` would impose a hard 64 KiB cap that
654            // could reject segments the caller explicitly permitted via ReaderConfig.
655            let tok = Tokenizer::with_limit(
656                raw.bytes.as_slice(),
657                self.ssa,
658                self.config.max_segment_bytes,
659            );
660            let mut parser_iter = Parser::new(tok);
661            match parser_iter.next() {
662                Some(Ok(s)) => {
663                    let seg = OwnedSegment::from(s).offset(raw.start_offset);
664                    if let Some(error) = self.check_limits(&seg.tag) {
665                        self.state = StreamState::Done;
666                        return Some(Err(error));
667                    }
668                    self.account(&seg.tag);
669                    return Some(Ok(seg));
670                }
671                Some(Err(e)) => {
672                    self.state = StreamState::Done;
673                    return Some(Err(e));
674                }
675                None => {} // Empty segment — loop back.
676            }
677        }
678    }
679}
680
681/// Parse EDIFACT from a buffered reader as a streaming iterator.
682pub fn from_bufread_stream<R: BufRead>(reader: R) -> OwnedSegmentStream<R> {
683    OwnedSegmentStream::new(reader)
684}
685
686/// Parse EDIFACT from a buffered reader as a streaming iterator with custom config.
687pub fn from_bufread_stream_with_config<R: BufRead>(
688    reader: R,
689    config: ReaderConfig,
690) -> OwnedSegmentStream<R> {
691    OwnedSegmentStream::with_config(reader, config)
692}
693
694/// Parse EDIFACT from an arbitrary reader as a streaming iterator.
695pub fn from_reader_stream<R: Read>(reader: R) -> OwnedSegmentStream<BufReader<R>> {
696    from_bufread_stream(BufReader::new(reader))
697}
698
699/// Parse EDIFACT from an arbitrary reader as a streaming iterator with custom config.
700///
701/// # Example
702/// ```
703/// use edifact_rs::{ReaderConfig, from_reader_with_config};
704///
705/// let cfg = ReaderConfig::default().max_segment_bytes(4_096);
706/// let segs: Vec<_> = from_reader_with_config(b"BGM+220+1+9'".as_ref(), cfg)
707///     .collect::<Result<_, _>>()
708///     .unwrap();
709/// assert_eq!(segs[0].tag, "BGM");
710/// ```
711pub fn from_reader_with_config<R: Read>(
712    reader: R,
713    config: ReaderConfig,
714) -> OwnedSegmentStream<BufReader<R>> {
715    from_bufread_stream_with_config(BufReader::new(reader), config)
716}
717
718fn read_next_raw_segment<R: BufRead>(
719    reader: &mut R,
720    ssa: &mut crate::tokenizer::ServiceStringAdvice,
721    scanned_header: &mut bool,
722    stream_offset: &mut usize,
723    max_segment_bytes: usize,
724) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
725    loop {
726        let Some((first_offset, first)) = read_next_non_ws_byte(reader, stream_offset)? else {
727            return Ok(None);
728        };
729
730        if !*scanned_header && first == b'U' {
731            let second = read_required_byte(reader, stream_offset)?;
732            let third = read_required_byte(reader, stream_offset)?;
733            if second == b'N' && third == b'A' {
734                let mut una = [0u8; 9];
735                una[0] = b'U';
736                una[1] = b'N';
737                una[2] = b'A';
738                for slot in una.iter_mut().skip(3) {
739                    *slot = read_required_byte(reader, stream_offset)?;
740                }
741                *ssa = crate::tokenizer::ServiceStringAdvice {
742                    component_sep: una[3],
743                    element_sep: una[4],
744                    decimal_mark: una[5],
745                    release_char: una[6],
746                    repetition_sep: una[7],
747                    segment_term: una[8],
748                };
749                if !ssa.is_valid() {
750                    return Err(EdifactError::InvalidUna);
751                }
752                *scanned_header = true;
753                continue;
754            }
755
756            *scanned_header = true;
757            return read_remainder_of_segment(
758                reader,
759                ssa,
760                crate::tokenizer::RawSegment {
761                    bytes: vec![first, second, third],
762                    start_offset: first_offset,
763                },
764                stream_offset,
765                max_segment_bytes,
766            );
767        }
768
769        *scanned_header = true;
770        return read_remainder_of_segment(
771            reader,
772            ssa,
773            crate::tokenizer::RawSegment {
774                bytes: vec![first],
775                start_offset: first_offset,
776            },
777            stream_offset,
778            max_segment_bytes,
779        );
780    }
781}
782
783fn read_remainder_of_segment<R: BufRead>(
784    reader: &mut R,
785    ssa: &crate::tokenizer::ServiceStringAdvice,
786    mut out: crate::tokenizer::RawSegment,
787    stream_offset: &mut usize,
788    max_segment_bytes: usize,
789) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
790    let mut escaped = false;
791    loop {
792        // Strictly greater-than, matching the fast path (`pos > max_segment_bytes`)
793        // and the tokenizer guard.  Using `>=` here made a segment of exactly
794        // `max_segment_bytes` bytes parse on the fast path but fail on the slow
795        // path, so identical input succeeded or failed depending only on whether
796        // it happened to straddle a read-buffer boundary.
797        if out.bytes.len() > max_segment_bytes {
798            return Err(EdifactError::SegmentTooLong {
799                offset: out.start_offset,
800                limit: max_segment_bytes,
801            });
802        }
803        let Some(byte) = read_next_byte(reader, stream_offset)? else {
804            return if out.bytes.is_empty() {
805                Ok(None)
806            } else if escaped {
807                Err(EdifactError::InvalidReleaseSequence {
808                    offset: out.start_offset + out.bytes.len().saturating_sub(1),
809                })
810            } else {
811                Err(EdifactError::UnexpectedEof {
812                    offset: out.start_offset + out.bytes.len(),
813                })
814            };
815        };
816
817        if !escaped && byte == ssa.segment_term {
818            return Ok(Some(out));
819        }
820
821        if !escaped && byte == ssa.release_char {
822            escaped = true;
823            out.bytes.push(byte);
824            continue;
825        }
826
827        escaped = false;
828        out.bytes.push(byte);
829    }
830}
831
832fn read_next_byte<R: BufRead>(
833    reader: &mut R,
834    stream_offset: &mut usize,
835) -> Result<Option<u8>, EdifactError> {
836    let buf = reader.fill_buf()?;
837    if buf.is_empty() {
838        return Ok(None);
839    }
840
841    let byte = buf[0];
842    reader.consume(1);
843    // Saturating add: on 32-bit targets `stream_offset` is a `usize` clamped from a
844    // `u64` field.  Plain `+= 1` would wrap to 0 once the counter reaches `usize::MAX`
845    // and corrupt subsequent span diagnostics / `bytes_consumed` accounting.
846    // Using an explicit local avoids relying on `&mut` auto-deref evaluation order.
847    let next_offset = stream_offset.saturating_add(1);
848    *stream_offset = next_offset;
849    Ok(Some(byte))
850}
851
852fn read_required_byte<R: BufRead>(
853    reader: &mut R,
854    stream_offset: &mut usize,
855) -> Result<u8, EdifactError> {
856    read_next_byte(reader, stream_offset)?.ok_or(EdifactError::UnexpectedEof {
857        offset: *stream_offset,
858    })
859}
860
861fn read_next_non_ws_byte<R: BufRead>(
862    reader: &mut R,
863    stream_offset: &mut usize,
864) -> Result<Option<(usize, u8)>, EdifactError> {
865    loop {
866        let current_offset = *stream_offset;
867        let Some(byte) = read_next_byte(reader, stream_offset)? else {
868            return Ok(None);
869        };
870        if !matches!(byte, b' ' | b'\t' | b'\r' | b'\n') {
871            return Ok(Some((current_offset, byte)));
872        }
873    }
874}
875
876#[cfg(test)]
877mod tests {
878    use super::*;
879    use crate::tokenizer::ServiceStringAdvice;
880
881    fn parse_all(input: &[u8]) -> Vec<Segment<'_>> {
882        let ssa = ServiceStringAdvice::from_bytes_unchecked(input);
883        let tok = Tokenizer::new(input, ssa);
884        Parser::new(tok)
885            .collect::<Result<Vec<_>, _>>()
886            .expect("parse failed")
887    }
888
889    #[test]
890    fn parses_unb_unz() {
891        let input = b"UNB+UNOA:1+SENDER+RECEIVER+200101:0900+1'UNZ+0+1'";
892        let segs = parse_all(input);
893        assert_eq!(segs.len(), 2);
894        assert_eq!(segs[0].tag, "UNB");
895        assert_eq!(segs[1].tag, "UNZ");
896        assert_eq!(segs[0].tag_span, Span::new(0, 3));
897        assert_eq!(segs[0].span, Span::new(0, 41));
898    }
899
900    #[test]
901    fn element_access() {
902        let input = b"BGM+220+ORDER123+9'";
903        let segs = parse_all(input);
904        assert_eq!(segs[0].element_str(0), Some("220"));
905        assert_eq!(segs[0].element_str(1), Some("ORDER123"));
906    }
907
908    #[test]
909    fn component_access() {
910        let input = b"DTM+137:20200101:102'";
911        let segs = parse_all(input);
912        let dtm = &segs[0];
913        assert_eq!(dtm.get_element(0).unwrap().get_component(0), Some("137"));
914        assert_eq!(
915            dtm.get_element(0).unwrap().get_component(1),
916            Some("20200101")
917        );
918        assert_eq!(dtm.get_element(0).unwrap().get_component(2), Some("102"));
919    }
920
921    #[test]
922    fn release_char_resolved() {
923        let input = b"FTX+AAA++test?+value'";
924        let segs = parse_all(input);
925        assert_eq!(segs[0].element_str(2), Some("test+value"));
926        assert_eq!(
927            segs[0].get_element(2).unwrap().component_span(0),
928            Some(Span::new(9, 20))
929        );
930    }
931
932    #[test]
933    fn reader_path_preserves_custom_una_delimiters() {
934        let input = b"UNA:;.? 'BGM;220;test?;value'";
935        let segments = super::from_bufread(std::io::BufReader::new(std::io::Cursor::new(input)))
936            .expect("reader parse should succeed");
937        let bgm = segments
938            .iter()
939            .find(|segment| segment.tag == "BGM")
940            .expect("BGM segment should be present");
941        assert_eq!(bgm.elements[0].components[0].0, "220");
942        assert_eq!(bgm.elements[1].components[0].0, "test;value");
943    }
944
945    #[test]
946    fn arbitrary_bytes_no_panic() {
947        // This is the stable no-panic property — arbitrary input must not panic
948        let garbage: &[u8] = b"\xff\x00\x01\x02ABC+++'''???";
949        let _ = crate::from_bytes(garbage).collect::<Vec<_>>();
950    }
951
952    #[test]
953    fn from_reader_handles_chunk_boundaries() {
954        let input = b"UNA:+.? 'BGM+220+test?+value'UNT+2+1'";
955        let reader = std::io::BufReader::with_capacity(5, std::io::Cursor::new(input));
956        let parsed = from_bufread(reader).expect("reader parsing should succeed");
957        assert_eq!(parsed.len(), 2);
958        assert_eq!(parsed[0].tag, "BGM");
959        assert_eq!(parsed[0].elements[1].components[0].0, "test+value");
960        assert_eq!(parsed[1].tag, "UNT");
961    }
962
963    #[test]
964    fn from_reader_without_una_uses_default_delimiters() {
965        let input = b"BGM+220+X'UNT+2+1'";
966        let parsed =
967            from_reader(std::io::Cursor::new(input)).expect("reader parsing should succeed");
968        assert_eq!(parsed.len(), 2);
969        assert_eq!(parsed[0].tag, "BGM");
970        assert_eq!(parsed[0].elements[0].components[0].0, "220");
971        assert_eq!(parsed[1].span, Span::new(10, 18));
972    }
973
974    #[test]
975    fn dangling_release_sequence_is_error() {
976        let input = b"FTX+AAA++dangling?";
977        let err = crate::from_bytes(input)
978            .collect::<Result<Vec<_>, _>>()
979            .expect_err("expected dangling release to fail");
980
981        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
982    }
983
984    #[test]
985    fn from_reader_reports_dangling_release_sequence() {
986        let input = b"FTX+AAA++dangling?";
987        let err = from_reader(std::io::Cursor::new(input))
988            .expect_err("expected dangling release from reader path");
989        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
990    }
991
992    #[test]
993    fn from_reader_rejects_invalid_una() {
994        let input = b"UNA::.? 'BGM:220'";
995        let err = from_reader(std::io::Cursor::new(input))
996            .expect_err("invalid UNA should fail reader parsing");
997        assert!(matches!(err, EdifactError::InvalidUna));
998    }
999}