pgdumpx 0.2.0

Read-only, bounded inspection, extraction, and row scanning for PostgreSQL custom-format dumps
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
use crate::{Limits, PgDumpError, ScanLimits};
use std::{
    fmt,
    io::{BufRead, BufReader, Read},
    iter::FusedIterator,
};

const INITIAL_ROW_CAPACITY_BYTES: usize = 8 * 1024;
const COPY_TERMINATOR: &[u8] = b"\\.";

/// A borrowed logical field from a PostgreSQL COPY text row.
///
/// This type is byte-oriented: valid UTF-8 is not required. [`FieldRef::Bytes`]
/// contains the logical bytes after COPY-text backslash decoding, while the exact
/// unescaped `\N` field spelling is represented by [`FieldRef::Null`]. An empty
/// field is therefore `Bytes(b"")`, not `Null`.
///
/// The bytes borrow the reusable storage of the [`Row`] that produced them and
/// cannot outlive that row.
///
/// # Example
///
/// ```
/// use pgdumpx::{CopyRowReader, FieldRef, PgDumpError};
/// use std::io::Cursor;
///
/// # fn main() -> Result<(), PgDumpError> {
/// let input = b"hello\\tworld\t\\N\n\\.\n";
/// let mut rows = CopyRowReader::new(Cursor::new(input));
/// let row = rows.next_row()?.expect("one data row");
///
/// assert_eq!(row.field(0), Some(FieldRef::Bytes(b"hello\tworld")));
/// assert_eq!(row.field(1), Some(FieldRef::Null));
/// assert!(rows.next_row()?.is_none());
/// # Ok(())
/// # }
/// ```
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum FieldRef<'a> {
    /// The raw field spelling was exactly PostgreSQL's `\N` NULL marker.
    Null,
    /// Logical field bytes after PostgreSQL COPY backslash decoding.
    Bytes(&'a [u8]),
}

/// An owned logical field copied from one matched COPY text row.
///
/// [`OwnedField::Bytes`] has the same post-unescape byte semantics as
/// [`FieldRef::Bytes`], but owns its storage so it can outlive the row reader.
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum OwnedField {
    /// A PostgreSQL COPY NULL value.
    Null,
    /// Logical field bytes after PostgreSQL COPY backslash decoding.
    Bytes(Vec<u8>),
}

/// One owned COPY text row that can outlive its streaming reader.
///
/// `find_first` materializes only the matching row into this type; non-matching
/// rows continue to use the lending [`Row`] representation. Field order is the
/// original COPY order and values keep the byte-oriented semantics of [`FieldRef`].
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OwnedRow {
    fields: Vec<OwnedField>,
}

#[allow(clippy::len_without_is_empty)]
impl OwnedRow {
    /// Returns the number of fields in this row.
    pub fn len(&self) -> usize {
        self.fields.len()
    }

    /// Returns one owned field by zero-based index.
    pub fn field(&self, index: usize) -> Option<&OwnedField> {
        self.fields.get(index)
    }

    /// Returns all owned fields in COPY column order.
    pub fn fields(&self) -> &[OwnedField] {
        &self.fields
    }

    pub(crate) fn try_from_borrowed(row: &Row<'_>) -> Result<Self, PgDumpError> {
        let mut fields = Vec::new();
        fields.try_reserve_exact(row.len()).map_err(|_| {
            PgDumpError::CopyFieldAllocationFailed {
                row: row.number,
                requested: u64::try_from(row.len()).unwrap_or(u64::MAX),
            }
        })?;

        for field in row.fields() {
            let owned = match field {
                FieldRef::Null => OwnedField::Null,
                FieldRef::Bytes(bytes) => {
                    let mut owned = Vec::new();
                    owned.try_reserve_exact(bytes.len()).map_err(|_| {
                        PgDumpError::CopyRowAllocationFailed {
                            row: row.number,
                            requested: u64::try_from(bytes.len()).unwrap_or(u64::MAX),
                        }
                    })?;
                    owned.extend_from_slice(bytes);
                    OwnedField::Bytes(owned)
                }
            };
            fields.push(owned);
        }

        Ok(Self { fields })
    }
}

