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    // `for_segment` (not `with_limit`) because this slice is one bare segment:
559    // the whole-interchange constructors would mistake a segment tagged `UNA`
560    // for a service string advice and skip nine bytes of it.  The configured
561    // `max_segment_bytes` is passed through so the tokenizer never applies a
562    // tighter cap than the caller allowed.
563    let tok = Tokenizer::for_segment(&buf[..pos + 1], ssa, max_segment_bytes);
564    let mut parser_iter = Parser::new(tok);
565    match parser_iter.next() {
566        None => FastSegment::Skip(pos + 1),
567        Some(Err(e)) => FastSegment::Err(e),
568        Some(Ok(s)) => FastSegment::Parsed(OwnedSegment::from(s).offset(seg_start), pos + 1),
569    }
570    // `buf` borrow released here — `reader.consume()` is safe to call in the caller.
571}
572
573// ── Iterator impl ─────────────────────────────────────────────────────────────
574
575impl<R: BufRead> Iterator for OwnedSegmentStream<R> {
576    type Item = Result<OwnedSegment, EdifactError>;
577
578    fn next(&mut self) -> Option<Self::Item> {
579        if self.state == StreamState::Done {
580            return None;
581        }
582
583        loop {
584            // ── Fast path (after UNA has been consumed) ───────────────────
585            if self.state == StreamState::Running {
586                let seg_start = self.stream_offset;
587                match try_fast_segment(
588                    &mut self.reader,
589                    self.ssa,
590                    // Saturate rather than wrap on 32-bit targets; Span offsets
591                    // are `usize` so streams > 4 GiB on 32-bit produce clamped
592                    // (but monotonic) diagnostic positions.
593                    seg_start.min(usize::MAX as u64) as usize,
594                    self.config.max_segment_bytes,
595                ) {
596                    FastSegment::Parsed(seg, n) => {
597                        let n = n as u64;
598                        self.reader.consume(n as usize);
599                        self.stream_offset += n;
600                        self.bytes_consumed = self.stream_offset;
601                        if let Some(error) = self.check_limits(&seg.tag) {
602                            self.state = StreamState::Done;
603                            return Some(Err(error));
604                        }
605                        self.account(&seg.tag);
606                        return Some(Ok(seg));
607                    }
608                    FastSegment::Skip(n) => {
609                        let n = n as u64;
610                        self.reader.consume(n as usize);
611                        self.stream_offset += n;
612                        self.bytes_consumed = self.stream_offset;
613                        continue;
614                    }
615                    FastSegment::Eof => return None,
616                    FastSegment::Err(e) => {
617                        self.state = StreamState::Done;
618                        return Some(Err(e));
619                    }
620                    FastSegment::NeedMore => {
621                        // Segment spans buffer boundary — fall through to slow path.
622                    }
623                }
624            }
625
626            // ── Slow path: byte accumulation (also handles UNA header) ────
627            let mut scanned = self.state != StreamState::Init;
628            // `read_next_raw_segment` tracks offset as `usize` for segment
629            // start positions; sync back to the `u64` field afterward.
630            // Saturate rather than wrap on 32-bit targets (same rationale as
631            // the fast-path cast above).
632            let mut slow_offset: usize = self.stream_offset.min(usize::MAX as u64) as usize;
633            let mut raw = match read_next_raw_segment(
634                &mut self.reader,
635                &mut self.ssa,
636                &mut scanned,
637                &mut slow_offset,
638                self.config.max_segment_bytes,
639            ) {
640                Ok(Some(r)) => r,
641                Ok(None) => return None,
642                Err(e) => {
643                    self.state = StreamState::Done;
644                    return Some(Err(e));
645                }
646            };
647            self.stream_offset = slow_offset as u64;
648            if scanned {
649                self.state = StreamState::Running;
650            }
651            self.bytes_consumed = self.stream_offset;
652
653            raw.bytes.push(self.ssa.segment_term);
654            // `for_segment`, and with the configured `max_segment_bytes`: this
655            // buffer holds one bare segment, so neither the UNA skip nor the
656            // 64 KiB default of `Tokenizer::new` applies.
657            let tok = Tokenizer::for_segment(
658                raw.bytes.as_slice(),
659                self.ssa,
660                self.config.max_segment_bytes,
661            );
662            let mut parser_iter = Parser::new(tok);
663            match parser_iter.next() {
664                Some(Ok(s)) => {
665                    let seg = OwnedSegment::from(s).offset(raw.start_offset);
666                    if let Some(error) = self.check_limits(&seg.tag) {
667                        self.state = StreamState::Done;
668                        return Some(Err(error));
669                    }
670                    self.account(&seg.tag);
671                    return Some(Ok(seg));
672                }
673                Some(Err(e)) => {
674                    self.state = StreamState::Done;
675                    return Some(Err(e));
676                }
677                None => {} // Empty segment — loop back.
678            }
679        }
680    }
681}
682
683/// Parse EDIFACT from a buffered reader as a streaming iterator.
684pub fn from_bufread_stream<R: BufRead>(reader: R) -> OwnedSegmentStream<R> {
685    OwnedSegmentStream::new(reader)
686}
687
688/// Parse EDIFACT from a buffered reader as a streaming iterator with custom config.
689pub fn from_bufread_stream_with_config<R: BufRead>(
690    reader: R,
691    config: ReaderConfig,
692) -> OwnedSegmentStream<R> {
693    OwnedSegmentStream::with_config(reader, config)
694}
695
696/// Parse EDIFACT from an arbitrary reader as a streaming iterator.
697pub fn from_reader_stream<R: Read>(reader: R) -> OwnedSegmentStream<BufReader<R>> {
698    from_bufread_stream(BufReader::new(reader))
699}
700
701/// Parse EDIFACT from an arbitrary reader as a streaming iterator with custom config.
702///
703/// # Example
704/// ```
705/// use edifact_rs::{ReaderConfig, from_reader_with_config};
706///
707/// let cfg = ReaderConfig::default().max_segment_bytes(4_096);
708/// let segs: Vec<_> = from_reader_with_config(b"BGM+220+1+9'".as_ref(), cfg)
709///     .collect::<Result<_, _>>()
710///     .unwrap();
711/// assert_eq!(segs[0].tag, "BGM");
712/// ```
713pub fn from_reader_with_config<R: Read>(
714    reader: R,
715    config: ReaderConfig,
716) -> OwnedSegmentStream<BufReader<R>> {
717    from_bufread_stream_with_config(BufReader::new(reader), config)
718}
719
720fn read_next_raw_segment<R: BufRead>(
721    reader: &mut R,
722    ssa: &mut crate::tokenizer::ServiceStringAdvice,
723    scanned_header: &mut bool,
724    stream_offset: &mut usize,
725    max_segment_bytes: usize,
726) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
727    loop {
728        let Some((first_offset, first)) = read_next_non_ws_byte(reader, stream_offset)? else {
729            return Ok(None);
730        };
731
732        if !*scanned_header && first == b'U' {
733            let second = read_required_byte(reader, stream_offset)?;
734            let third = read_required_byte(reader, stream_offset)?;
735            if second == b'N' && third == b'A' {
736                let mut una = [0u8; 9];
737                una[0] = b'U';
738                una[1] = b'N';
739                una[2] = b'A';
740                for slot in una.iter_mut().skip(3) {
741                    *slot = read_required_byte(reader, stream_offset)?;
742                }
743                *ssa = crate::tokenizer::ServiceStringAdvice {
744                    component_sep: una[3],
745                    element_sep: una[4],
746                    decimal_mark: una[5],
747                    release_char: una[6],
748                    repetition_sep: una[7],
749                    segment_term: una[8],
750                };
751                if !ssa.is_valid() {
752                    return Err(EdifactError::InvalidUna);
753                }
754                *scanned_header = true;
755                continue;
756            }
757
758            *scanned_header = true;
759            return read_remainder_of_segment(
760                reader,
761                ssa,
762                crate::tokenizer::RawSegment {
763                    bytes: vec![first, second, third],
764                    start_offset: first_offset,
765                },
766                stream_offset,
767                max_segment_bytes,
768            );
769        }
770
771        *scanned_header = true;
772        return read_remainder_of_segment(
773            reader,
774            ssa,
775            crate::tokenizer::RawSegment {
776                bytes: vec![first],
777                start_offset: first_offset,
778            },
779            stream_offset,
780            max_segment_bytes,
781        );
782    }
783}
784
785fn read_remainder_of_segment<R: BufRead>(
786    reader: &mut R,
787    ssa: &crate::tokenizer::ServiceStringAdvice,
788    mut out: crate::tokenizer::RawSegment,
789    stream_offset: &mut usize,
790    max_segment_bytes: usize,
791) -> Result<Option<crate::tokenizer::RawSegment>, EdifactError> {
792    let mut escaped = false;
793    loop {
794        // Strictly greater-than, matching the fast path (`pos > max_segment_bytes`)
795        // and the tokenizer guard.  Using `>=` here made a segment of exactly
796        // `max_segment_bytes` bytes parse on the fast path but fail on the slow
797        // path, so identical input succeeded or failed depending only on whether
798        // it happened to straddle a read-buffer boundary.
799        if out.bytes.len() > max_segment_bytes {
800            return Err(EdifactError::SegmentTooLong {
801                offset: out.start_offset,
802                limit: max_segment_bytes,
803            });
804        }
805        let Some(byte) = read_next_byte(reader, stream_offset)? else {
806            return if out.bytes.is_empty() {
807                Ok(None)
808            } else if escaped {
809                Err(EdifactError::InvalidReleaseSequence {
810                    offset: out.start_offset + out.bytes.len().saturating_sub(1),
811                })
812            } else {
813                Err(EdifactError::UnexpectedEof {
814                    offset: out.start_offset + out.bytes.len(),
815                })
816            };
817        };
818
819        if !escaped && byte == ssa.segment_term {
820            return Ok(Some(out));
821        }
822
823        if !escaped && byte == ssa.release_char {
824            escaped = true;
825            out.bytes.push(byte);
826            continue;
827        }
828
829        escaped = false;
830        out.bytes.push(byte);
831    }
832}
833
834fn read_next_byte<R: BufRead>(
835    reader: &mut R,
836    stream_offset: &mut usize,
837) -> Result<Option<u8>, EdifactError> {
838    let buf = reader.fill_buf()?;
839    if buf.is_empty() {
840        return Ok(None);
841    }
842
843    let byte = buf[0];
844    reader.consume(1);
845    // Saturating add: on 32-bit targets `stream_offset` is a `usize` clamped from a
846    // `u64` field.  Plain `+= 1` would wrap to 0 once the counter reaches `usize::MAX`
847    // and corrupt subsequent span diagnostics / `bytes_consumed` accounting.
848    // Using an explicit local avoids relying on `&mut` auto-deref evaluation order.
849    let next_offset = stream_offset.saturating_add(1);
850    *stream_offset = next_offset;
851    Ok(Some(byte))
852}
853
854fn read_required_byte<R: BufRead>(
855    reader: &mut R,
856    stream_offset: &mut usize,
857) -> Result<u8, EdifactError> {
858    read_next_byte(reader, stream_offset)?.ok_or(EdifactError::UnexpectedEof {
859        offset: *stream_offset,
860    })
861}
862
863fn read_next_non_ws_byte<R: BufRead>(
864    reader: &mut R,
865    stream_offset: &mut usize,
866) -> Result<Option<(usize, u8)>, EdifactError> {
867    loop {
868        let current_offset = *stream_offset;
869        let Some(byte) = read_next_byte(reader, stream_offset)? else {
870            return Ok(None);
871        };
872        if !matches!(byte, b' ' | b'\t' | b'\r' | b'\n') {
873            return Ok(Some((current_offset, byte)));
874        }
875    }
876}
877
878#[cfg(test)]
879mod tests {
880    use super::*;
881    use crate::tokenizer::ServiceStringAdvice;
882
883    fn parse_all(input: &[u8]) -> Vec<Segment<'_>> {
884        let ssa = ServiceStringAdvice::from_bytes_unchecked(input);
885        let tok = Tokenizer::new(input, ssa);
886        Parser::new(tok)
887            .collect::<Result<Vec<_>, _>>()
888            .expect("parse failed")
889    }
890
891    #[test]
892    fn parses_unb_unz() {
893        let input = b"UNB+UNOA:1+SENDER+RECEIVER+200101:0900+1'UNZ+0+1'";
894        let segs = parse_all(input);
895        assert_eq!(segs.len(), 2);
896        assert_eq!(segs[0].tag, "UNB");
897        assert_eq!(segs[1].tag, "UNZ");
898        assert_eq!(segs[0].tag_span, Span::new(0, 3));
899        assert_eq!(segs[0].span, Span::new(0, 41));
900    }
901
902    #[test]
903    fn element_access() {
904        let input = b"BGM+220+ORDER123+9'";
905        let segs = parse_all(input);
906        assert_eq!(segs[0].element_str(0), Some("220"));
907        assert_eq!(segs[0].element_str(1), Some("ORDER123"));
908    }
909
910    #[test]
911    fn component_access() {
912        let input = b"DTM+137:20200101:102'";
913        let segs = parse_all(input);
914        let dtm = &segs[0];
915        assert_eq!(dtm.get_element(0).unwrap().get_component(0), Some("137"));
916        assert_eq!(
917            dtm.get_element(0).unwrap().get_component(1),
918            Some("20200101")
919        );
920        assert_eq!(dtm.get_element(0).unwrap().get_component(2), Some("102"));
921    }
922
923    #[test]
924    fn release_char_resolved() {
925        let input = b"FTX+AAA++test?+value'";
926        let segs = parse_all(input);
927        assert_eq!(segs[0].element_str(2), Some("test+value"));
928        assert_eq!(
929            segs[0].get_element(2).unwrap().component_span(0),
930            Some(Span::new(9, 20))
931        );
932    }
933
934    #[test]
935    fn reader_path_preserves_custom_una_delimiters() {
936        let input = b"UNA:;.? 'BGM;220;test?;value'";
937        let segments = super::from_bufread(std::io::BufReader::new(std::io::Cursor::new(input)))
938            .expect("reader parse should succeed");
939        let bgm = segments
940            .iter()
941            .find(|segment| segment.tag == "BGM")
942            .expect("BGM segment should be present");
943        assert_eq!(bgm.elements[0].components[0].0, "220");
944        assert_eq!(bgm.elements[1].components[0].0, "test;value");
945    }
946
947    #[test]
948    fn arbitrary_bytes_no_panic() {
949        // This is the stable no-panic property — arbitrary input must not panic
950        let garbage: &[u8] = b"\xff\x00\x01\x02ABC+++'''???";
951        let _ = crate::from_bytes(garbage).collect::<Vec<_>>();
952    }
953
954    #[test]
955    fn from_reader_handles_chunk_boundaries() {
956        let input = b"UNA:+.? 'BGM+220+test?+value'UNT+2+1'";
957        let reader = std::io::BufReader::with_capacity(5, std::io::Cursor::new(input));
958        let parsed = from_bufread(reader).expect("reader parsing should succeed");
959        assert_eq!(parsed.len(), 2);
960        assert_eq!(parsed[0].tag, "BGM");
961        assert_eq!(parsed[0].elements[1].components[0].0, "test+value");
962        assert_eq!(parsed[1].tag, "UNT");
963    }
964
965    #[test]
966    fn from_reader_without_una_uses_default_delimiters() {
967        let input = b"BGM+220+X'UNT+2+1'";
968        let parsed =
969            from_reader(std::io::Cursor::new(input)).expect("reader parsing should succeed");
970        assert_eq!(parsed.len(), 2);
971        assert_eq!(parsed[0].tag, "BGM");
972        assert_eq!(parsed[0].elements[0].components[0].0, "220");
973        assert_eq!(parsed[1].span, Span::new(10, 18));
974    }
975
976    #[test]
977    fn dangling_release_sequence_is_error() {
978        let input = b"FTX+AAA++dangling?";
979        let err = crate::from_bytes(input)
980            .collect::<Result<Vec<_>, _>>()
981            .expect_err("expected dangling release to fail");
982
983        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
984    }
985
986    #[test]
987    fn from_reader_reports_dangling_release_sequence() {
988        let input = b"FTX+AAA++dangling?";
989        let err = from_reader(std::io::Cursor::new(input))
990            .expect_err("expected dangling release from reader path");
991        assert!(matches!(err, EdifactError::InvalidReleaseSequence { .. }));
992    }
993
994    #[test]
995    fn a_segment_tagged_una_parses_the_same_on_both_paths() {
996        // `UNA` is three ASCII uppercase letters, so it is a syntactically legal
997        // segment tag.  The reader re-tokenizes each segment from its own slice,
998        // where the whole-interchange "skip nine bytes of service string advice"
999        // heuristic is wrong: it ate the tag and everything after it, and the
1000        // identical bytes that parsed cleanly through `from_bytes` came back as
1001        // `InvalidSegmentTag` through a reader.
1002        let input = b"BGM+220'UNA+XXXXXX'BGM+221'";
1003
1004        let from_slice: Vec<_> = crate::from_bytes(input)
1005            .collect::<Result<Vec<_>, _>>()
1006            .expect("slice path");
1007        let from_reader =
1008            from_reader(std::io::Cursor::new(&input[..])).expect("reader path must agree");
1009
1010        assert_eq!(
1011            from_slice.iter().map(|s| s.tag).collect::<Vec<_>>(),
1012            from_reader
1013                .iter()
1014                .map(|s| s.tag.as_str())
1015                .collect::<Vec<_>>(),
1016        );
1017        assert_eq!(from_reader[1].element_str(0), Some("XXXXXX"));
1018    }
1019
1020    #[test]
1021    fn from_reader_rejects_invalid_una() {
1022        let input = b"UNA::.? 'BGM:220'";
1023        let err = from_reader(std::io::Cursor::new(input))
1024            .expect_err("invalid UNA should fail reader parsing");
1025        assert!(matches!(err, EdifactError::InvalidUna));
1026    }
1027}