/// A borrowed COPY text row backed by reusable parser storage.
///
/// The row and all [`FieldRef::Bytes`] slices obtained from it remain valid only
/// until the originating [`CopyRowReader`] (or a wrapper such as
/// [`crate::TableRowReader`]) is mutably borrowed again. This lending shape is
/// why row readers expose `next_row(&mut self)` instead of implementing
/// [`Iterator`]. Use [`OwnedRow`] when data must survive reader advancement.
pub struct Row<'a> {
    number: u64,
    bytes: &'a [u8],
    fields: &'a [FieldSpan],
}

impl Row<'_> {
    /// Returns the number of logical fields in this row.
    pub fn len(&self) -> usize {
        self.fields.len()
    }

    /// Returns whether this row has no fields.
    pub fn is_empty(&self) -> bool {
        self.fields.is_empty()
    }

    /// Returns one logical field by zero-based index.
    ///
    /// `None` means only that `index` is outside this row; SQL NULL is represented
    /// explicitly as [`FieldRef::Null`].
    pub fn field(&self, index: usize) -> Option<FieldRef<'_>> {
        self.fields.get(index).map(|span| span.resolve(self.bytes))
    }

    /// Iterates over all borrowed logical fields in COPY column order.
    ///
    /// The yielded fields borrow this row and therefore share its lending lifetime.
    pub fn fields(&self) -> impl ExactSizeIterator<Item = FieldRef<'_>> + FusedIterator + '_ {
        self.fields.iter().map(move |span| span.resolve(self.bytes))
    }
}

impl fmt::Debug for Row<'_> {
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        formatter.debug_list().entries(self.fields()).finish()
    }
}

#[derive(Clone, Copy, Debug)]
enum FieldSpan {
    Null,
    Bytes { start: usize, end: usize },
}

impl FieldSpan {
    fn resolve<'a>(self, bytes: &'a [u8]) -> FieldRef<'a> {
        match self {
            Self::Null => FieldRef::Null,
            Self::Bytes { start, end } => FieldRef::Bytes(&bytes[start..end]),
        }
    }
}

#[derive(Clone, Copy, Debug)]
struct ScanBudget {
    limits: ScanLimits,
    start_consumed: u64,
    consumed_rows: u64,
}

impl ScanBudget {
    const fn new(limits: ScanLimits, start_consumed: u64) -> Self {
        Self {
            limits,
            start_consumed,
            consumed_rows: 0,
        }
    }

    fn check_bytes(&self, row: u64, current: u64) -> Result<(), PgDumpError> {
        let consumed = current
            .checked_sub(self.start_consumed)
            .ok_or(PgDumpError::ArithmeticOverflow { offset: current })?;
        if let Some(limit) = self.limits.max_decompressed_bytes() {
            if consumed > limit {
                return Err(PgDumpError::ScanDecompressedByteLimitExceeded {
                    row,
                    limit,
                    consumed,
                    byte_offset: current,
                });
            }
        }
        Ok(())
    }

    fn next_row_count(&self, row: u64) -> Result<u64, PgDumpError> {
        self.consumed_rows
            .checked_add(1)
            .ok_or(PgDumpError::ScanRowCountOverflow {
                row,
                consumed: self.consumed_rows,
            })
    }

    fn check_rows(&self, row: u64, consumed: u64) -> Result<(), PgDumpError> {
        if let Some(limit) = self.limits.max_rows() {
            if consumed > limit {
                return Err(PgDumpError::ScanRowLimitExceeded {
                    row,
                    limit,
                    consumed,
                });
            }
        }
        Ok(())
    }
}

/// A lending, byte-oriented parser for PostgreSQL COPY text rows.
///
/// Input is consumed incrementally from any [`Read`] implementation. The parser
/// buffers only the current physical row plus decoded logical fields; it never
/// buffers the complete COPY stream. Archive framing and decompression are not
/// performed by this standalone parser; [`crate::Archive::table_rows`] composes
/// those layers for custom-format table data.
///
/// Structural [`Limits`] and total-work [`ScanLimits`] are distinct. Structural
/// limits bound one row/field layout. Constructor-level scan limits accumulate
/// across subsequent row operations on this reader. Operation-level limits passed
/// to [`CopyRowReader::find_first_with_limits`] add a second budget measured from
/// that call's current stream position.
///
/// Any error returned while reading or parsing a row makes this reader terminal.
/// The current record may already be partially consumed, and pgdumpx does not try
/// to resynchronize malformed COPY input or retry a source-I/O failure. Subsequent
/// row-reading and search calls therefore return `Ok(None)` without exposing bytes
/// from the rejected record or any later record.
pub struct CopyRowReader<R> {
    input: CopyInput<R>,
    limits: Limits,
    scan_budget: ScanBudget,
    raw_row: Vec<u8>,
    logical_bytes: Vec<u8>,
    field_spans: Vec<FieldSpan>,
    next_row_number: u64,
    finished: bool,
}

impl<R: Read> CopyRowReader<R> {
    /// Creates a reader with finite default structural limits and unlimited scan work.
    pub fn new(reader: R) -> Self {
        Self::with_limits(reader, Limits::default())
    }

    /// Creates a reader using caller-supplied structural limits.
    ///
    /// Only the row-byte and fields-per-row members are used by this standalone
    /// parser. The same [`Limits`] values are used by [`crate::Archive::table_rows`].
    /// Scan work remains unlimited unless configured separately.
    pub fn with_limits(reader: R, limits: Limits) -> Self {
        Self::with_limits_and_scan_limits(reader, limits, ScanLimits::unlimited())
    }

    /// Creates a reader with finite default structural limits and scan work budgets.
    ///
    /// These scan budgets are reader-wide: consumed rows and parser-consumed bytes
    /// accumulate across successive `next_row` and search calls.
    pub fn with_scan_limits(reader: R, scan_limits: ScanLimits) -> Self {
        Self::with_limits_and_scan_limits(reader, Limits::default(), scan_limits)
    }

    /// Creates a COPY reader using both structural and reader-wide scan limits.
    ///
    /// A scan row limit of `N` permits at most `N` complete rows. A scan byte limit
    /// of `N` permits at most `N` decompressed bytes consumed by this parser. The
    /// first byte or row that would make the corresponding count exceed `N` causes
    /// a typed resource error and is not exposed as a row.
    pub fn with_limits_and_scan_limits(reader: R, limits: Limits, scan_limits: ScanLimits) -> Self {
        let input = CopyInput::new(reader);
        let scan_budget = ScanBudget::new(scan_limits, input.consumed());
        Self {
            input,
            limits,
            scan_budget,
            raw_row: Vec::new(),
            logical_bytes: Vec::new(),
            field_spans: Vec::new(),
            next_row_number: 1,
            finished: false,
        }
    }

    #[cfg(test)]
    pub(crate) fn with_limits_and_consumed(reader: R, limits: Limits, consumed: u64) -> Self {
        let mut parser = Self::with_limits(reader, limits);
        parser.input.consumed = consumed;
        parser.scan_budget.start_consumed = consumed;
        parser
    }

    #[cfg(test)]
    pub(crate) fn with_scan_state_for_test(
        reader: R,
        limits: Limits,
        scan_limits: ScanLimits,
        consumed_rows: u64,
    ) -> Self {
        let mut parser = Self::with_limits_and_scan_limits(reader, limits, scan_limits);
        parser.scan_budget.consumed_rows = consumed_rows;
        parser
    }

    /// Parses and lends the next logical row.
    ///
    /// `Ok(None)` is returned after a standalone `\.` terminator or after the
    /// underlying COPY stream reaches EOF. A returned row borrows reusable parser
    /// storage, so Rust prevents another mutable operation on this reader while
    /// that row is still used.
    ///
    /// Any error makes the reader terminal because the rejected physical record may
    /// have been partially consumed. Later `next_row` or search calls return
    /// `Ok(None)`; the original failing call retains its existing typed error context.
    ///
    /// ```compile_fail
    /// use pgdumpx::{CopyRowReader, PgDumpError};
    /// use std::io::Cursor;
    ///
    /// fn cannot_hold_two_rows() -> Result<(), PgDumpError> {
    ///     let mut rows = CopyRowReader::new(Cursor::new(b"a\nb\n"));
    ///     let first = rows.next_row()?.unwrap();
    ///     let _second = rows.next_row()?;
    ///     println!("{first:?}");
    ///     Ok(())
    /// }
    /// ```
    pub fn next_row(&mut self) -> Result<Option<Row<'_>>, PgDumpError> {
        let mut operation_budget = None;
        self.next_row_with_budget(&mut operation_budget)
    }

    /// Sequentially scans from the current reader position for the first match.
    ///
    /// The predicate is called once for each fully parsed row that is within the
    /// reader-wide scan budget. On the first `true`, only that row is copied into
    /// [`OwnedRow`] and scanning stops immediately. If the remaining stream ends
    /// without a match, this returns `Ok(None)`.
    ///
    /// This is not an indexed lookup. Worst-case unrestricted work is proportional
    /// to the remaining COPY stream. No additional operation-level budget is added
    /// by this convenience method; constructor-level [`ScanLimits`], if any, still
    /// apply.
    pub fn find_first<F>(&mut self, predicate: F) -> Result<Option<OwnedRow>, PgDumpError>
    where
        F: FnMut(&Row<'_>) -> bool,
    {
        self.find_first_with_limits(ScanLimits::unlimited(), predicate)
    }

    /// Sequentially scans for the first match with an additional operation budget.
    ///
    /// `scan_limits` is measured from this call's current stream position and is
    /// enforced in addition to any reader-wide limits supplied at construction.
    /// Decompressed-byte accounting advances as physical COPY bytes are consumed by
    /// the parser, including separators and terminator bytes; decoder or buffered
    /// read-ahead that has not been consumed does not count. Complete rows are
    /// counted before decoding or predicate invocation, so a row that crosses either
    /// budget is never exposed to `predicate`. The matching row itself counts toward
    /// both budgets, and no bytes after a successful match are consumed by the scan.
    pub fn find_first_with_limits<F>(
        &mut self,
        scan_limits: ScanLimits,
        mut predicate: F,
    ) -> Result<Option<OwnedRow>, PgDumpError>
    where
        F: FnMut(&Row<'_>) -> bool,
    {
        let mut operation_budget = Some(ScanBudget::new(scan_limits, self.consumed_input_bytes()));
        while let Some(row) = self.next_row_with_budget(&mut operation_budget)? {
            if predicate(&row) {
                return OwnedRow::try_from_borrowed(&row).map(Some);
            }
        }
        Ok(None)
    }

    pub(crate) const fn consumed_input_bytes(&self) -> u64 {
        self.input.consumed()
    }

    fn next_row_with_budget(
        &mut self,
        operation_budget: &mut Option<ScanBudget>,
    ) -> Result<Option<Row<'_>>, PgDumpError> {
        if self.finished {
            return Ok(None);
        }

        // Assume terminal until a complete line-delimited row has been parsed. Any
        // early-returning error then leaves the reader safely terminal without
        // replacing the original typed error or attempting record resynchronization.
        self.finished = true;

        self.raw_row.clear();
        self.logical_bytes.clear();
        self.field_spans.clear();

        let row = self.next_row_number;
        let row_start = self.consumed_input_bytes();
        let record_end = self.read_record(row, operation_budget)?;

        if self.raw_row.is_empty() && record_end == RecordEnd::Eof {
            return Ok(None);
        }

        if self.raw_row == COPY_TERMINATOR {
            if record_end == RecordEnd::Line {
                return Ok(None);
            }
            return Err(PgDumpError::MalformedCopyTerminator {
                row,
                byte_offset: row_start,
            });
        }

        self.count_scanned_row(row, operation_budget)?;

        let field_count = inspect_field_layout(
            &self.raw_row,
            self.limits.max_fields_per_row(),
            row,
            row_start,
        )?;
        self.prepare_decoded_storage(row, field_count)?;
        decode_fields(
            &self.raw_row,
            &mut self.logical_bytes,
            &mut self.field_spans,
            row,
            row_start,
        )?;

        if record_end == RecordEnd::Line {
            self.next_row_number = row
                .checked_add(1)
                .ok_or(PgDumpError::CopyRowNumberOverflow { row })?;
            self.finished = false;
        }

        Ok(Some(Row {
            number: row,
            bytes: &self.logical_bytes,
            fields: &self.field_spans,
        }))
    }

    fn read_record(
        &mut self,
        row: u64,
        operation_budget: &Option<ScanBudget>,
    ) -> Result<RecordEnd, PgDumpError> {
        let mut escaped = false;
        loop {
            let Some(byte) = self.next_input_byte(row, operation_budget)? else {
                if escaped {
                    return Err(PgDumpError::MalformedCopyEscape {
                        row,
                        byte_offset: self.input.consumed().saturating_sub(1),
                    });
                }
                return Ok(RecordEnd::Eof);
            };

            if escaped {
                self.push_raw_byte(row, byte)?;
                escaped = false;
                continue;
            }

            match byte {
                b'\\' => {
                    self.push_raw_byte(row, byte)?;
                    escaped = true;
                }
                b'\n' => return Ok(RecordEnd::Line),
                b'\r' => {
                    if self.input.peek_byte(row)? == Some(b'\n') {
                        let consumed = self.next_input_byte(row, operation_budget)?;
                        debug_assert_eq!(consumed, Some(b'\n'));
                    }
                    return Ok(RecordEnd::Line);
                }
                _ => self.push_raw_byte(row, byte)?,
            }
        }
    }

    fn next_input_byte(
        &mut self,
        row: u64,
        operation_budget: &Option<ScanBudget>,
    ) -> Result<Option<u8>, PgDumpError> {
        let byte = self.input.next_byte(row)?;
        if byte.is_some() {
            if let Err(error) = self.check_scan_bytes(row, operation_budget) {
                self.finished = true;
                return Err(error);
            }
        }
        Ok(byte)
    }

    fn check_scan_bytes(
        &self,
        row: u64,
        operation_budget: &Option<ScanBudget>,
    ) -> Result<(), PgDumpError> {
        let current = self.input.consumed();
        self.scan_budget.check_bytes(row, current)?;
        if let Some(budget) = operation_budget {
            budget.check_bytes(row, current)?;
        }
        Ok(())
    }

    fn count_scanned_row(
        &mut self,
        row: u64,
        operation_budget: &mut Option<ScanBudget>,
    ) -> Result<(), PgDumpError> {
        let result = self.try_count_scanned_row(row, operation_budget);
        if result.is_err() {
            self.finished = true;
        }
        result
    }

    fn try_count_scanned_row(
        &mut self,
        row: u64,
        operation_budget: &mut Option<ScanBudget>,
    ) -> Result<(), PgDumpError> {
        let next_reader_rows = self.scan_budget.next_row_count(row)?;
        let next_operation_rows = operation_budget
            .as_ref()
            .map(|budget| budget.next_row_count(row))
            .transpose()?;

        self.scan_budget.check_rows(row, next_reader_rows)?;
        if let (Some(budget), Some(consumed)) = (operation_budget.as_ref(), next_operation_rows) {
            budget.check_rows(row, consumed)?;
        }

        self.scan_budget.consumed_rows = next_reader_rows;
        if let (Some(budget), Some(consumed)) = (operation_budget.as_mut(), next_operation_rows) {
            budget.consumed_rows = consumed;
        }
        Ok(())
    }

    fn push_raw_byte(&mut self, row: u64, byte: u8) -> Result<(), PgDumpError> {
        let actual = self
            .raw_row
            .len()
            .checked_add(1)
            .ok_or(PgDumpError::ArithmeticOverflow {
                offset: self.input.consumed(),
            })?;
        let limit = self.limits.max_row_bytes();
        if actual > limit {
            return Err(PgDumpError::CopyRowByteLimitExceeded {
                row,
                limit: usize_to_u64(limit, self.input.consumed())?,
                actual: usize_to_u64(actual, self.input.consumed())?,
                byte_offset: self.input.consumed(),
            });
        }

        if actual > self.raw_row.capacity() {
            let proposed = if self.raw_row.capacity() == 0 {
                INITIAL_ROW_CAPACITY_BYTES
            } else {
                self.raw_row.capacity().saturating_mul(2)
            };
            let target = proposed.max(actual).min(limit);
            let additional =
                target
                    .checked_sub(self.raw_row.len())
                    .ok_or(PgDumpError::ArithmeticOverflow {
                        offset: self.input.consumed(),
                    })?;
            self.raw_row.try_reserve_exact(additional).map_err(|_| {
                PgDumpError::CopyRowAllocationFailed {
                    row,
                    requested: u64::try_from(target).unwrap_or(u64::MAX),
                }
            })?;
        }

        self.raw_row.push(byte);
        Ok(())
    }

    fn prepare_decoded_storage(&mut self, row: u64, field_count: usize) -> Result<(), PgDumpError> {
        if self.logical_bytes.capacity() < self.raw_row.len() {
            self.logical_bytes
                .try_reserve_exact(self.raw_row.len())
                .map_err(|_| PgDumpError::CopyRowAllocationFailed {
                    row,
                    requested: u64::try_from(self.raw_row.len()).unwrap_or(u64::MAX),
                })?;
        }

        if self.field_spans.capacity() < field_count {
            self.field_spans
                .try_reserve_exact(field_count)
                .map_err(|_| PgDumpError::CopyFieldAllocationFailed {
                    row,
                    requested: u64::try_from(field_count).unwrap_or(u64::MAX),
                })?;
        }
        Ok(())
    }
}

#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum RecordEnd {
    Line,
    Eof,
}

struct CopyInput<R> {
    reader: BufReader<R>,
    consumed: u64,
}

impl<R: Read> CopyInput<R> {
    fn new(reader: R) -> Self {
        Self {
            reader: BufReader::new(reader),
            consumed: 0,
        }
    }

    const fn consumed(&self) -> u64 {
        self.consumed
    }

    fn peek_byte(&mut self, row: u64) -> Result<Option<u8>, PgDumpError> {
        self.reader
            .fill_buf()
            .map(|buffer| buffer.first().copied())
            .map_err(|source| PgDumpError::CopyIo {
                row,
                consumed: self.consumed,
                source,
            })
    }

    fn next_byte(&mut self, row: u64) -> Result<Option<u8>, PgDumpError> {
        let byte = self.peek_byte(row)?;
        if byte.is_some() {
            let increment = 1;
            let next = self.consumed.checked_add(increment).ok_or(
                PgDumpError::CopyConsumedByteCountOverflow {
                    row,
                    consumed: self.consumed,
                    increment,
                },
            )?;
            self.reader.consume(1);
            self.consumed = next;
        }
        Ok(byte)
    }
}

fn inspect_field_layout(
    raw: &[u8],
    max_fields: usize,
    row: u64,
    row_start: u64,
) -> Result<usize, PgDumpError> {
    let mut fields = 1_usize;
    let mut index = 0_usize;
    while index < raw.len() {
        if raw[index] == b'\\' {
            let Some(escaped) = raw.get(index + 1).copied() else {
                return Err(PgDumpError::MalformedCopyEscape {
                    row,
                    byte_offset: byte_offset(row_start, index)?,
                });
            };
            if escaped == b'.' {
                return Err(PgDumpError::MalformedCopyTerminator {
                    row,
                    byte_offset: byte_offset(row_start, index)?,
                });
            }
            index += 2;
            continue;
        }

        if raw[index] == b'\t' {
            fields = fields
                .checked_add(1)
                .ok_or(PgDumpError::ArithmeticOverflow {
                    offset: byte_offset(row_start, index)?,
                })?;
            if fields > max_fields {
                return Err(PgDumpError::CopyFieldCountLimitExceeded {
                    row,
                    limit: usize_to_u64(max_fields, row_start)?,
                    actual: usize_to_u64(fields, row_start)?,
                    byte_offset: byte_offset(row_start, index)?,
                });
            }
        }
        index += 1;
    }

    if fields > max_fields {
        return Err(PgDumpError::CopyFieldCountLimitExceeded {
            row,
            limit: usize_to_u64(max_fields, row_start)?,
            actual: usize_to_u64(fields, row_start)?,
            byte_offset: row_start,
        });
    }
    Ok(fields)
}

fn decode_fields(
    raw: &[u8],
    logical: &mut Vec<u8>,
    spans: &mut Vec<FieldSpan>,
    row: u64,
    row_start: u64,
) -> Result<(), PgDumpError> {
    let mut field_start = 0_usize;
    let mut index = 0_usize;
    while index < raw.len() {
        match raw[index] {
            b'\\' => index += 2,
            b'\t' => {
                decode_field(
                    &raw[field_start..index],
                    logical,
                    spans,
                    row,
                    byte_offset(row_start, field_start)?,
                )?;
                field_start = index + 1;
                index += 1;
            }
            _ => index += 1,
        }
    }
    decode_field(
        &raw[field_start..],
        logical,
        spans,
        row,
        byte_offset(row_start, field_start)?,
    )
}

fn decode_field(
    raw: &[u8],
    logical: &mut Vec<u8>,
    spans: &mut Vec<FieldSpan>,
    row: u64,
    field_offset: u64,
) -> Result<(), PgDumpError> {
    if raw == b"\\N" {
        spans.push(FieldSpan::Null);
        return Ok(());
    }

    let start = logical.len();
    let mut index = 0_usize;
    while index < raw.len() {
        let byte = raw[index];
        if byte != b'\\' {
            logical.push(byte);
            index += 1;
            continue;
        }

        let Some(escaped) = raw.get(index + 1).copied() else {
            return Err(PgDumpError::MalformedCopyEscape {
                row,
                byte_offset: byte_offset(field_offset, index)?,
            });
        };
        index += 2;

        let decoded = match escaped {
            b'0'..=b'7' => {
                let mut value = u16::from(escaped - b'0');
                let mut digits = 1_u8;
                while digits < 3 {
                    let Some(next) = raw.get(index).copied() else {
                        break;
                    };
                    if !(b'0'..=b'7').contains(&next) {
                        break;
                    }
                    value = (value << 3) + u16::from(next - b'0');
                    index += 1;
                    digits += 1;
                }
                (value & 0xff) as u8
            }
            b'x' => {
                let Some(first) = raw.get(index).copied().and_then(hex_value) else {
                    logical.push(b'x');
                    continue;
                };
                index += 1;
                let mut value = first;
                if let Some(second) = raw.get(index).copied().and_then(hex_value) {
                    value = (value << 4) | second;
                    index += 1;
                }
                value
            }
            b'b' => 0x08,
            b'f' => 0x0c,
            b'n' => b'\n',
            b'r' => b'\r',
            b't' => b'\t',
            b'v' => 0x0b,
            other => other,
        };
        logical.push(decoded);
    }

    spans.push(FieldSpan::Bytes {
        start,
        end: logical.len(),
    });
    Ok(())
}

fn byte_offset(start: u64, index: usize) -> Result<u64, PgDumpError> {
    let index =
        u64::try_from(index).map_err(|_| PgDumpError::ArithmeticOverflow { offset: start })?;
    start
        .checked_add(index)
        .ok_or(PgDumpError::ArithmeticOverflow { offset: start })
}

fn usize_to_u64(value: usize, offset: u64) -> Result<u64, PgDumpError> {
    u64::try_from(value).map_err(|_| PgDumpError::ArithmeticOverflow { offset })
}

const fn hex_value(byte: u8) -> Option<u8> {
    match byte {
        b'0'..=b'9' => Some(byte - b'0'),
        b'a'..=b'f' => Some(byte - b'a' + 10),
        b'A'..=b'F' => Some(byte - b'A' + 10),
        _ => None,
    }
}