1#![allow(clippy::missing_inline_in_public_items)]
11
12use crate::options::{DecodeOptions, EncodeOptions, RecordFormat, ZonedEncodingFormat};
13use crate::zoned_overpunch::ZeroSignPolicy;
14use base64::Engine;
15use copybook_core::{Error, ErrorCode, Result, Schema};
16use serde_json::Value;
17use std::collections::HashMap;
18use std::convert::TryFrom;
19use std::io::{BufRead, BufReader, Read, Write};
20use std::sync::Arc;
21use tracing::info;
22
23mod envelope;
24mod run_summary;
25mod telemetry;
26mod warnings;
27
28use envelope::{RecordMetadata, build_json_envelope};
29pub use run_summary::{MAX_CAPTURED_FAILURES, RecordFailure, RunSummary};
30pub use warnings::increment_warning_counter;
31use warnings::{reset_warning_counter, warning_count};
32
33const MAX_WORKERS: usize = 64;
34
35#[derive(Clone, Copy)]
36enum RawCapture {
37 Record,
38 RecordRdw,
39}
40
41impl RawCapture {
42 const fn as_str(self) -> &'static str {
43 match self {
44 Self::Record => "record",
45 Self::RecordRdw => "record+rdw",
46 }
47 }
48}
49
50struct RawRecord {
51 b64: String,
52 capture: RawCapture,
53}
54
55fn parse_raw_rdw_frame(frame: &[u8]) -> Result<(u16, &[u8])> {
56 let (raw_header, raw_payload) = frame.split_at_checked(4).ok_or_else(|| {
57 Error::new(
58 ErrorCode::CBKF102_RECORD_LENGTH_INVALID,
59 format!(
60 "Raw RDW record is {} bytes; expected at least a 4-byte header",
61 frame.len()
62 ),
63 )
64 })?;
65 let header_bytes: [u8; 4] = raw_header.try_into().map_err(|_| {
66 Error::new(
67 ErrorCode::CBKF102_RECORD_LENGTH_INVALID,
68 "Raw RDW record does not contain a complete 4-byte header",
69 )
70 })?;
71 let header = copybook_rdw::RdwHeader::from_bytes(header_bytes);
72 let declared_payload_len = usize::from(header.length());
73 if declared_payload_len != raw_payload.len() {
74 return Err(Error::new(
75 ErrorCode::CBKF102_RECORD_LENGTH_INVALID,
76 format!(
77 "Raw RDW header declares {declared_payload_len} payload bytes, but {} bytes follow",
78 raw_payload.len()
79 ),
80 ));
81 }
82 Ok((header.reserved(), raw_payload))
83}
84
85fn validate_captured_raw_rdw(frame: &[u8], expected_payload: &[u8]) -> Result<()> {
86 let (_, raw_payload) = parse_raw_rdw_frame(frame)?;
87 if raw_payload != expected_payload {
88 return Err(Error::new(
89 ErrorCode::CBKF102_RECORD_LENGTH_INVALID,
90 "Raw RDW payload does not match the decoded record payload",
91 ));
92 }
93 Ok(())
94}
95
96fn captured_raw_record(
97 data: &[u8],
98 supplied_raw: Option<&[u8]>,
99 mode: crate::options::RawMode,
100 format: crate::options::RecordFormat,
101) -> Result<Option<RawRecord>> {
102 let (bytes, capture) = match mode {
103 crate::options::RawMode::Off | crate::options::RawMode::Field => return Ok(None),
104 crate::options::RawMode::Record => (data, RawCapture::Record),
105 crate::options::RawMode::RecordRDW => {
106 let frame = supplied_raw.ok_or_else(|| {
107 Error::new(
108 ErrorCode::CBKF102_RECORD_LENGTH_INVALID,
109 "RawMode::RecordRDW requires an RDW header plus payload",
110 )
111 })?;
112 if format == crate::options::RecordFormat::Vb {
115 let (_, raw_payload) = parse_vb_raw_rdw_frame(frame)?;
116 if raw_payload != data {
117 return Err(Error::new(
118 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
119 "Raw VB payload does not match the decoded record payload",
120 ));
121 }
122 } else {
123 validate_captured_raw_rdw(frame, data)?;
124 }
125 (frame, RawCapture::RecordRdw)
126 }
127 };
128 Ok(Some(RawRecord {
129 b64: base64::engine::general_purpose::STANDARD.encode(bytes),
130 capture,
131 }))
132}
133
134#[inline]
144#[must_use = "Handle the Result or propagate the error"]
145pub fn decode_record(schema: &Schema, data: &[u8], options: &DecodeOptions) -> Result<Value> {
146 decode_record_with_raw_data(schema, data, options, None, 0)
147}
148
149#[inline]
185#[must_use = "Handle the Result or propagate the error"]
186pub fn decode_record_with_scratch(
187 schema: &Schema,
188 data: &[u8],
189 options: &DecodeOptions,
190 scratch: &mut crate::memory::ScratchBuffers,
191) -> Result<Value> {
192 decode_record_with_scratch_and_raw(schema, data, options, None, 0, None, scratch)
193}
194
195fn decode_record_with_scratch_and_raw(
197 schema: &Schema,
198 data: &[u8],
199 options: &DecodeOptions,
200 raw_data: Option<&[u8]>,
201 record_index: u64,
202 record_offset: Option<u64>,
203 scratch: &mut crate::memory::ScratchBuffers,
204) -> Result<Value> {
205 use serde_json::Map;
206
207 let mut fields_map = Map::new();
208 let mut encoding_acc = Vec::new();
209 let record_raw = captured_raw_record(data, raw_data, options.emit_raw, options.format)?;
210
211 process_fields_recursive_with_scratch(
212 &schema.fields,
213 data,
214 &mut fields_map,
215 options,
216 scratch,
217 record_index,
218 &mut encoding_acc,
219 )?;
220
221 Ok(build_json_envelope(
222 fields_map,
223 schema,
224 options,
225 record_index,
226 &RecordMetadata {
227 length: data.len(),
228 offset: record_offset,
229 },
230 record_raw,
231 encoding_acc,
232 ))
233}
234
235#[inline]
240#[must_use = "Handle the Result or propagate the error"]
241pub fn decode_record_with_raw_data(
242 schema: &Schema,
243 data: &[u8],
244 options: &DecodeOptions,
245 raw_data_with_header: Option<&[u8]>,
246 record_index: u64,
247) -> Result<Value> {
248 decode_record_with_raw_data_at_offset(
249 schema,
250 data,
251 options,
252 raw_data_with_header,
253 record_index,
254 None,
255 )
256}
257
258fn decode_record_with_raw_data_at_offset(
259 schema: &Schema,
260 data: &[u8],
261 options: &DecodeOptions,
262 raw_data_with_header: Option<&[u8]>,
263 record_index: u64,
264 record_offset: Option<u64>,
265) -> Result<Value> {
266 use serde_json::Map;
267
268 let record_raw =
271 captured_raw_record(data, raw_data_with_header, options.emit_raw, options.format)?;
272
273 let mut fields_map = Map::new();
274 let mut scratch_buffers: Option<crate::memory::ScratchBuffers> = None;
275 let mut encoding_acc = Vec::new();
276
277 process_fields_recursive(
278 &schema.fields,
279 data,
280 &mut fields_map,
281 options,
282 &mut scratch_buffers,
283 record_index,
284 &mut encoding_acc,
285 )?;
286
287 Ok(build_json_envelope(
288 fields_map,
289 schema,
290 options,
291 record_index,
292 &RecordMetadata {
293 length: data.len(),
294 offset: record_offset,
295 },
296 record_raw,
297 encoding_acc,
298 ))
299}
300
301fn process_fields_recursive(
306 fields: &[copybook_core::Field],
307 data: &[u8],
308 json_obj: &mut serde_json::Map<String, Value>,
309 options: &DecodeOptions,
310 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
311 record_index: u64,
312 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
313) -> Result<()> {
314 use copybook_core::FieldKind;
315
316 let total_fields = fields.len();
317 let mut deferred_group_views = Vec::new();
318
319 for (field_index, field) in fields.iter().enumerate() {
320 match (&field.kind, &field.occurs) {
321 (_, Some(occurs)) => {
322 process_array_field(
323 field,
324 occurs,
325 data,
326 json_obj,
327 options,
328 fields,
329 scratch_buffers,
330 record_index,
331 encoding_acc,
332 )?;
333 }
334 (FieldKind::Group, None) if field.level > 1 => {
335 let mut group_obj = serde_json::Map::new();
336 let metadata_start = encoding_acc.len();
337 process_fields_recursive(
338 &field.children,
339 data,
340 &mut group_obj,
341 options,
342 scratch_buffers,
343 record_index,
344 encoding_acc,
345 )?;
346 if is_scalar_target_group_redefine(field, fields) {
347 let group_value = Value::Object(group_obj);
348 if let Value::Object(group_fields) = &group_value {
349 insert_decoded_group_fields(
350 json_obj,
351 group_fields,
352 &mut encoding_acc[metadata_start..],
353 options.emit_filler,
354 );
355 }
356 deferred_group_views.push((field.name.clone(), group_value));
357 } else if field.redefines_of.is_none() {
358 insert_decoded_field(json_obj, &field.name, Value::Object(group_obj));
359 }
360 }
361 (FieldKind::Group, None) => {
362 process_fields_recursive(
363 &field.children,
364 data,
365 json_obj,
366 options,
367 scratch_buffers,
368 record_index,
369 encoding_acc,
370 )?;
371 }
372 _ => {
373 process_scalar_field_standard(
374 field,
375 field_index,
376 total_fields,
377 data,
378 json_obj,
379 options,
380 scratch_buffers,
381 record_index,
382 encoding_acc,
383 )?;
384 }
385 }
386 }
387
388 for (name, value) in deferred_group_views {
389 insert_decoded_field(json_obj, &name, value);
390 }
391
392 Ok(())
393}
394
395fn insert_decoded_field(json_obj: &mut serde_json::Map<String, Value>, name: &str, value: Value) {
401 let _ = insert_decoded_field_with_key(json_obj, name, value);
402}
403
404fn insert_decoded_group_fields(
407 json_obj: &mut serde_json::Map<String, Value>,
408 group_fields: &serde_json::Map<String, Value>,
409 encoding_metadata: &mut [(String, ZonedEncodingFormat)],
410 allow_filler_collision: bool,
411) {
412 let mut emitted_keys = Vec::new();
413 let mut metadata_used = vec![false; encoding_metadata.len()];
414 for (name, value) in group_fields {
415 if let Some(field_name) = name.strip_suffix("_raw_b64")
416 && let Some((_, emitted_key)) = emitted_keys
417 .iter()
418 .rev()
419 .find(|(original, _)| original == field_name)
420 {
421 json_obj.insert(format!("{emitted_key}_raw_b64"), value.clone());
422 continue;
423 }
424
425 let emitted_key = if allow_filler_collision {
426 insert_decoded_field_with_key_emitting_filler(json_obj, name, value.clone())
427 } else {
428 insert_decoded_field_with_key(json_obj, name, value.clone())
429 }
430 .unwrap_or_else(|| name.clone());
431 if let Some((metadata_index, (metadata_name, _))) = encoding_metadata
432 .iter_mut()
433 .enumerate()
434 .find(|(index, (metadata_name, _))| !metadata_used[*index] && metadata_name == name)
435 {
436 metadata_used[metadata_index] = true;
437 metadata_name.clone_from(&emitted_key);
438 }
439 if !name.ends_with("_raw_b64") {
440 emitted_keys.push((name.clone(), emitted_key));
441 }
442 }
443}
444
445fn insert_decoded_field_with_key(
447 json_obj: &mut serde_json::Map<String, Value>,
448 name: &str,
449 value: Value,
450) -> Option<String> {
451 insert_decoded_field_with_key_mode(json_obj, name, value, false)
452}
453
454fn insert_decoded_field_with_key_emitting_filler(
455 json_obj: &mut serde_json::Map<String, Value>,
456 name: &str,
457 value: Value,
458) -> Option<String> {
459 insert_decoded_field_with_key_mode(json_obj, name, value, true)
460}
461
462fn insert_decoded_field_with_key_mode(
463 json_obj: &mut serde_json::Map<String, Value>,
464 name: &str,
465 value: Value,
466 allow_filler_collision: bool,
467) -> Option<String> {
468 if !allow_filler_collision
469 && (name.eq_ignore_ascii_case("FILLER") || name.starts_with("_filler_"))
470 {
471 json_obj.insert(name.to_owned(), value);
472 return None;
473 }
474
475 match json_obj.entry(name.to_owned()) {
476 serde_json::map::Entry::Vacant(entry) => {
477 entry.insert(value);
478 return None;
479 }
480 serde_json::map::Entry::Occupied(_) => {}
481 }
482
483 let (base_name, has_duplicate_suffix) = duplicate_name_base(name);
484 let mut candidate = if has_duplicate_suffix && json_obj.contains_key(base_name) {
485 base_name.to_owned()
486 } else {
487 name.to_owned()
488 };
489 let mut suffix = 2;
490 while json_obj.contains_key(&candidate) {
491 candidate = format!("{base_name}__dup{suffix}");
492 suffix += 1;
493 }
494 json_obj.insert(candidate.clone(), value);
495 Some(candidate)
496}
497
498fn duplicate_name_base(name: &str) -> (&str, bool) {
500 let Some((base, suffix)) = name.rsplit_once("__dup") else {
501 return (name, false);
502 };
503 let Ok(number) = suffix.parse::<usize>() else {
504 return (name, false);
505 };
506 if base.is_empty() || number < 2 {
507 return (name, false);
508 }
509 (base, true)
510}
511
512fn process_fields_recursive_with_scratch(
515 fields: &[copybook_core::Field],
516 data: &[u8],
517 json_obj: &mut serde_json::Map<String, Value>,
518 options: &DecodeOptions,
519 scratch: &mut crate::memory::ScratchBuffers,
520 record_index: u64,
521 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
522) -> Result<()> {
523 use copybook_core::FieldKind;
524
525 let mut deferred_group_views = Vec::new();
526
527 for field in fields {
528 if is_filler_field(field) && !options.emit_filler {
529 continue;
530 }
531
532 match (&field.kind, &field.occurs) {
533 (_, Some(occurs)) => {
534 process_array_field_with_scratch(
535 field,
536 occurs,
537 data,
538 json_obj,
539 options,
540 fields,
541 scratch,
542 record_index,
543 encoding_acc,
544 )?;
545 }
546 (FieldKind::Group, None) if field.level > 1 => {
547 let mut group_obj = serde_json::Map::new();
548 let metadata_start = encoding_acc.len();
549 process_fields_recursive_with_scratch(
550 &field.children,
551 data,
552 &mut group_obj,
553 options,
554 scratch,
555 record_index,
556 encoding_acc,
557 )?;
558 if is_scalar_target_group_redefine(field, fields) {
559 let group_value = Value::Object(group_obj);
560 if let Value::Object(group_fields) = &group_value {
561 insert_decoded_group_fields(
562 json_obj,
563 group_fields,
564 &mut encoding_acc[metadata_start..],
565 options.emit_filler,
566 );
567 }
568 deferred_group_views.push((field.name.clone(), group_value));
569 } else if field.redefines_of.is_none() {
570 insert_decoded_field(json_obj, &field.name, Value::Object(group_obj));
571 }
572 }
573 (FieldKind::Group, None) => {
574 process_fields_recursive_with_scratch(
575 &field.children,
576 data,
577 json_obj,
578 options,
579 scratch,
580 record_index,
581 encoding_acc,
582 )?;
583 }
584 _ => {
585 process_scalar_field_with_scratch(
586 field,
587 data,
588 json_obj,
589 options,
590 scratch,
591 record_index,
592 encoding_acc,
593 )?;
594 }
595 }
596 }
597
598 for (name, value) in deferred_group_views {
599 insert_decoded_field(json_obj, &name, value);
600 }
601
602 Ok(())
603}
604
605#[inline]
615#[allow(clippy::too_many_arguments)]
616fn process_scalar_field_standard(
617 field: ©book_core::Field,
618 field_index: usize,
619 total_fields: usize,
620 data: &[u8],
621 json_obj: &mut serde_json::Map<String, Value>,
622 options: &DecodeOptions,
623 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
624 record_index: u64,
625 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
626) -> Result<()> {
627 if matches!(field.kind, copybook_core::FieldKind::Renames { .. }) {
629 let Some(resolved) = &field.resolved_renames else {
630 return Err(Error::new(
631 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
632 format!(
633 "RENAMES field '{name}' has no resolved metadata",
634 name = field.name
635 ),
636 ));
637 };
638
639 let alias_start = resolved.offset as usize;
640 let alias_end = alias_start + resolved.length as usize;
641
642 if alias_end > data.len() {
643 return Err(Error::new(
644 ErrorCode::CBKD301_RECORD_TOO_SHORT,
645 format!(
646 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
647 name = field.name,
648 offset = resolved.offset,
649 length = resolved.length,
650 data_len = data.len()
651 ),
652 ));
653 }
654
655 let alias_data = &data[alias_start..alias_end];
656 let text = crate::charset::ebcdic_to_utf8(
657 alias_data,
658 options.codepage,
659 options.on_decode_unmappable,
660 )?;
661 insert_decoded_field(json_obj, &field.name, Value::String(text));
662 return Ok(());
663 }
664
665 let field_start = field.offset as usize;
666 let mut field_end = field_start + field.len as usize;
667
668 if options.format.is_variable()
669 && field_index + 1 == total_fields
670 && matches!(field.kind, copybook_core::FieldKind::Alphanum { .. })
671 && data.len() > field_end
672 {
673 field_end = data.len();
674 }
675
676 if field_start > data.len() {
677 return Err(Error::new(
678 ErrorCode::CBKD301_RECORD_TOO_SHORT,
679 format!(
680 "Field '{name}' starts beyond record boundary",
681 name = field.name
682 ),
683 ));
684 }
685
686 field_end = field_end.min(data.len());
687
688 if field_start >= field_end {
689 return Ok(());
690 }
691
692 let field_data = &data[field_start..field_end];
693 let value = decode_scalar_field_value_standard(field, field_data, options, scratch_buffers)
694 .map_err(|error| add_zoned_overflow_context(error, field, record_index))?;
695
696 let emitted_key = if options.emit_filler {
697 insert_decoded_field_with_key_emitting_filler(json_obj, &field.name, value)
698 } else {
699 insert_decoded_field_with_key(json_obj, &field.name, value)
700 };
701
702 if options.preserve_zoned_encoding {
704 let metadata_key = emitted_key.as_deref().unwrap_or(&field.name);
705 collect_zoned_encoding_info(field, metadata_key, field_data, options, encoding_acc);
706 }
707
708 if matches!(options.emit_raw, crate::options::RawMode::Field) {
710 let raw_key = emitted_key.map_or_else(
711 || format!("{}_raw_b64", field.name),
712 |key| format!("{key}_raw_b64"),
713 );
714 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(field_data);
715 json_obj.insert(raw_key, Value::String(raw_b64));
716 }
717
718 Ok(())
719}
720
721#[inline]
725fn process_scalar_field_with_scratch(
726 field: ©book_core::Field,
727 data: &[u8],
728 json_obj: &mut serde_json::Map<String, Value>,
729 options: &DecodeOptions,
730 scratch: &mut crate::memory::ScratchBuffers,
731 record_index: u64,
732 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
733) -> Result<()> {
734 if matches!(field.kind, copybook_core::FieldKind::Renames { .. }) {
736 let Some(resolved) = &field.resolved_renames else {
737 return Err(Error::new(
738 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
739 format!(
740 "RENAMES field '{name}' has no resolved metadata",
741 name = field.name
742 ),
743 ));
744 };
745
746 let alias_start = resolved.offset as usize;
747 let alias_end = alias_start + resolved.length as usize;
748
749 if alias_end > data.len() {
750 return Err(Error::new(
751 ErrorCode::CBKD301_RECORD_TOO_SHORT,
752 format!(
753 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
754 name = field.name,
755 offset = resolved.offset,
756 length = resolved.length,
757 data_len = data.len()
758 ),
759 ));
760 }
761
762 let alias_data = &data[alias_start..alias_end];
763 let text = crate::charset::ebcdic_to_utf8(
764 alias_data,
765 options.codepage,
766 options.on_decode_unmappable,
767 )?;
768 insert_decoded_field(json_obj, &field.name, Value::String(text));
769 return Ok(());
770 }
771
772 let field_start = field.offset as usize;
773 let mut field_end = field_start + field.len as usize;
774
775 if field_start > data.len() {
776 return Err(Error::new(
777 ErrorCode::CBKD301_RECORD_TOO_SHORT,
778 format!(
779 "Field '{name}' starts beyond record boundary",
780 name = field.name
781 ),
782 ));
783 }
784
785 if options.format.is_variable() {
786 field_end = field_end.min(data.len());
787 }
788
789 if field_start >= field_end {
790 return Ok(());
791 }
792
793 if field_end > data.len() {
794 return Err(Error::new(
795 ErrorCode::CBKD301_RECORD_TOO_SHORT,
796 format!(
797 "Field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
798 name = field.name,
799 offset = field.offset,
800 length = field.len,
801 data_len = data.len()
802 ),
803 ));
804 }
805
806 let field_data = &data[field_start..field_end];
807 let value = decode_scalar_field_value_with_scratch(field, field_data, options, scratch)
808 .map_err(|error| add_zoned_overflow_context(error, field, record_index))?;
809
810 let emitted_key = if options.emit_filler {
811 insert_decoded_field_with_key_emitting_filler(json_obj, &field.name, value)
812 } else {
813 insert_decoded_field_with_key(json_obj, &field.name, value)
814 };
815
816 if options.preserve_zoned_encoding {
818 let metadata_key = emitted_key.as_deref().unwrap_or(&field.name);
819 collect_zoned_encoding_info(field, metadata_key, field_data, options, encoding_acc);
820 }
821
822 if matches!(options.emit_raw, crate::options::RawMode::Field) {
824 let raw_key = emitted_key.map_or_else(
825 || format!("{}_raw_b64", field.name),
826 |key| format!("{key}_raw_b64"),
827 );
828 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(field_data);
829 json_obj.insert(raw_key, Value::String(raw_b64));
830 }
831
832 Ok(())
833}
834
835#[inline]
836fn add_zoned_overflow_context(
837 error: Error,
838 field: ©book_core::Field,
839 record_index: u64,
840) -> Error {
841 if error.code == ErrorCode::CBKD410_ZONED_OVERFLOW {
842 error
843 .with_record(record_index)
844 .with_field(field.path.clone())
845 .with_offset(u64::from(field.offset))
846 } else {
847 error
848 }
849}
850
851#[allow(clippy::too_many_arguments)]
853fn process_array_field(
854 field: ©book_core::Field,
855 occurs: ©book_core::Occurs,
856 data: &[u8],
857 json_obj: &mut serde_json::Map<String, Value>,
858 options: &DecodeOptions,
859 all_fields: &[copybook_core::Field],
860 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
861 record_index: u64,
862 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
863) -> Result<()> {
864 use copybook_core::FieldKind;
865
866 let count = resolve_array_count(
867 field,
868 occurs,
869 data,
870 options,
871 all_fields,
872 &mut ArrayScratch::Optional(scratch_buffers),
873 record_index,
874 )?;
875
876 let element_size = field.len as usize;
877 let array_start = field.offset as usize;
878 let total_array_size = element_size * count as usize;
879 let array_end = array_start + total_array_size;
880
881 if array_end > data.len() {
883 return Err(Error::new(
884 ErrorCode::CBKD301_RECORD_TOO_SHORT,
885 format!(
886 "Array '{}' requires {} bytes but only {} bytes available",
887 field.name,
888 total_array_size,
889 data.len().saturating_sub(array_start)
890 ),
891 ));
892 }
893
894 let array_metadata_start = encoding_acc.len();
897
898 let mut array_values = Vec::new();
900 let capture_raw = matches!(options.emit_raw, crate::options::RawMode::Field)
901 && !matches!(field.kind, FieldKind::Group | FieldKind::Condition { .. });
902 let mut raw_values = capture_raw.then(|| Vec::with_capacity(count as usize));
903 for i in 0..count {
904 let element_start = array_start + (i as usize * element_size);
905 let element_end = element_start + element_size;
906
907 let element_value = match &field.kind {
908 FieldKind::Group => {
909 let mut element_obj = serde_json::Map::new();
911 let element_base_offset = u32::try_from(element_start).map_err(|_| {
912 Error::new(
913 ErrorCode::CBKD301_RECORD_TOO_SHORT,
914 format!("Array element offset {element_start} exceeds supported range"),
915 )
916 })?;
917 let adjusted_children =
918 adjust_field_offsets(&field.children, element_base_offset, field.offset);
919 process_fields_recursive(
920 &adjusted_children,
921 data,
922 &mut element_obj,
923 options,
924 scratch_buffers,
925 record_index,
926 encoding_acc,
927 )?;
928 Value::Object(element_obj)
929 }
930 FieldKind::Condition { values } => condition_value(values, "CONDITION_ARRAY"),
931 _ => {
932 let element_data = &data[element_start..element_end];
933 if let Some(raw_values) = raw_values.as_mut() {
934 raw_values.push(Value::String(
935 base64::engine::general_purpose::STANDARD.encode(element_data),
936 ));
937 }
938 let val = decode_scalar_field_value_standard(
939 field,
940 element_data,
941 options,
942 scratch_buffers,
943 )
944 .map_err(|error| add_zoned_overflow_context(error, field, record_index))?;
945 if options.preserve_zoned_encoding {
946 collect_array_zoned_encoding_info(field, element_data, options, encoding_acc);
947 }
948 val
949 }
950 };
951
952 array_values.push(element_value);
953 }
954
955 let emitted_key = insert_decoded_array_field(
956 json_obj,
957 field,
958 array_values,
959 encoding_acc,
960 array_metadata_start,
961 );
962 insert_decoded_array_raw_sidecar(json_obj, field, emitted_key, raw_values);
963 Ok(())
964}
965
966#[allow(clippy::too_many_arguments)]
968fn process_array_field_with_scratch(
969 field: ©book_core::Field,
970 occurs: ©book_core::Occurs,
971 data: &[u8],
972 json_obj: &mut serde_json::Map<String, Value>,
973 options: &DecodeOptions,
974 all_fields: &[copybook_core::Field],
975 scratch: &mut crate::memory::ScratchBuffers,
976 record_index: u64,
977 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
978) -> Result<()> {
979 use copybook_core::FieldKind;
980 use serde_json::Value;
981
982 let count = resolve_array_count(
983 field,
984 occurs,
985 data,
986 options,
987 all_fields,
988 &mut ArrayScratch::Direct(scratch),
989 record_index,
990 )?;
991
992 let element_size = field.len as usize;
993 let array_start = field.offset as usize;
994 let total_array_size = element_size * count as usize;
995 let array_end = array_start + total_array_size;
996
997 if array_end > data.len() {
998 return Err(Error::new(
999 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1000 format!(
1001 "Array field '{}' with {} elements at offset {} requires {} bytes but record has {}",
1002 field.name,
1003 count,
1004 array_start,
1005 total_array_size,
1006 data.len() - array_start
1007 ),
1008 ));
1009 }
1010
1011 let array_metadata_start = encoding_acc.len();
1012 let mut array_values = Vec::new();
1013 let capture_raw = matches!(options.emit_raw, crate::options::RawMode::Field)
1014 && !matches!(field.kind, FieldKind::Group | FieldKind::Condition { .. });
1015 let mut raw_values = capture_raw.then(|| Vec::with_capacity(count as usize));
1016
1017 for i in 0..count {
1018 let element_offset = array_start + (i as usize * element_size);
1019 let element_data = &data[element_offset..element_offset + element_size];
1020
1021 let element_value = match &field.kind {
1022 FieldKind::Group => {
1023 let mut group_obj = serde_json::Map::new();
1025
1026 let element_offset_u32 = u32::try_from(element_offset).map_err(|_| {
1030 Error::new(
1031 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1032 format!("Array element offset {element_offset} exceeds supported range"),
1033 )
1034 })?;
1035 let adjusted_children =
1036 adjust_field_offsets(&field.children, element_offset_u32, field.offset);
1037
1038 process_fields_recursive_with_scratch(
1039 &adjusted_children,
1040 data,
1041 &mut group_obj,
1042 options,
1043 scratch,
1044 record_index,
1045 encoding_acc,
1046 )?;
1047 Value::Object(group_obj)
1048 }
1049 FieldKind::Condition { values } => condition_value(values, "CONDITION_ARRAY"),
1050 _ => {
1051 if let Some(raw_values) = raw_values.as_mut() {
1052 raw_values.push(Value::String(
1053 base64::engine::general_purpose::STANDARD.encode(element_data),
1054 ));
1055 }
1056 let val =
1057 decode_scalar_field_value_with_scratch(field, element_data, options, scratch)
1058 .map_err(|error| add_zoned_overflow_context(error, field, record_index))?;
1059 if options.preserve_zoned_encoding {
1060 collect_array_zoned_encoding_info(field, element_data, options, encoding_acc);
1061 }
1062 val
1063 }
1064 };
1065
1066 array_values.push(element_value);
1067 }
1068
1069 let emitted_key = insert_decoded_array_field(
1070 json_obj,
1071 field,
1072 array_values,
1073 encoding_acc,
1074 array_metadata_start,
1075 );
1076 insert_decoded_array_raw_sidecar(json_obj, field, emitted_key, raw_values);
1077 Ok(())
1078}
1079
1080enum ArrayScratch<'a> {
1082 Optional(&'a mut Option<crate::memory::ScratchBuffers>),
1083 Direct(&'a mut crate::memory::ScratchBuffers),
1084}
1085
1086fn resolve_array_count(
1087 field: ©book_core::Field,
1088 occurs: ©book_core::Occurs,
1089 data: &[u8],
1090 options: &DecodeOptions,
1091 all_fields: &[copybook_core::Field],
1092 scratch: &mut ArrayScratch<'_>,
1093 record_index: u64,
1094) -> Result<u32> {
1095 match occurs {
1096 copybook_core::Occurs::Fixed { count } => Ok(*count),
1097 copybook_core::Occurs::ODO {
1098 min,
1099 max,
1100 counter_path,
1101 } => {
1102 let scratch = match scratch {
1103 ArrayScratch::Optional(slot) => {
1104 slot.get_or_insert_with(crate::memory::ScratchBuffers::new)
1105 }
1106 ArrayScratch::Direct(value) => value,
1107 };
1108 let counter_value = find_and_read_counter_field(
1109 counter_path,
1110 all_fields,
1111 data,
1112 options,
1113 scratch,
1114 record_index,
1115 )?;
1116 let counter_field = find_field_by_path(all_fields, counter_path)?;
1117 let validation_context = crate::odo_redefines::OdoValidationContext {
1118 field_path: field.path.clone(),
1119 counter_path: counter_path.clone(),
1120 record_index,
1121 byte_offset: u64::from(counter_field.offset),
1122 };
1123 let validation = crate::odo_redefines::validate_odo_decode(
1124 counter_value,
1125 *min,
1126 *max,
1127 &validation_context,
1128 options,
1129 )?;
1130 if let Some(warning) = validation.warning {
1131 tracing::warn!("{}", warning);
1132 increment_warning_counter();
1133 }
1134 Ok(validation.actual_count)
1135 }
1136 }
1137}
1138
1139fn find_and_read_counter_field(
1141 counter_path: &str,
1142 all_fields: &[copybook_core::Field],
1143 data: &[u8],
1144 options: &DecodeOptions,
1145 scratch: &mut crate::memory::ScratchBuffers,
1146 record_index: u64,
1147) -> Result<u32> {
1148 let counter_field = find_field_by_path(all_fields, counter_path)?;
1150
1151 let field_start = counter_field.offset as usize;
1153 let field_end = field_start + counter_field.len as usize;
1154
1155 if field_end > data.len() {
1156 return Err(Error::new(
1157 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1158 format!("Counter field '{counter_path}' extends beyond record"),
1159 ));
1160 }
1161
1162 let field_data = &data[field_start..field_end];
1163
1164 match &counter_field.kind {
1166 copybook_core::FieldKind::ZonedDecimal {
1167 digits,
1168 scale,
1169 signed,
1170 sign_separate,
1171 } => {
1172 let count = if let Some(sign_sep) = sign_separate {
1173 let decimal = crate::numeric::decode_zoned_decimal_sign_separate(
1174 field_data,
1175 *digits,
1176 *scale,
1177 sign_sep,
1178 options.codepage,
1179 )?;
1180 decimal_counter_to_u32(&decimal, counter_path)?
1181 } else {
1182 let decimal_str = crate::numeric::decode_zoned_decimal_to_string_with_scratch(
1183 field_data,
1184 *digits,
1185 *scale,
1186 *signed,
1187 options.codepage,
1188 counter_field.blank_when_zero,
1189 scratch,
1190 )
1191 .map_err(|error| add_zoned_overflow_context(error, counter_field, record_index))?;
1192 decimal_str.parse::<u32>().map_err(|_| {
1193 Error::new(
1194 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1195 format!("ODO counter '{counter_path}' has invalid value: {decimal_str}"),
1196 )
1197 })?
1198 };
1199
1200 Ok(count)
1201 }
1202 copybook_core::FieldKind::BinaryInt { bits, signed } => {
1203 let int_value = crate::numeric::decode_binary_int(field_data, *bits, *signed)?;
1204 if int_value < 0 {
1205 return Err(Error::new(
1206 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1207 format!("ODO counter '{counter_path}' has negative value: {int_value}"),
1208 ));
1209 }
1210 Ok(u32::try_from(int_value).map_err(|_| {
1211 Error::new(
1212 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1213 format!("ODO counter '{counter_path}' exceeds supported range: {int_value}"),
1214 )
1215 })?)
1216 }
1217 copybook_core::FieldKind::PackedDecimal {
1218 digits,
1219 scale,
1220 signed,
1221 } => {
1222 let decimal_str = crate::numeric::decode_packed_decimal_to_string_with_scratch(
1223 field_data, *digits, *scale, *signed, scratch,
1224 )?;
1225 let count = decimal_str.parse::<u32>().map_err(|_| {
1226 Error::new(
1227 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1228 format!("ODO counter '{counter_path}' has invalid value: {decimal_str}"),
1229 )
1230 })?;
1231 Ok(count)
1232 }
1233 _ => Err(Error::new(
1234 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1235 format!("ODO counter '{counter_path}' has unsupported type"),
1236 )),
1237 }
1238}
1239
1240fn find_field_by_path<'a>(
1242 fields: &'a [copybook_core::Field],
1243 path: &str,
1244) -> Result<&'a copybook_core::Field> {
1245 for field in fields {
1246 if field.path == path || field.name == path {
1247 return Ok(field);
1248 }
1249 if let Ok(found) = find_field_by_path(&field.children, path) {
1251 return Ok(found);
1252 }
1253 }
1254
1255 Err(Error::new(
1256 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1257 format!("ODO counter field '{path}' not found"),
1258 ))
1259}
1260
1261fn adjust_field_offsets(
1267 fields: &[copybook_core::Field],
1268 base_offset: u32,
1269 source_base_offset: u32,
1270) -> Vec<copybook_core::Field> {
1271 fields
1272 .iter()
1273 .map(|field| {
1274 let mut adjusted_field = field.clone();
1275 let relative_offset = field.offset.saturating_sub(source_base_offset);
1276 adjusted_field.offset = base_offset.saturating_add(relative_offset);
1277 if !adjusted_field.children.is_empty() {
1278 adjusted_field.children =
1279 adjust_field_offsets(&adjusted_field.children, base_offset, source_base_offset);
1280 }
1281 adjusted_field
1282 })
1283 .collect()
1284}
1285
1286#[inline]
1288fn is_filler_field(field: ©book_core::Field) -> bool {
1289 field.name.eq_ignore_ascii_case("FILLER") || field.name.starts_with("_filler_")
1290}
1291
1292#[inline]
1297fn collect_zoned_encoding_info(
1298 field: ©book_core::Field,
1299 emitted_key: &str,
1300 field_data: &[u8],
1301 options: &DecodeOptions,
1302 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
1303) {
1304 if let copybook_core::FieldKind::ZonedDecimal { digits, signed, .. } = &field.kind
1305 && let Ok((_, Some(info))) = crate::numeric::decode_zoned_decimal_with_encoding(
1306 field_data,
1307 *digits,
1308 0, *signed,
1310 options.codepage,
1311 field.blank_when_zero,
1312 true,
1313 )
1314 && !info.has_mixed_encoding
1315 {
1316 encoding_acc.push((emitted_key.to_owned(), info.detected_format));
1317 }
1318}
1319
1320fn numeric_string_to_value(s: String, options: &DecodeOptions) -> Value {
1326 use crate::options::JsonNumberMode;
1327 match options.json_number_mode {
1328 JsonNumberMode::Lossless => Value::String(s),
1329 JsonNumberMode::Native => {
1330 if !s.contains('.') && !s.contains('e') && !s.contains('E') {
1332 if let Ok(n) = s.parse::<i64>() {
1333 return Value::Number(serde_json::Number::from(n));
1334 }
1335 if let Ok(n) = s.parse::<u64>() {
1336 return Value::Number(serde_json::Number::from(n));
1337 }
1338 }
1339 if let Ok(f) = s.parse::<f64>()
1341 && let Some(n) = serde_json::Number::from_f64(f)
1342 {
1343 return Value::Number(n);
1344 }
1345 Value::String(s)
1347 }
1348 }
1349}
1350
1351#[allow(clippy::too_many_lines)]
1353fn decode_scalar_field_value_standard(
1354 field: ©book_core::Field,
1355 field_data: &[u8],
1356 options: &DecodeOptions,
1357 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
1358) -> Result<Value> {
1359 use copybook_core::FieldKind;
1360
1361 match &field.kind {
1362 FieldKind::Alphanum { .. } => {
1363 let text = crate::charset::ebcdic_to_utf8(
1364 field_data,
1365 options.codepage,
1366 options.on_decode_unmappable,
1367 )?;
1368 Ok(Value::String(text))
1369 }
1370 FieldKind::ZonedDecimal {
1371 digits,
1372 scale,
1373 signed,
1374 sign_separate,
1375 } => {
1376 if let Some(sign_sep) = sign_separate {
1377 let decimal = crate::numeric::decode_zoned_decimal_sign_separate(
1378 field_data,
1379 *digits,
1380 *scale,
1381 sign_sep,
1382 options.codepage,
1383 )?;
1384 Ok(zoned_decimal_to_json_value(
1385 &decimal,
1386 *digits,
1387 *scale,
1388 field.blank_when_zero,
1389 options,
1390 ))
1391 } else if options.preserve_zoned_encoding {
1392 let (decimal, _encoding_info) = crate::numeric::decode_zoned_decimal_with_encoding(
1394 field_data,
1395 *digits,
1396 *scale,
1397 *signed,
1398 options.codepage,
1399 field.blank_when_zero,
1400 true, )?;
1402
1403 Ok(zoned_decimal_to_json_value(
1406 &decimal,
1407 *digits,
1408 *scale,
1409 field.blank_when_zero,
1410 options,
1411 ))
1412 } else {
1413 let decimal = crate::numeric::decode_zoned_decimal(
1415 field_data,
1416 *digits,
1417 *scale,
1418 *signed,
1419 options.codepage,
1420 field.blank_when_zero,
1421 )?;
1422 Ok(zoned_decimal_to_json_value(
1423 &decimal,
1424 *digits,
1425 *scale,
1426 field.blank_when_zero,
1427 options,
1428 ))
1429 }
1430 }
1431 FieldKind::BinaryInt { bits, signed } => {
1432 let int_value = crate::numeric::decode_binary_int(field_data, *bits, *signed)?;
1433 let scratch = scratch_buffers.get_or_insert_with(crate::memory::ScratchBuffers::new);
1434 let formatted =
1435 crate::numeric::format_binary_int_to_string_with_scratch(int_value, scratch);
1436 Ok(numeric_string_to_value(formatted, options))
1437 }
1438 FieldKind::PackedDecimal {
1439 digits,
1440 scale,
1441 signed,
1442 } => {
1443 let scratch = scratch_buffers.get_or_insert_with(crate::memory::ScratchBuffers::new);
1444 let decimal_str = crate::numeric::decode_packed_decimal_to_string_with_scratch(
1445 field_data, *digits, *scale, *signed, scratch,
1446 )?;
1447 Ok(numeric_string_to_value(decimal_str, options))
1448 }
1449 FieldKind::Group => {
1450 Err(Error::new(
1452 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1453 format!(
1454 "Cannot process group field '{name}' as scalar",
1455 name = field.name
1456 ),
1457 ))
1458 }
1459 FieldKind::Condition { values } => {
1460 Ok(condition_value(values, "CONDITION"))
1463 }
1464 FieldKind::Renames { .. } => {
1465 let Some(resolved) = &field.resolved_renames else {
1467 return Err(Error::new(
1468 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1469 format!(
1470 "RENAMES field '{name}' has no resolved metadata",
1471 name = field.name
1472 ),
1473 ));
1474 };
1475 let alias_start = resolved.offset as usize;
1477 let alias_end = alias_start + resolved.length as usize;
1478
1479 if alias_end > field_data.len() {
1480 return Err(Error::new(
1481 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1482 format!(
1483 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
1484 name = field.name,
1485 offset = resolved.offset,
1486 length = resolved.length,
1487 data_len = field_data.len()
1488 ),
1489 ));
1490 }
1491
1492 if resolved.members.len() == 1 {
1495 let alias_data = &field_data[alias_start..alias_end];
1497 let text = crate::charset::ebcdic_to_utf8(
1499 alias_data,
1500 options.codepage,
1501 options.on_decode_unmappable,
1502 )?;
1503 return Ok(Value::String(text));
1504 }
1505 let alias_data = &field_data[alias_start..alias_end];
1507 let text = crate::charset::ebcdic_to_utf8(
1508 alias_data,
1509 options.codepage,
1510 options.on_decode_unmappable,
1511 )?;
1512 Ok(Value::String(text))
1513 }
1514 FieldKind::EditedNumeric {
1515 pic_string, scale, ..
1516 } => {
1517 let raw_str = crate::charset::ebcdic_to_utf8(
1519 field_data,
1520 options.codepage,
1521 options.on_decode_unmappable,
1522 )?;
1523
1524 let pattern = crate::edited_pic::tokenize_edited_pic(pic_string)?;
1526
1527 let numeric_value = crate::edited_pic::decode_edited_numeric(
1529 &raw_str,
1530 &pattern,
1531 *scale,
1532 field.blank_when_zero,
1533 )?;
1534
1535 Ok(numeric_string_to_value(
1537 numeric_value.to_decimal_string(),
1538 options,
1539 ))
1540 }
1541 FieldKind::FloatSingle => {
1542 let value =
1543 crate::numeric::decode_float_single_with_format(field_data, options.float_format)?;
1544 if value.is_nan() || value.is_infinite() {
1545 Ok(Value::Null)
1546 } else {
1547 Ok(Value::Number(
1548 serde_json::Number::from_f64(f64::from(value))
1549 .unwrap_or_else(|| serde_json::Number::from(0)),
1550 ))
1551 }
1552 }
1553 FieldKind::FloatDouble => {
1554 let value =
1555 crate::numeric::decode_float_double_with_format(field_data, options.float_format)?;
1556 if value.is_nan() || value.is_infinite() {
1557 Ok(Value::Null)
1558 } else {
1559 Ok(Value::Number(
1560 serde_json::Number::from_f64(value)
1561 .unwrap_or_else(|| serde_json::Number::from(0)),
1562 ))
1563 }
1564 }
1565 }
1566}
1567
1568fn is_scalar_target_group_redefine(
1573 field: ©book_core::Field,
1574 sibling_fields: &[copybook_core::Field],
1575) -> bool {
1576 let Some(target_path) = field.redefines_of.as_deref() else {
1577 return false;
1578 };
1579
1580 matches!(field.kind, copybook_core::FieldKind::Group)
1581 && find_field_by_path(sibling_fields, target_path)
1582 .is_ok_and(|target| !matches!(target.kind, copybook_core::FieldKind::Group))
1583}
1584
1585#[allow(clippy::too_many_lines)]
1587fn decode_scalar_field_value_with_scratch(
1588 field: ©book_core::Field,
1589 field_data: &[u8],
1590 options: &DecodeOptions,
1591 scratch: &mut crate::memory::ScratchBuffers,
1592) -> Result<Value> {
1593 use copybook_core::FieldKind;
1594
1595 match &field.kind {
1596 FieldKind::Alphanum { .. } => {
1597 let text = crate::charset::ebcdic_to_utf8(
1598 field_data,
1599 options.codepage,
1600 options.on_decode_unmappable,
1601 )?;
1602 Ok(Value::String(text))
1603 }
1604 FieldKind::ZonedDecimal {
1605 digits,
1606 scale,
1607 signed,
1608 sign_separate,
1609 } => {
1610 if let Some(sign_sep) = sign_separate {
1611 let decimal = crate::numeric::decode_zoned_decimal_sign_separate(
1612 field_data,
1613 *digits,
1614 *scale,
1615 sign_sep,
1616 options.codepage,
1617 )?;
1618 Ok(zoned_decimal_to_json_value(
1619 &decimal,
1620 *digits,
1621 *scale,
1622 field.blank_when_zero,
1623 options,
1624 ))
1625 } else {
1626 let decimal_str = crate::numeric::decode_zoned_decimal_to_string_with_scratch(
1627 field_data,
1628 *digits,
1629 *scale,
1630 *signed,
1631 options.codepage,
1632 field.blank_when_zero,
1633 scratch,
1634 )?;
1635 Ok(numeric_string_to_value(decimal_str, options))
1636 }
1637 }
1638 FieldKind::BinaryInt { bits, signed } => {
1639 let int_value = crate::numeric::decode_binary_int(field_data, *bits, *signed)?;
1640 let formatted =
1641 crate::numeric::format_binary_int_to_string_with_scratch(int_value, scratch);
1642 Ok(numeric_string_to_value(formatted, options))
1643 }
1644 FieldKind::PackedDecimal {
1645 digits,
1646 scale,
1647 signed,
1648 } => {
1649 let decimal_str = crate::numeric::decode_packed_decimal_to_string_with_scratch(
1650 field_data, *digits, *scale, *signed, scratch,
1651 )?;
1652 Ok(numeric_string_to_value(decimal_str, options))
1653 }
1654 FieldKind::Group => Err(Error::new(
1655 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1656 format!(
1657 "Cannot process group field '{name}' as scalar",
1658 name = field.name
1659 ),
1660 )),
1661 FieldKind::Condition { values } => Ok(condition_value(values, "CONDITION")),
1662 FieldKind::Renames { .. } => {
1663 let Some(resolved) = &field.resolved_renames else {
1665 return Err(Error::new(
1666 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1667 format!(
1668 "RENAMES field '{name}' has no resolved metadata",
1669 name = field.name
1670 ),
1671 ));
1672 };
1673 let alias_start = resolved.offset as usize;
1675 let alias_end = alias_start + resolved.length as usize;
1676
1677 if alias_end > field_data.len() {
1678 return Err(Error::new(
1679 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1680 format!(
1681 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
1682 name = field.name,
1683 offset = resolved.offset,
1684 length = resolved.length,
1685 data_len = field_data.len()
1686 ),
1687 ));
1688 }
1689
1690 let alias_data = &field_data[alias_start..alias_end];
1692 let text = crate::charset::ebcdic_to_utf8(
1693 alias_data,
1694 options.codepage,
1695 options.on_decode_unmappable,
1696 )?;
1697 Ok(Value::String(text))
1698 }
1699 FieldKind::EditedNumeric {
1700 pic_string, scale, ..
1701 } => {
1702 let raw_str = crate::charset::ebcdic_to_utf8(
1704 field_data,
1705 options.codepage,
1706 options.on_decode_unmappable,
1707 )?;
1708
1709 let pattern = crate::edited_pic::tokenize_edited_pic(pic_string)?;
1711
1712 let numeric_value = crate::edited_pic::decode_edited_numeric(
1714 &raw_str,
1715 &pattern,
1716 *scale,
1717 field.blank_when_zero,
1718 )?;
1719
1720 Ok(numeric_string_to_value(
1722 numeric_value.to_decimal_string(),
1723 options,
1724 ))
1725 }
1726 FieldKind::FloatSingle => {
1727 let value =
1728 crate::numeric::decode_float_single_with_format(field_data, options.float_format)?;
1729 if value.is_nan() || value.is_infinite() {
1730 Ok(Value::Null)
1731 } else {
1732 Ok(Value::Number(
1733 serde_json::Number::from_f64(f64::from(value))
1734 .unwrap_or_else(|| serde_json::Number::from(0)),
1735 ))
1736 }
1737 }
1738 FieldKind::FloatDouble => {
1739 let value =
1740 crate::numeric::decode_float_double_with_format(field_data, options.float_format)?;
1741 if value.is_nan() || value.is_infinite() {
1742 Ok(Value::Null)
1743 } else {
1744 Ok(Value::Number(
1745 serde_json::Number::from_f64(value)
1746 .unwrap_or_else(|| serde_json::Number::from(0)),
1747 ))
1748 }
1749 }
1750 }
1751}
1752
1753#[inline]
1757fn condition_value(values: &[String], prefix: &str) -> Value {
1758 if values.is_empty() {
1759 Value::String(prefix.to_owned())
1760 } else {
1761 Value::String(format!("{prefix}({})", values.join("|")))
1762 }
1763}
1764
1765#[inline]
1796#[must_use = "Handle the Result or propagate the error"]
1797pub fn encode_record(schema: &Schema, json: &Value, options: &EncodeOptions) -> Result<Vec<u8>> {
1798 let root_obj = json.as_object().ok_or_else(|| {
1799 Error::new(
1800 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1801 "Expected JSON object for record envelope",
1802 )
1803 })?;
1804 let encoding_metadata = root_obj
1805 .get("_encoding_metadata")
1806 .and_then(Value::as_object);
1807 let fields_value = if let Some(fields_val) = root_obj.get("fields") {
1808 fields_val.as_object().ok_or_else(|| {
1809 Error::new(
1810 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1811 "`fields` must be a JSON object",
1812 )
1813 })?;
1814 fields_val
1815 } else {
1816 json
1817 };
1818
1819 if let Some(raw_replay) =
1820 encode_raw_replay(root_obj, fields_value, schema, encoding_metadata, options)?
1821 {
1822 return Ok(raw_replay);
1823 }
1824
1825 validate_lib_api_redefines_encoding(schema, fields_value, options)?;
1827 validate_lib_api_odo_encoding(schema, fields_value, options)?;
1828
1829 match options.format {
1830 RecordFormat::Fixed => {
1831 let payload = encode_fields_to_bytes(schema, fields_value, encoding_metadata, options)?;
1832 Ok(payload)
1833 }
1834 RecordFormat::RDW => {
1835 let payload = encode_fields_to_bytes(schema, fields_value, encoding_metadata, options)?;
1836
1837 let rdw_record = crate::record::RDWRecord::try_new(payload)?;
1839 let mut result = Vec::new();
1840 result.extend_from_slice(&rdw_record.header);
1841 result.extend_from_slice(&rdw_record.payload);
1842 Ok(result)
1843 }
1844 RecordFormat::Vb => {
1845 let payload = encode_fields_to_bytes(schema, fields_value, encoding_metadata, options)?;
1846
1847 let mut block = Vec::new();
1849 let mut writer = crate::record::VbBlockWriter::new(&mut block);
1850 writer.write_record_from_payload(&payload, 0)?;
1851 writer.finish()?;
1852 Ok(block)
1853 }
1854 }
1855}
1856
1857fn parse_raw_capture(root: &serde_json::Map<String, Value>) -> Result<Option<RawCapture>> {
1858 match root.get("raw_capture") {
1859 None => Ok(None),
1860 Some(Value::String(value)) if value == "record" => Ok(Some(RawCapture::Record)),
1861 Some(Value::String(value)) if value == "record+rdw" => Ok(Some(RawCapture::RecordRdw)),
1862 Some(value) => Err(Error::new(
1863 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1864 format!("Invalid raw_capture {value}; expected 'record' or 'record+rdw'"),
1865 )),
1866 }
1867}
1868
1869fn encode_raw_replay(
1870 root: &serde_json::Map<String, Value>,
1871 fields: &Value,
1872 schema: &Schema,
1873 encoding_metadata: Option<&serde_json::Map<String, Value>>,
1874 options: &EncodeOptions,
1875) -> Result<Option<Vec<u8>>> {
1876 if !options.use_raw {
1877 return Ok(None);
1878 }
1879 let Some(raw_str) = root
1880 .get("raw_b64")
1881 .or_else(|| root.get("__raw_b64"))
1882 .and_then(Value::as_str)
1883 else {
1884 return Ok(None);
1885 };
1886 let capture = parse_raw_capture(root)?;
1887 let raw_data = base64::engine::general_purpose::STANDARD
1888 .decode(raw_str)
1889 .map_err(|error| {
1890 Error::new(
1891 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1892 format!("Invalid base64 in raw_b64: {error}"),
1893 )
1894 })?;
1895
1896 match options.format {
1897 RecordFormat::Fixed => encode_fixed_raw_replay(raw_data, capture),
1898 RecordFormat::RDW => encode_rdw_raw_replay(
1899 raw_data,
1900 capture,
1901 fields,
1902 schema,
1903 encoding_metadata,
1904 options,
1905 ),
1906 RecordFormat::Vb => encode_vb_raw_replay(raw_data, capture),
1907 }
1908 .map(Some)
1909}
1910
1911fn encode_vb_raw_replay(raw_data: Vec<u8>, capture: Option<RawCapture>) -> Result<Vec<u8>> {
1917 let framed = if matches!(capture, Some(RawCapture::Record)) {
1918 let rdw_len = raw_data.len() + crate::record::RDW_HEADER_LEN;
1920 let rdw_len = u16::try_from(rdw_len).map_err(|_| {
1921 Error::new(
1922 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
1923 format!("VB raw replay payload too large: {rdw_len} bytes"),
1924 )
1925 })?;
1926 let mut framed = Vec::with_capacity(rdw_len as usize);
1927 framed.extend_from_slice(&rdw_len.to_be_bytes());
1928 framed.extend_from_slice(&[0, 0]);
1929 framed.extend_from_slice(&raw_data);
1930 framed
1931 } else {
1932 parse_vb_raw_rdw_frame(&raw_data)?;
1937 raw_data
1938 };
1939 let block_len = framed.len() + crate::record::BDW_HEADER_LEN;
1940 let header = crate::record::BdwHeader::from_block_len(block_len)?;
1941 let mut block = Vec::with_capacity(block_len);
1942 block.extend_from_slice(&header.bytes());
1943 block.extend_from_slice(&framed);
1944 Ok(block)
1945}
1946
1947fn parse_vb_raw_rdw_frame(frame: &[u8]) -> Result<(u16, &[u8])> {
1957 let (raw_header, raw_payload) = frame.split_at_checked(4).ok_or_else(|| {
1958 Error::new(
1959 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
1960 format!(
1961 "Raw VB record is {} bytes; expected at least a 4-byte RDW header",
1962 frame.len()
1963 ),
1964 )
1965 })?;
1966 let header_bytes: [u8; 4] = raw_header.try_into().map_err(|_| {
1967 Error::new(
1968 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
1969 "Raw VB record does not contain a complete 4-byte RDW header",
1970 )
1971 })?;
1972 let declared_len = usize::from(u16::from_be_bytes([header_bytes[0], header_bytes[1]]));
1973 if declared_len < crate::record::RDW_HEADER_LEN
1974 || declared_len != raw_payload.len() + crate::record::RDW_HEADER_LEN
1975 {
1976 return Err(Error::new(
1977 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
1978 format!(
1979 "Raw VB RDW header declares {declared_len} bytes (inclusive), but {} payload bytes follow",
1980 raw_payload.len()
1981 ),
1982 ));
1983 }
1984 let reserved = u16::from_be_bytes([header_bytes[2], header_bytes[3]]);
1985 Ok((reserved, raw_payload))
1986}
1987
1988fn encode_fixed_raw_replay(raw_data: Vec<u8>, capture: Option<RawCapture>) -> Result<Vec<u8>> {
1989 if matches!(capture, Some(RawCapture::RecordRdw)) {
1990 return Err(Error::new(
1991 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1992 "raw_capture 'record+rdw' conflicts with fixed record format",
1993 ));
1994 }
1995 Ok(raw_data)
1996}
1997
1998fn encode_rdw_raw_replay(
1999 raw_data: Vec<u8>,
2000 capture: Option<RawCapture>,
2001 fields: &Value,
2002 schema: &Schema,
2003 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2004 options: &EncodeOptions,
2005) -> Result<Vec<u8>> {
2006 if matches!(capture, Some(RawCapture::Record)) {
2007 return Ok(crate::record::RDWRecord::try_with_reserved(raw_data, 0)?.as_bytes());
2008 }
2009
2010 let (reserved, raw_payload) = parse_raw_rdw_frame(&raw_data)?;
2013 let field_payload = encode_fields_to_bytes(schema, fields, encoding_metadata, options)?;
2014 if field_payload == raw_payload {
2015 return Ok(raw_data);
2016 }
2017 Ok(crate::record::RDWRecord::try_with_reserved(field_payload, reserved)?.as_bytes())
2018}
2019
2020fn validate_lib_api_redefines_encoding(
2022 schema: &Schema,
2023 json_value: &Value,
2024 options: &EncodeOptions,
2025) -> Result<()> {
2026 let redefines_context = crate::odo_redefines::build_redefines_context(schema, json_value);
2027
2028 for (cluster_path, non_null_views) in &redefines_context.cluster_views {
2029 let field_path = non_null_views
2030 .first()
2031 .cloned()
2032 .unwrap_or_else(|| cluster_path.clone());
2033
2034 let byte_offset = non_null_views
2035 .iter()
2036 .find_map(|view| schema.find_field(view).map(|field| u64::from(field.offset)))
2037 .or_else(|| {
2038 schema
2039 .find_field(cluster_path)
2040 .map(|field| u64::from(field.offset))
2041 })
2042 .unwrap_or(0);
2043
2044 crate::odo_redefines::validate_redefines_encoding(
2045 &redefines_context,
2046 cluster_path,
2047 &field_path,
2048 json_value,
2049 options.use_raw,
2050 0,
2051 byte_offset,
2052 )?;
2053 }
2054
2055 Ok(())
2056}
2057
2058fn validate_lib_api_odo_encoding(
2060 schema: &Schema,
2061 json_value: &Value,
2062 options: &EncodeOptions,
2063) -> Result<()> {
2064 let Some(tail_odo) = &schema.tail_odo else {
2065 return Ok(());
2066 };
2067
2068 let fields_value = if let Some(fields_value) = json_value.get("fields") {
2069 fields_value
2070 } else {
2071 json_value
2072 };
2073
2074 let has_wrapper = json_value.get("fields").is_some();
2075 if !fields_value.is_object() {
2076 if has_wrapper {
2077 return Err(Error::new(
2078 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2079 "`fields` must be a JSON object",
2080 ));
2081 }
2082 return Ok(());
2083 }
2084
2085 let array_field =
2086 crate::odo_redefines::find_field_by_path_or_unique_name(schema, &tail_odo.array_path)
2087 .ok_or_else(|| {
2088 Error::new(
2089 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
2090 format!(
2091 "ODO array field '{}' not found in schema",
2092 tail_odo.array_path
2093 ),
2094 )
2095 .with_context(
2096 crate::odo_redefines::create_comprehensive_error_context(
2097 0,
2098 &tail_odo.array_path,
2099 0,
2100 None,
2101 ),
2102 )
2103 })?;
2104
2105 let counter_field =
2106 crate::odo_redefines::find_field_by_path_or_unique_name(schema, &tail_odo.counter_path)
2107 .ok_or_else(|| {
2108 crate::odo_redefines::handle_missing_counter_field(
2109 &tail_odo.counter_path,
2110 &tail_odo.array_path,
2111 schema,
2112 0,
2113 0,
2114 )
2115 })?;
2116
2117 if let Some(array) = json_lookup_array(fields_value, &array_field.path)
2118 .or_else(|| json_lookup_array(fields_value, &tail_odo.array_path))
2119 {
2120 let Some(counter_json_value) = json_lookup_value(fields_value, &counter_field.path)
2121 .or_else(|| json_lookup_value(fields_value, &tail_odo.counter_path))
2122 else {
2123 return Err(crate::odo_redefines::handle_missing_counter_field(
2124 &counter_field.path,
2125 &array_field.path,
2126 schema,
2127 0,
2128 u64::from(counter_field.offset),
2129 ));
2130 };
2131
2132 if let Some(counter_count) = json_counter_value_as_usize(counter_json_value)
2138 && counter_count != array.len()
2139 {
2140 return Err(Error::new(
2141 ErrorCode::CBKE521_ARRAY_LEN_OOB,
2142 format!(
2143 "ODO counter '{}' value ({counter_count}) does not match array '{}' length ({})",
2144 counter_field.path,
2145 array_field.path,
2146 array.len()
2147 ),
2148 )
2149 .with_context(crate::odo_redefines::create_comprehensive_error_context(
2150 0,
2151 &array_field.path,
2152 u64::from(array_field.offset),
2153 Some(format!(
2154 "counter_field={}, counter_value={counter_count}, array_length={}",
2155 counter_field.path,
2156 array.len()
2157 )),
2158 )));
2159 }
2160
2161 let context = crate::odo_redefines::OdoValidationContext {
2162 field_path: array_field.path.clone(),
2163 counter_path: counter_field.path.clone(),
2164 record_index: 0,
2165 byte_offset: u64::from(array_field.offset),
2166 };
2167
2168 crate::odo_redefines::validate_odo_encode(
2169 array.len(),
2170 tail_odo.min_count,
2171 tail_odo.max_count,
2172 &context,
2173 options,
2174 )?;
2175 }
2176
2177 Ok(())
2178}
2179
2180fn json_counter_value_as_usize(value: &Value) -> Option<usize> {
2182 match value {
2183 Value::Number(n) => n.as_u64().and_then(|v| usize::try_from(v).ok()),
2184 Value::String(s) => s.trim().parse::<usize>().ok(),
2185 _ => None,
2186 }
2187}
2188
2189fn json_lookup_value<'a>(value: &'a Value, field_path: &str) -> Option<&'a Value> {
2190 json_lookup_exact_value(value, field_path).or_else(|| {
2191 let (_, path_without_root) = field_path.split_once('.')?;
2192 json_lookup_exact_value(value, path_without_root)
2193 })
2194}
2195
2196fn json_lookup_exact_value<'a>(value: &'a Value, field_path: &str) -> Option<&'a Value> {
2197 let mut current = value;
2198 for segment in field_path.split('.') {
2199 current = current.as_object()?.get(segment)?;
2200 }
2201 Some(current)
2202}
2203
2204fn json_lookup_array<'a>(value: &'a Value, field_path: &str) -> Option<&'a Vec<Value>> {
2205 let leaf = field_path.split('.').next_back().unwrap_or("");
2206 match json_lookup_value(value, field_path) {
2207 Some(Value::Array(array)) => Some(array),
2208 _ => {
2209 if let Value::Object(obj) = value {
2210 obj.get(leaf).and_then(|candidate| candidate.as_array())
2211 } else {
2212 None
2213 }
2214 }
2215 }
2216}
2217
2218fn encode_fields_to_bytes(
2220 schema: &Schema,
2221 json: &Value,
2222 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2223 options: &EncodeOptions,
2224) -> Result<Vec<u8>> {
2225 let maximum_record_length = schema.lrecl_fixed.unwrap_or_else(|| {
2226 schema.fields.iter().map(|f| f.len).sum::<u32>()
2228 }) as usize;
2229 let record_length = if options.format.is_variable() {
2230 rdw_record_length_for_json(schema, json).unwrap_or(maximum_record_length)
2231 } else {
2232 maximum_record_length
2233 };
2234
2235 let mut buffer = vec![0u8; record_length];
2236
2237 if let Some(obj) = json.as_object() {
2238 encode_fields_recursive(
2239 &schema.fields,
2240 obj,
2241 "",
2242 &EncodeFieldsContext {
2243 encoding_metadata,
2244 flattened: false,
2245 flattened_prior_names: None,
2246 },
2247 &mut buffer,
2248 0,
2249 options,
2250 )?;
2251 }
2252
2253 Ok(buffer)
2254}
2255
2256fn rdw_record_length_for_json(schema: &Schema, json: &Value) -> Option<usize> {
2263 let array_field = schema
2264 .all_fields()
2265 .into_iter()
2266 .find(|field| matches!(field.occurs, Some(copybook_core::Occurs::ODO { .. })))?;
2267 let Some(copybook_core::Occurs::ODO { max, .. }) = array_field.occurs else {
2268 return None;
2269 };
2270 let field_offset = usize::try_from(array_field.offset).ok()?;
2271 let field_length = usize::try_from(array_field.len).ok()?;
2272 let maximum_array_end = array_field
2273 .offset
2274 .checked_add(array_field.len.checked_mul(max)?)?;
2275 let maximum_array_end = usize::try_from(maximum_array_end).ok()?;
2276 let schema_length = usize::try_from(schema.lrecl_fixed?).ok()?;
2277 let trailing_length = schema_length.checked_sub(maximum_array_end)?;
2278 let array = json_lookup_array(json, &array_field.path)
2279 .or_else(|| json_lookup_array(json, &array_field.name))?;
2280 let count = array.len();
2281 let array_bytes = field_length.checked_mul(count)?;
2282 field_offset
2283 .checked_add(array_bytes)?
2284 .checked_add(trailing_length)
2285}
2286
2287fn encode_fields_recursive(
2289 fields: &[copybook_core::Field],
2290 json_obj: &serde_json::Map<String, Value>,
2291 path_prefix: &str,
2292 context: &EncodeFieldsContext<'_>,
2293 buffer: &mut [u8],
2294 offset: usize,
2295 options: &EncodeOptions,
2296) -> Result<usize> {
2297 let mut current_offset = offset;
2298 let mut name_occurrences = HashMap::new();
2299
2300 for field in fields {
2301 let occurrence = name_occurrences
2302 .get(field.name.as_str())
2303 .copied()
2304 .unwrap_or(0);
2305 let json_field_name = emitted_field_name(
2306 json_obj,
2307 &field.name,
2308 occurrence,
2309 context.flattened,
2310 context.flattened_prior_names,
2311 );
2312 let field_path = if path_prefix.is_empty() {
2313 field.name.clone()
2314 } else {
2315 format!("{path_prefix}.{}", field.name)
2316 };
2317
2318 let field_names = FieldNames {
2319 path: &field_path,
2320 json: &json_field_name,
2321 prior_name_occurrences: &name_occurrences,
2322 };
2323 let field_offset = if field.redefines_of.is_some() {
2324 match usize::try_from(field.offset) {
2325 Ok(offset) => offset,
2326 Err(_) => current_offset,
2327 }
2328 } else {
2329 current_offset
2330 };
2331 let encoded_offset = encode_single_field(
2332 field,
2333 &field_names,
2334 json_obj,
2335 context.encoding_metadata,
2336 buffer,
2337 field_offset,
2338 options,
2339 )?;
2340 current_offset = if field.redefines_of.is_some() {
2341 current_offset.max(encoded_offset)
2342 } else {
2343 encoded_offset
2344 };
2345
2346 let emitted_group = json_obj.contains_key(&field.name);
2347 if field.redefines_of.is_none() || emitted_group {
2348 name_occurrences.insert(field.name.as_str(), occurrence + 1);
2349 }
2350
2351 if matches!(field.kind, copybook_core::FieldKind::Group)
2352 && field.redefines_of.is_some()
2353 && !json_obj.contains_key(&field.name)
2354 {
2355 for child in &field.children {
2356 if matches!(child.kind, copybook_core::FieldKind::Group)
2357 && child.redefines_of.is_some()
2358 && !json_obj.contains_key(&child.name)
2359 {
2360 continue;
2361 }
2362 let has_emitted_child = json_obj.contains_key(&child.name)
2363 || json_obj.contains_key(&format!("{}__dup2", child.name));
2364 if has_emitted_child {
2365 name_occurrences
2366 .entry(child.name.as_str())
2367 .and_modify(|count| *count += 1)
2368 .or_insert(1);
2369 }
2370 }
2371 }
2372 }
2373
2374 Ok(current_offset)
2375}
2376
2377#[inline]
2378fn collect_array_zoned_encoding_info(
2379 field: ©book_core::Field,
2380 field_data: &[u8],
2381 options: &DecodeOptions,
2382 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
2383) {
2384 collect_zoned_encoding_info(field, &field.name, field_data, options, encoding_acc);
2385}
2386
2387fn insert_decoded_array_field(
2388 json_obj: &mut serde_json::Map<String, Value>,
2389 field: ©book_core::Field,
2390 array_values: Vec<Value>,
2391 encoding_acc: &mut [(String, ZonedEncodingFormat)],
2392 metadata_start: usize,
2393) -> Option<String> {
2394 let emitted_key =
2395 insert_decoded_field_with_key(json_obj, &field.name, Value::Array(array_values));
2396 finalize_array_zoned_metadata(
2397 &field.kind,
2398 encoding_acc,
2399 metadata_start,
2400 emitted_key.as_deref().unwrap_or(&field.name),
2401 );
2402 emitted_key
2403}
2404
2405fn insert_decoded_array_raw_sidecar(
2406 json_obj: &mut serde_json::Map<String, Value>,
2407 field: ©book_core::Field,
2408 emitted_key: Option<String>,
2409 raw_values: Option<Vec<Value>>,
2410) {
2411 let Some(raw_values) = raw_values else {
2412 return;
2413 };
2414 let raw_key = emitted_key.map_or_else(
2415 || format!("{}_raw_b64", field.name),
2416 |key| format!("{key}_raw_b64"),
2417 );
2418 json_obj.insert(raw_key, Value::Array(raw_values));
2419}
2420
2421fn finalize_array_zoned_metadata(
2422 field_kind: ©book_core::FieldKind,
2423 encoding_acc: &mut [(String, ZonedEncodingFormat)],
2424 metadata_start: usize,
2425 emitted_key: &str,
2426) {
2427 if !matches!(field_kind, copybook_core::FieldKind::ZonedDecimal { .. }) {
2428 return;
2429 }
2430 let metadata_key = emitted_key.to_owned();
2431 for (key, _) in &mut encoding_acc[metadata_start..] {
2432 key.clone_from(&metadata_key);
2433 }
2434}
2435
2436struct FieldNames<'a> {
2437 path: &'a str,
2438 json: &'a str,
2439 prior_name_occurrences: &'a HashMap<&'a str, usize>,
2440}
2441
2442struct EncodeFieldsContext<'a> {
2443 encoding_metadata: Option<&'a serde_json::Map<String, Value>>,
2444 flattened: bool,
2445 flattened_prior_names: Option<&'a HashMap<&'a str, usize>>,
2446}
2447
2448fn emitted_field_name(
2449 json_obj: &serde_json::Map<String, Value>,
2450 field_name: &str,
2451 occurrence: usize,
2452 flattened: bool,
2453 flattened_prior_names: Option<&HashMap<&str, usize>>,
2454) -> String {
2455 let candidate = if occurrence == 0 {
2456 field_name.to_owned()
2457 } else {
2458 format!("{field_name}__dup{}", occurrence + 1)
2459 };
2460 if flattened
2461 && occurrence == 0
2462 && json_obj.contains_key(&candidate)
2463 && flattened_prior_names.is_some_and(|names| names.contains_key(field_name))
2464 {
2465 let duplicate = format!("{field_name}__dup2");
2466 if json_obj.contains_key(&duplicate) {
2467 duplicate
2468 } else {
2469 candidate
2470 }
2471 } else if json_obj.contains_key(&candidate) {
2472 candidate
2473 } else {
2474 field_name.to_owned()
2475 }
2476}
2477
2478#[inline]
2483#[allow(clippy::too_many_lines)]
2484fn encode_single_field(
2485 field: ©book_core::Field,
2486 field_names: &FieldNames<'_>,
2487 json_obj: &serde_json::Map<String, Value>,
2488 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2489 buffer: &mut [u8],
2490 current_offset: usize,
2491 options: &EncodeOptions,
2492) -> Result<usize> {
2493 use copybook_core::FieldKind;
2494
2495 if let Some(occurs) = &field.occurs {
2496 return encode_occurs_field(
2497 field,
2498 occurs,
2499 field_names,
2500 json_obj,
2501 encoding_metadata,
2502 buffer,
2503 current_offset,
2504 options,
2505 );
2506 }
2507
2508 match &field.kind {
2509 FieldKind::Group => encode_group_field(
2510 field,
2511 field_names,
2512 json_obj,
2513 encoding_metadata,
2514 buffer,
2515 current_offset,
2516 options,
2517 ),
2518 FieldKind::Alphanum { .. } => encode_alphanum_field(
2519 field,
2520 field_names.json,
2521 json_obj,
2522 buffer,
2523 current_offset,
2524 options,
2525 ),
2526 FieldKind::ZonedDecimal {
2527 digits,
2528 scale,
2529 signed,
2530 sign_separate,
2531 } => {
2532 if let Some(sign_sep) = sign_separate {
2533 let field_len = field.len as usize;
2534 if let Some(text) = json_obj.get(field_names.json).and_then(|v| v.as_str()) {
2535 crate::numeric::encode_zoned_decimal_sign_separate(
2536 text,
2537 *digits,
2538 *scale,
2539 sign_sep,
2540 options.codepage,
2541 &mut buffer[current_offset..current_offset + field_len],
2542 )?;
2543 }
2544 Ok(current_offset + field_len)
2545 } else {
2546 encode_zoned_decimal_field(
2547 field,
2548 field_names.path,
2549 field_names.json,
2550 json_obj,
2551 encoding_metadata,
2552 buffer,
2553 current_offset,
2554 options,
2555 DecimalSpec {
2556 digits: *digits,
2557 scale: *scale,
2558 signed: *signed,
2559 },
2560 )
2561 }
2562 }
2563 FieldKind::PackedDecimal {
2564 digits,
2565 scale,
2566 signed,
2567 } => encode_packed_decimal_field(
2568 field,
2569 field_names.json,
2570 json_obj,
2571 buffer,
2572 current_offset,
2573 options,
2574 DecimalSpec {
2575 digits: *digits,
2576 scale: *scale,
2577 signed: *signed,
2578 },
2579 ),
2580 FieldKind::BinaryInt { bits, signed } => encode_binary_int_field(
2581 field,
2582 field_names.json,
2583 json_obj,
2584 buffer,
2585 current_offset,
2586 options,
2587 BinarySpec {
2588 bits: *bits,
2589 signed: *signed,
2590 },
2591 ),
2592 FieldKind::Condition { .. } => Ok(current_offset),
2593 FieldKind::Renames { .. } => {
2594 Ok(current_offset)
2598 }
2599 FieldKind::EditedNumeric {
2600 pic_string, scale, ..
2601 } => {
2602 if let Some(text) = encodable_numeric_text(
2604 json_obj,
2605 field,
2606 &field.name,
2607 "an edited numeric string",
2608 options.coerce_numbers,
2609 )? {
2610 let pattern = crate::edited_pic::tokenize_edited_pic(pic_string)?;
2612
2613 let encoded = crate::edited_pic::encode_edited_numeric(
2615 &text,
2616 &pattern,
2617 *scale,
2618 field.blank_when_zero,
2619 )?;
2620
2621 let bytes = crate::charset::utf8_to_ebcdic(&encoded, options.codepage)?;
2623 let field_len = field.len as usize;
2624 let copy_len = bytes.len().min(field_len);
2625
2626 if current_offset + field_len <= buffer.len() {
2627 buffer[current_offset..current_offset + copy_len]
2628 .copy_from_slice(&bytes[..copy_len]);
2629 let space = crate::charset::space_byte(options.codepage);
2631 buffer[current_offset + copy_len..current_offset + field_len].fill(space);
2632 }
2633 }
2634 Ok(current_offset + field.len as usize)
2635 }
2636 FieldKind::FloatSingle => {
2637 let field_len = field.len as usize;
2638 if let Some(val) = json_obj.get(&field.name) {
2639 let f = match val {
2640 Value::Number(n) => {
2641 let f64_val = n.as_f64().unwrap_or(0.0);
2642 if f64_val.is_finite()
2644 && (f64_val > f64::from(f32::MAX) || f64_val < f64::from(f32::MIN))
2645 {
2646 return Err(Error::new(
2647 ErrorCode::CBKE531_FLOAT_ENCODE_OVERFLOW,
2648 format!("Value overflow for COMP-1 field '{}'", field.name),
2649 ));
2650 }
2651 #[allow(clippy::cast_possible_truncation)]
2653 {
2654 f64_val as f32
2655 }
2656 }
2657 Value::String(s) => s.parse::<f32>().map_err(|e| {
2658 Error::new(
2659 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2660 format!(
2661 "Cannot parse '{}' as f32 for field '{}': {}",
2662 s, field.name, e
2663 ),
2664 )
2665 })?,
2666 Value::Null => f32::NAN,
2667 _ => {
2668 return Err(Error::new(
2669 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2670 format!("Expected number for COMP-1 field '{}'", field.name),
2671 ));
2672 }
2673 };
2674 if current_offset + field_len <= buffer.len() {
2675 crate::numeric::encode_float_single_with_format(
2676 f,
2677 &mut buffer[current_offset..current_offset + field_len],
2678 options.float_format,
2679 )?;
2680 }
2681 }
2682 Ok(current_offset + field_len)
2683 }
2684 FieldKind::FloatDouble => {
2685 let field_len = field.len as usize;
2686 if let Some(val) = json_obj.get(&field.name) {
2687 let f = match val {
2688 Value::Number(n) => n.as_f64().unwrap_or(0.0),
2689 Value::String(s) => s.parse::<f64>().map_err(|e| {
2690 Error::new(
2691 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2692 format!(
2693 "Cannot parse '{}' as f64 for field '{}': {}",
2694 s, field.name, e
2695 ),
2696 )
2697 })?,
2698 Value::Null => f64::NAN,
2699 _ => {
2700 return Err(Error::new(
2701 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2702 format!("Expected number for COMP-2 field '{}'", field.name),
2703 ));
2704 }
2705 };
2706 if current_offset + field_len <= buffer.len() {
2707 crate::numeric::encode_float_double_with_format(
2708 f,
2709 &mut buffer[current_offset..current_offset + field_len],
2710 options.float_format,
2711 )?;
2712 }
2713 }
2714 Ok(current_offset + field_len)
2715 }
2716 }
2717}
2718
2719#[allow(clippy::too_many_arguments)]
2720fn encode_occurs_field(
2721 field: ©book_core::Field,
2722 occurs: ©book_core::Occurs,
2723 field_names: &FieldNames<'_>,
2724 json_obj: &serde_json::Map<String, Value>,
2725 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2726 buffer: &mut [u8],
2727 current_offset: usize,
2728 options: &EncodeOptions,
2729) -> Result<usize> {
2730 let max_count = occurs_max_count(occurs);
2731 let element_len = field.len as usize;
2732 let allocation_len = element_len
2733 .checked_mul(max_count as usize)
2734 .ok_or_else(|| Error::new(ErrorCode::CBKS141_RECORD_TOO_LARGE, "OCCURS size overflow"))?;
2735
2736 let Some(array) = json_obj.get(field_names.json).and_then(Value::as_array) else {
2737 return Ok(current_offset + allocation_len);
2738 };
2739
2740 validate_occurs_array_len(array.len(), occurs, field)?;
2741
2742 for (index, element) in array.iter().enumerate() {
2743 let element_offset = current_offset + index * element_len;
2744 encode_occurs_element(
2745 field,
2746 field_names,
2747 element,
2748 encoding_metadata,
2749 buffer,
2750 element_offset,
2751 options,
2752 )?;
2753 }
2754
2755 Ok(current_offset + allocation_len)
2756}
2757
2758fn occurs_max_count(occurs: ©book_core::Occurs) -> u32 {
2759 match occurs {
2760 copybook_core::Occurs::Fixed { count } => *count,
2761 copybook_core::Occurs::ODO { max, .. } => *max,
2762 }
2763}
2764
2765fn validate_occurs_array_len(
2766 actual_len: usize,
2767 occurs: ©book_core::Occurs,
2768 field: ©book_core::Field,
2769) -> Result<()> {
2770 match occurs {
2771 copybook_core::Occurs::Fixed { count } if actual_len != *count as usize => Err(Error::new(
2772 ErrorCode::CBKE521_ARRAY_LEN_OOB,
2773 format!(
2774 "Array length {} doesn't match fixed OCCURS count {} for field '{}'",
2775 actual_len, count, field.path
2776 ),
2777 )
2778 .with_field(field.path.clone())),
2779 copybook_core::Occurs::ODO { max, .. } if actual_len > *max as usize => Err(Error::new(
2780 ErrorCode::CBKE521_ARRAY_LEN_OOB,
2781 format!(
2782 "Array length {} exceeds ODO max {} for field '{}'",
2783 actual_len, max, field.path
2784 ),
2785 )
2786 .with_field(field.path.clone())),
2787 copybook_core::Occurs::Fixed { .. } | copybook_core::Occurs::ODO { .. } => Ok(()),
2788 }
2789}
2790
2791fn encode_occurs_element(
2792 field: ©book_core::Field,
2793 field_names: &FieldNames<'_>,
2794 element: &Value,
2795 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2796 buffer: &mut [u8],
2797 element_offset: usize,
2798 options: &EncodeOptions,
2799) -> Result<()> {
2800 use copybook_core::FieldKind;
2801
2802 let mut element_field = field.clone();
2803 element_field.occurs = None;
2804
2805 if let FieldKind::Group = &field.kind {
2806 let element_obj = element.as_object().ok_or_else(|| {
2807 Error::new(
2808 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2809 format!("Expected object element for OCCURS group '{}'", field.path),
2810 )
2811 .with_field(field.path.clone())
2812 })?;
2813 encode_fields_recursive(
2814 &element_field.children,
2815 element_obj,
2816 field_names.path,
2817 &EncodeFieldsContext {
2818 encoding_metadata,
2819 flattened: false,
2820 flattened_prior_names: None,
2821 },
2822 buffer,
2823 element_offset,
2824 options,
2825 )?;
2826 } else {
2827 let mut element_obj = serde_json::Map::new();
2828 element_obj.insert(field_names.json.to_owned(), element.clone());
2829 encode_single_field(
2830 &element_field,
2831 field_names,
2832 &element_obj,
2833 encoding_metadata,
2834 buffer,
2835 element_offset,
2836 options,
2837 )?;
2838 }
2839
2840 Ok(())
2841}
2842
2843#[inline]
2845fn encode_group_field(
2846 field: ©book_core::Field,
2847 field_names: &FieldNames<'_>,
2848 json_obj: &serde_json::Map<String, Value>,
2849 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2850 buffer: &mut [u8],
2851 current_offset: usize,
2852 options: &EncodeOptions,
2853) -> Result<usize> {
2854 if let Some(sub_obj) = json_obj.get(field_names.json).and_then(|v| v.as_object()) {
2855 encode_fields_recursive(
2856 &field.children,
2857 sub_obj,
2858 field_names.path,
2859 &EncodeFieldsContext {
2860 encoding_metadata,
2861 flattened: false,
2862 flattened_prior_names: None,
2863 },
2864 buffer,
2865 current_offset,
2866 options,
2867 )
2868 } else {
2869 encode_fields_recursive(
2870 &field.children,
2871 json_obj,
2872 field_names.path,
2873 &EncodeFieldsContext {
2874 encoding_metadata,
2875 flattened: field.redefines_of.is_some(),
2876 flattened_prior_names: field
2877 .redefines_of
2878 .is_some()
2879 .then_some(field_names.prior_name_occurrences),
2880 },
2881 buffer,
2882 current_offset,
2883 options,
2884 )
2885 }
2886}
2887
2888#[inline]
2890fn encode_alphanum_field(
2891 field: ©book_core::Field,
2892 json_field_name: &str,
2893 json_obj: &serde_json::Map<String, Value>,
2894 buffer: &mut [u8],
2895 current_offset: usize,
2896 options: &EncodeOptions,
2897) -> Result<usize> {
2898 let field_len = field.len as usize;
2899
2900 if let Some(text) = json_obj
2901 .get(json_field_name)
2902 .and_then(|value| value.as_str())
2903 {
2904 let bytes = crate::charset::utf8_to_ebcdic(text, options.codepage)?;
2906 if bytes.len() > field_len {
2907 return Err(Error::new(
2908 ErrorCode::CBKE515_STRING_LENGTH_VIOLATION,
2909 format!(
2910 "Encoded byte length {} exceeds field capacity {} for alphanumeric field {}",
2911 bytes.len(),
2912 field_len,
2913 field.path
2914 ),
2915 )
2916 .with_field(field.path.clone()));
2917 }
2918
2919 let copy_len = bytes.len();
2920
2921 if current_offset + field_len <= buffer.len() {
2922 buffer[current_offset..current_offset + copy_len].copy_from_slice(&bytes);
2923 let space = crate::charset::space_byte(options.codepage);
2925 buffer[current_offset + copy_len..current_offset + field_len].fill(space);
2926 }
2927 }
2928
2929 Ok(current_offset + field_len)
2930}
2931
2932#[derive(Copy, Clone)]
2933struct DecimalSpec {
2934 digits: u16,
2935 scale: i16,
2936 signed: bool,
2937}
2938
2939fn resolve_preserved_zoned_format(
2940 metadata: &serde_json::Map<String, Value>,
2941 field_path: &str,
2942 field_name: &str,
2943) -> Option<ZonedEncodingFormat> {
2944 let candidates = [field_name, field_path];
2945 for key in candidates {
2946 if let Some(format) = metadata
2947 .get(key)
2948 .and_then(parse_zoned_encoding_metadata_value)
2949 {
2950 return Some(format);
2951 }
2952 }
2953 None
2954}
2955
2956fn parse_zoned_encoding_metadata_value(value: &Value) -> Option<ZonedEncodingFormat> {
2957 match value {
2958 Value::String(s) => parse_zoned_encoding_format_str(s),
2959 Value::Object(map) => map
2960 .get("zoned_encoding")
2961 .and_then(Value::as_str)
2962 .and_then(parse_zoned_encoding_format_str),
2963 _ => None,
2964 }
2965}
2966
2967fn parse_zoned_encoding_format_str(value: &str) -> Option<ZonedEncodingFormat> {
2968 match value.trim().to_ascii_lowercase().as_str() {
2969 "ascii" => Some(ZonedEncodingFormat::Ascii),
2970 "ebcdic" => Some(ZonedEncodingFormat::Ebcdic),
2971 "auto" => Some(ZonedEncodingFormat::Auto),
2972 _ => None,
2973 }
2974}
2975
2976fn coerce_to_str(value: &Value, coerce: bool) -> Option<String> {
2982 match value {
2983 Value::String(s) => Some(s.clone()),
2984 Value::Number(n) if coerce => Some(n.to_string()),
2985 _ => None,
2986 }
2987}
2988
2989fn json_type_name(value: &Value) -> &'static str {
2991 match value {
2992 Value::Null => "null",
2993 Value::Bool(_) => "boolean",
2994 Value::Number(_) => "number",
2995 Value::String(_) => "string",
2996 Value::Array(_) => "array",
2997 Value::Object(_) => "object",
2998 }
2999}
3000
3001fn encodable_numeric_text(
3008 json_obj: &serde_json::Map<String, Value>,
3009 field: ©book_core::Field,
3010 json_field_name: &str,
3011 expected: &str,
3012 coerce_numbers: bool,
3013) -> Result<Option<String>> {
3014 let Some(value) = json_obj.get(json_field_name) else {
3015 return Ok(None);
3016 };
3017 if value.is_null() {
3018 return Ok(None);
3019 }
3020 if let Some(text) = coerce_to_str(value, coerce_numbers) {
3021 return Ok(Some(text));
3022 }
3023
3024 let hint = if value.is_number() {
3025 " (pass --coerce-numbers to accept JSON numbers here)"
3026 } else {
3027 ""
3028 };
3029 Err(Error::new(
3030 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
3031 format!(
3032 "Field '{}' expects {expected}, found {}{hint}",
3033 field.name,
3034 json_type_name(value),
3035 ),
3036 ))
3037}
3038
3039#[inline]
3040#[allow(clippy::too_many_arguments)]
3041fn encode_zoned_decimal_field(
3042 field: ©book_core::Field,
3043 field_path: &str,
3044 json_field_name: &str,
3045 json_obj: &serde_json::Map<String, Value>,
3046 encoding_metadata: Option<&serde_json::Map<String, Value>>,
3047 buffer: &mut [u8],
3048 current_offset: usize,
3049 options: &EncodeOptions,
3050 spec: DecimalSpec,
3051) -> Result<usize> {
3052 let field_len = field.len as usize;
3053
3054 if let Some(text) = encodable_numeric_text(
3055 json_obj,
3056 field,
3057 json_field_name,
3058 "a zoned decimal string",
3059 options.coerce_numbers,
3060 )? {
3061 if field.blank_when_zero && options.bwz_encode {
3064 let encoded = crate::numeric::encode_zoned_decimal_with_bwz(
3065 &text,
3066 spec.digits,
3067 spec.scale,
3068 spec.signed,
3069 options.codepage,
3070 options.bwz_encode,
3071 )?;
3072 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
3073 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
3074 }
3075 return Ok(current_offset + field_len);
3076 }
3077
3078 let preserved_format = encoding_metadata
3079 .and_then(|meta| resolve_preserved_zoned_format(meta, field_path, json_field_name));
3080 let resolved_format = options
3081 .zoned_encoding_override
3082 .or(preserved_format)
3083 .unwrap_or(options.preferred_zoned_encoding);
3084 let (effective_format, zero_policy) = match resolved_format {
3086 ZonedEncodingFormat::Ascii => (ZonedEncodingFormat::Ascii, ZeroSignPolicy::Positive),
3087 ZonedEncodingFormat::Ebcdic => (ZonedEncodingFormat::Ebcdic, ZeroSignPolicy::Preferred),
3088 ZonedEncodingFormat::Auto => {
3089 if options.codepage.is_ascii() {
3090 (ZonedEncodingFormat::Ascii, ZeroSignPolicy::Positive)
3091 } else {
3092 (ZonedEncodingFormat::Ebcdic, ZeroSignPolicy::Preferred)
3093 }
3094 }
3095 };
3096
3097 let encoded = crate::numeric::encode_zoned_decimal_with_format_and_policy(
3098 &text,
3099 spec.digits,
3100 spec.scale,
3101 spec.signed,
3102 options.codepage,
3103 Some(effective_format),
3104 zero_policy,
3105 )?;
3106
3107 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
3108 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
3109 }
3110 }
3111
3112 Ok(current_offset + field_len)
3113}
3114
3115#[inline]
3116fn encode_packed_decimal_field(
3117 field: ©book_core::Field,
3118 json_field_name: &str,
3119 json_obj: &serde_json::Map<String, Value>,
3120 buffer: &mut [u8],
3121 current_offset: usize,
3122 options: &EncodeOptions,
3123 spec: DecimalSpec,
3124) -> Result<usize> {
3125 let field_len = field.len as usize;
3126
3127 if let Some(text) = encodable_numeric_text(
3128 json_obj,
3129 field,
3130 json_field_name,
3131 "a packed decimal string",
3132 options.coerce_numbers,
3133 )? {
3134 let encoded =
3135 crate::numeric::encode_packed_decimal(&text, spec.digits, spec.scale, spec.signed)?;
3136 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
3137 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
3138 }
3139 }
3140
3141 Ok(current_offset + field_len)
3142}
3143
3144#[derive(Copy, Clone)]
3145struct BinarySpec {
3146 bits: u16,
3147 signed: bool,
3148}
3149
3150#[inline]
3151fn encode_binary_int_field(
3152 field: ©book_core::Field,
3153 json_field_name: &str,
3154 json_obj: &serde_json::Map<String, Value>,
3155 buffer: &mut [u8],
3156 current_offset: usize,
3157 options: &EncodeOptions,
3158 spec: BinarySpec,
3159) -> Result<usize> {
3160 let field_len = field.len as usize;
3161
3162 let value = match json_obj.get(json_field_name) {
3163 None => None,
3164 Some(v) if v.is_null() => None,
3165 Some(v) => Some(v),
3166 };
3167
3168 if let Some(value) = value {
3169 let num = if let Some(n) = value.as_i64() {
3171 n
3172 } else {
3173 let text = coerce_to_str(value, options.coerce_numbers).ok_or_else(|| {
3174 Error::new(
3175 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
3176 format!(
3177 "Field '{}' expects an integer, found {}",
3178 field.name,
3179 json_type_name(value),
3180 ),
3181 )
3182 })?;
3183 text.parse::<i64>().map_err(|e| {
3184 Error::new(
3185 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
3186 format!(
3187 "Field '{}' expects an integer, found '{text}': {e}",
3188 field.name,
3189 ),
3190 )
3191 })?
3192 };
3193 let encoded = crate::numeric::encode_binary_int(num, spec.bits, spec.signed)?;
3194 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
3195 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
3196 }
3197 }
3198
3199 Ok(current_offset + field_len)
3200}
3201
3202#[inline]
3231#[must_use = "Handle the Result or propagate the error"]
3232pub fn decode_file_to_jsonl(
3233 schema: &Schema,
3234 input: impl Read,
3235 mut output: impl Write,
3236 options: &DecodeOptions,
3237) -> Result<RunSummary> {
3238 let start_time = std::time::Instant::now();
3239 let mut summary = RunSummary::with_threads(effective_worker_count(options.threads));
3240 summary.set_schema_fingerprint(schema.fingerprint.clone());
3241
3242 reset_warning_counter();
3243
3244 match options.format {
3245 RecordFormat::Fixed => {
3246 process_fixed_records(schema, input, &mut output, options, &mut summary)?;
3247 }
3248 RecordFormat::RDW => {
3249 process_rdw_records(schema, input, &mut output, options, &mut summary)?;
3250 }
3251 RecordFormat::Vb => {
3252 process_vb_records(schema, input, &mut output, options, &mut summary)?;
3253 }
3254 }
3255
3256 let elapsed_ms = start_time.elapsed().as_millis();
3257 summary.processing_time_ms = u64::try_from(elapsed_ms).unwrap_or(u64::MAX);
3258 summary.calculate_throughput();
3259 summary.warnings = warning_count();
3260 telemetry::record_completion(
3261 summary.processing_time_seconds(),
3262 summary.throughput_mbps,
3263 options,
3264 );
3265 info!(
3266 target: "copybook::decode",
3267 records_processed = summary.records_processed,
3268 records_with_errors = summary.records_with_errors,
3269 warnings = summary.warnings,
3270 bytes_processed = summary.bytes_processed,
3271 elapsed_ms = summary.processing_time_ms,
3272 throughput_mibps = summary.throughput_mbps,
3273 schema_fingerprint = %summary.schema_fingerprint,
3274 codepage = %options.codepage,
3275 format = ?options.format,
3276 strict_mode = options.strict_mode,
3277 raw_mode = ?options.emit_raw,
3278 );
3279
3280 Ok(summary)
3281}
3282
3283fn process_fixed_records<R: Read, W: Write>(
3284 schema: &Schema,
3285 reader: R,
3286 output: &mut W,
3287 options: &DecodeOptions,
3288 summary: &mut RunSummary,
3289) -> Result<()> {
3290 if options.threads > 1 {
3291 return process_fixed_records_parallel(schema, reader, output, options, summary);
3292 }
3293
3294 let mut reader = crate::file::fixed::reader(reader, schema)?;
3295 let mut scratch = crate::memory::ScratchBuffers::new();
3296 let mut record_index = 0u64;
3297 let mut record_offset = 0u64;
3298
3299 while let Some(record_data) = reader.read_record()? {
3300 record_index += 1;
3301 let current_offset = record_offset;
3302 record_offset = record_offset.saturating_add(record_data.len() as u64);
3303 crate::file::fixed::validate_record_length(
3304 schema,
3305 reader.lrecl(),
3306 reader.record_count(),
3307 &record_data,
3308 )?;
3309 summary.bytes_processed += record_data.len() as u64;
3310 telemetry::record_read(record_data.len(), options);
3311
3312 let raw_data_for_decode = match options.emit_raw {
3313 crate::options::RawMode::Record => Some(record_data.clone()),
3314 _ => None,
3315 };
3316
3317 match decode_record_with_scratch_and_raw(
3318 schema,
3319 &record_data,
3320 options,
3321 raw_data_for_decode.as_deref(),
3322 record_index,
3323 Some(current_offset),
3324 &mut scratch,
3325 ) {
3326 Ok(json_value) => {
3327 write_json_record(output, &json_value)?;
3328 summary.records_processed += 1;
3329 }
3330 Err(error) => {
3331 summary.note_failure(record_index, &error);
3332 telemetry::record_error(error.family_prefix());
3333 if options.strict_mode {
3334 return Err(error);
3335 }
3336 }
3337 }
3338 }
3339
3340 Ok(())
3341}
3342
3343struct DecodeWork {
3344 payload: Vec<u8>,
3345 raw_data: Option<Vec<u8>>,
3346 record_index: u64,
3347 record_offset: u64,
3348}
3349
3350struct DecodeOutcome {
3351 result: Result<Value>,
3352 warnings: u64,
3353 record_index: u64,
3355}
3356
3357fn effective_worker_count(requested: usize) -> usize {
3358 requested.clamp(1, MAX_WORKERS)
3359}
3360
3361fn decode_worker_pool(
3362 schema: &Schema,
3363 options: &DecodeOptions,
3364) -> crate::memory::WorkerPool<DecodeWork, DecodeOutcome> {
3365 let workers = effective_worker_count(options.threads);
3366 let channel_capacity = workers.saturating_mul(4).max(1);
3367 let max_window_size = workers.saturating_mul(2).max(1);
3368 let schema = Arc::new(schema.clone());
3369 let options = Arc::new(options.clone());
3370
3371 crate::memory::WorkerPool::new(
3372 workers,
3373 channel_capacity,
3374 max_window_size,
3375 move |work: DecodeWork, scratch: &mut crate::memory::ScratchBuffers| {
3376 let warning_count_before = warning_count();
3377 let result = decode_record_with_scratch_and_raw(
3378 &schema,
3379 &work.payload,
3380 &options,
3381 work.raw_data.as_deref(),
3382 work.record_index,
3383 Some(work.record_offset),
3384 scratch,
3385 );
3386 let warnings = warning_count().saturating_sub(warning_count_before);
3387 DecodeOutcome {
3388 result,
3389 warnings,
3390 record_index: work.record_index,
3391 }
3392 },
3393 )
3394}
3395
3396fn process_decode_batch<W: Write>(
3397 pool: &mut crate::memory::WorkerPool<DecodeWork, DecodeOutcome>,
3398 batch_len: usize,
3399 output: &mut W,
3400 options: &DecodeOptions,
3401 summary: &mut RunSummary,
3402) -> Result<()> {
3403 let mut first_error = None;
3404
3405 for _ in 0..batch_len {
3406 let outcome = pool
3407 .recv_ordered()
3408 .map_err(|error| Error::new(ErrorCode::CBKI001_INVALID_STATE, error.to_string()))?
3409 .ok_or_else(|| {
3410 Error::new(
3411 ErrorCode::CBKI001_INVALID_STATE,
3412 "decode worker pool ended before the submitted batch completed",
3413 )
3414 })?;
3415
3416 for _ in 0..outcome.warnings {
3417 increment_warning_counter();
3418 }
3419
3420 if first_error.is_some() {
3421 continue;
3422 }
3423
3424 let record_index = outcome.record_index;
3425 match outcome.result {
3426 Ok(json_value) => {
3427 write_json_record(output, &json_value)?;
3428 summary.records_processed += 1;
3429 }
3430 Err(error) => {
3431 summary.note_failure(record_index, &error);
3432 telemetry::record_error(error.family_prefix());
3433 if options.strict_mode {
3434 first_error = Some(error);
3435 }
3436 }
3437 }
3438 }
3439
3440 first_error.map_or(Ok(()), Err)
3441}
3442
3443fn process_fixed_records_parallel<R: Read, W: Write>(
3444 schema: &Schema,
3445 reader: R,
3446 output: &mut W,
3447 options: &DecodeOptions,
3448 summary: &mut RunSummary,
3449) -> Result<()> {
3450 let mut reader = crate::file::fixed::reader(reader, schema)?;
3451 let workers = effective_worker_count(options.threads);
3452 let batch_capacity = workers.saturating_mul(4).max(1);
3453 let mut pool = decode_worker_pool(schema, options);
3454 let mut record_index = 0_u64;
3455 let mut record_offset = 0_u64;
3456 let mut batch_len = 0_usize;
3457
3458 loop {
3459 let record = match reader.read_record() {
3460 Ok(record) => record,
3461 Err(error) => {
3462 let pending_result = if batch_len > 0 {
3463 process_decode_batch(&mut pool, batch_len, output, options, summary)
3464 } else {
3465 Ok(())
3466 };
3467 let _ = pool.shutdown();
3468 pending_result?;
3469 return Err(error);
3470 }
3471 };
3472 let Some(record_data) = record else { break };
3473
3474 record_index += 1;
3475 let current_offset = record_offset;
3476 record_offset = record_offset.saturating_add(record_data.len() as u64);
3477 crate::file::fixed::validate_record_length(
3478 schema,
3479 reader.lrecl(),
3480 reader.record_count(),
3481 &record_data,
3482 )?;
3483 summary.bytes_processed += record_data.len() as u64;
3484 telemetry::record_read(record_data.len(), options);
3485 let raw_data = match options.emit_raw {
3486 crate::options::RawMode::Record => Some(record_data.clone()),
3487 _ => None,
3488 };
3489
3490 if let Err(error) = pool.submit(DecodeWork {
3491 payload: record_data,
3492 raw_data,
3493 record_index,
3494 record_offset: current_offset,
3495 }) {
3496 let pending_result = if batch_len > 0 {
3497 process_decode_batch(&mut pool, batch_len, output, options, summary)
3498 } else {
3499 Ok(())
3500 };
3501 let _ = pool.shutdown();
3502 pending_result?;
3503 return Err(Error::new(
3504 ErrorCode::CBKI001_INVALID_STATE,
3505 error.to_string(),
3506 ));
3507 }
3508 batch_len += 1;
3509
3510 if batch_len == batch_capacity {
3511 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
3512 batch_len = 0;
3513 if let Err(error) = result {
3514 let _ = pool.shutdown();
3515 return Err(error);
3516 }
3517 }
3518 }
3519
3520 if batch_len > 0 {
3521 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
3522 if let Err(error) = result {
3523 let _ = pool.shutdown();
3524 return Err(error);
3525 }
3526 }
3527
3528 pool.shutdown().map_err(|error| {
3529 Error::new(
3530 ErrorCode::CBKI001_INVALID_STATE,
3531 format!("decode worker pool shutdown failed: {error}"),
3532 )
3533 })
3534}
3535
3536fn process_rdw_records<R: Read, W: Write>(
3537 schema: &Schema,
3538 reader: R,
3539 output: &mut W,
3540 options: &DecodeOptions,
3541 summary: &mut RunSummary,
3542) -> Result<()> {
3543 if options.threads > 1 {
3544 return process_rdw_records_parallel(schema, reader, output, options, summary);
3545 }
3546
3547 let mut reader = crate::record::RDWRecordReader::new(reader, options.strict_mode);
3548 let mut scratch = crate::memory::ScratchBuffers::new();
3549 let mut record_index = 0u64;
3550 let mut record_offset = 0u64;
3551
3552 while let Some(rdw_record) = reader.read_record()? {
3553 record_index += 1;
3554 let record_bytes = rdw_record.header.len() + rdw_record.payload.len();
3555 let current_offset = record_offset;
3556 record_offset = record_offset.saturating_add(record_bytes as u64);
3557 summary.bytes_processed += record_bytes as u64;
3558 telemetry::record_read(record_bytes, options);
3559 if rdw_record.reserved() != 0 {
3560 increment_warning_counter();
3561 }
3562
3563 if let Some(schema_lrecl) = schema.lrecl_fixed
3571 && schema.tail_odo.is_none()
3572 && rdw_record.payload.len() < schema_lrecl as usize
3573 {
3574 let error = rdw_underflow_error(schema_lrecl, rdw_record.payload.len());
3575
3576 summary.note_failure(record_index, &error);
3577 telemetry::record_error(error.family_prefix());
3578 if options.strict_mode {
3579 return Err(error);
3580 }
3581 continue;
3582 }
3583
3584 let full_raw_data = rdw_raw_data(&rdw_record, options.emit_raw);
3585
3586 match decode_record_with_scratch_and_raw(
3587 schema,
3588 &rdw_record.payload,
3589 options,
3590 full_raw_data.as_deref(),
3591 record_index,
3592 Some(current_offset),
3593 &mut scratch,
3594 ) {
3595 Ok(json_value) => {
3596 write_json_record(output, &json_value)?;
3597 summary.records_processed += 1;
3598 }
3599 Err(error) => {
3600 summary.note_failure(record_index, &error);
3601 telemetry::record_error(error.family_prefix());
3602 if options.strict_mode {
3603 return Err(error);
3604 }
3605 }
3606 }
3607 }
3608
3609 Ok(())
3610}
3611
3612fn process_vb_records<R: Read, W: Write>(
3613 schema: &Schema,
3614 reader: R,
3615 output: &mut W,
3616 options: &DecodeOptions,
3617 summary: &mut RunSummary,
3618) -> Result<()> {
3619 if options.threads > 1 {
3620 return process_vb_records_parallel(schema, reader, output, options, summary);
3621 }
3622
3623 let mut reader = crate::record::VbBlockReader::new(reader, options.strict_mode);
3624 let mut scratch = crate::memory::ScratchBuffers::new();
3625 let mut record_index = 0u64;
3626
3627 while let Some(vb_record) = reader.read_record()? {
3628 record_index += 1;
3629 let record_bytes = vb_record.rdw.len() + vb_record.payload.len();
3630 let current_offset = vb_record.physical_offset;
3633 summary.bytes_processed += record_bytes as u64;
3634 telemetry::record_read(record_bytes, options);
3635 if vb_record.rdw_reserved != 0 {
3636 increment_warning_counter();
3637 }
3638
3639 if let Some(schema_lrecl) = schema.lrecl_fixed
3645 && schema.tail_odo.is_none()
3646 && vb_record.payload.len() < schema_lrecl as usize
3647 {
3648 let error = rdw_underflow_error(schema_lrecl, vb_record.payload.len());
3649
3650 summary.note_failure(record_index, &error);
3651 telemetry::record_error(error.family_prefix());
3652 if options.strict_mode {
3653 return Err(error);
3654 }
3655 continue;
3656 }
3657
3658 let full_raw_data = vb_raw_data(&vb_record, options.emit_raw);
3659
3660 match decode_record_with_scratch_and_raw(
3661 schema,
3662 &vb_record.payload,
3663 options,
3664 full_raw_data.as_deref(),
3665 record_index,
3666 Some(current_offset),
3667 &mut scratch,
3668 ) {
3669 Ok(json_value) => {
3670 write_json_record(output, &json_value)?;
3671 summary.records_processed += 1;
3672 }
3673 Err(error) => {
3674 summary.note_failure(record_index, &error);
3675 telemetry::record_error(error.family_prefix());
3676 if options.strict_mode {
3677 return Err(error);
3678 }
3679 }
3680 }
3681 }
3682
3683 Ok(())
3684}
3685
3686fn vb_raw_data(
3687 record: &crate::record::VbRecord,
3688 raw_mode: crate::options::RawMode,
3689) -> Option<Vec<u8>> {
3690 match raw_mode {
3691 crate::options::RawMode::RecordRDW => {
3692 let mut full_data = Vec::with_capacity(record.rdw.len() + record.payload.len());
3693 full_data.extend_from_slice(&record.rdw);
3694 full_data.extend_from_slice(&record.payload);
3695 Some(full_data)
3696 }
3697 crate::options::RawMode::Record => Some(record.payload.clone()),
3698 _ => None,
3699 }
3700}
3701
3702fn process_vb_records_parallel<R: Read, W: Write>(
3703 schema: &Schema,
3704 reader: R,
3705 output: &mut W,
3706 options: &DecodeOptions,
3707 summary: &mut RunSummary,
3708) -> Result<()> {
3709 let mut reader = crate::record::VbBlockReader::new(reader, options.strict_mode);
3710 let workers = effective_worker_count(options.threads);
3711 let batch_capacity = workers.saturating_mul(4).max(1);
3712 let mut pool = decode_worker_pool(schema, options);
3713 let mut record_index = 0_u64;
3714 let mut batch_len = 0_usize;
3715
3716 loop {
3717 let vb_record = match reader.read_record() {
3718 Ok(record) => record,
3719 Err(error) => {
3720 let pending_result = if batch_len > 0 {
3721 process_decode_batch(&mut pool, batch_len, output, options, summary)
3722 } else {
3723 Ok(())
3724 };
3725 let _ = pool.shutdown();
3726 pending_result?;
3727 return Err(error);
3728 }
3729 };
3730 let Some(vb_record) = vb_record else { break };
3731
3732 record_index += 1;
3733 let record_bytes = vb_record.rdw.len() + vb_record.payload.len();
3734 let current_offset = vb_record.physical_offset;
3737 summary.bytes_processed += record_bytes as u64;
3738 telemetry::record_read(record_bytes, options);
3739 if vb_record.rdw_reserved != 0 {
3740 increment_warning_counter();
3741 }
3742
3743 if let Some(schema_lrecl) = schema.lrecl_fixed
3744 && schema.tail_odo.is_none()
3745 && vb_record.payload.len() < schema_lrecl as usize
3746 {
3747 let error = rdw_underflow_error(schema_lrecl, vb_record.payload.len());
3748
3749 summary.note_failure(record_index, &error);
3750 telemetry::record_error(error.family_prefix());
3751 if options.strict_mode {
3752 let pending_result = if batch_len > 0 {
3753 process_decode_batch(&mut pool, batch_len, output, options, summary)
3754 } else {
3755 Ok(())
3756 };
3757 let _ = pool.shutdown();
3758 pending_result?;
3759 return Err(error);
3760 }
3761 continue;
3762 }
3763
3764 let full_raw_data = vb_raw_data(&vb_record, options.emit_raw);
3765
3766 if let Err(error) = pool.submit(DecodeWork {
3767 payload: vb_record.payload,
3768 raw_data: full_raw_data,
3769 record_index,
3770 record_offset: current_offset,
3771 }) {
3772 let pending_result = if batch_len > 0 {
3773 process_decode_batch(&mut pool, batch_len, output, options, summary)
3774 } else {
3775 Ok(())
3776 };
3777 let _ = pool.shutdown();
3778 pending_result?;
3779 return Err(Error::new(
3780 ErrorCode::CBKI001_INVALID_STATE,
3781 error.to_string(),
3782 ));
3783 }
3784 batch_len += 1;
3785
3786 if batch_len == batch_capacity {
3787 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
3788 batch_len = 0;
3789 if let Err(error) = result {
3790 let _ = pool.shutdown();
3791 return Err(error);
3792 }
3793 }
3794 }
3795
3796 if batch_len > 0 {
3797 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
3798 if let Err(error) = result {
3799 let _ = pool.shutdown();
3800 return Err(error);
3801 }
3802 }
3803
3804 pool.shutdown().map_err(|error| {
3805 Error::new(
3806 ErrorCode::CBKI001_INVALID_STATE,
3807 format!("decode worker pool shutdown failed: {error}"),
3808 )
3809 })
3810}
3811
3812fn rdw_underflow_error(schema_lrecl: u32, payload_len: usize) -> Error {
3813 Error::new(
3814 ErrorCode::CBKF221_RDW_UNDERFLOW,
3815 format!("RDW payload too short: {payload_len} bytes, schema requires {schema_lrecl} bytes"),
3816 )
3817}
3818
3819fn rdw_raw_data(
3820 record: &crate::record::RDWRecord,
3821 raw_mode: crate::options::RawMode,
3822) -> Option<Vec<u8>> {
3823 match raw_mode {
3824 crate::options::RawMode::RecordRDW => {
3825 let mut full_data = Vec::with_capacity(record.header.len() + record.payload.len());
3826 full_data.extend_from_slice(&record.header);
3827 full_data.extend_from_slice(&record.payload);
3828 Some(full_data)
3829 }
3830 crate::options::RawMode::Record => Some(record.payload.clone()),
3831 _ => None,
3832 }
3833}
3834
3835fn process_rdw_records_parallel<R: Read, W: Write>(
3836 schema: &Schema,
3837 reader: R,
3838 output: &mut W,
3839 options: &DecodeOptions,
3840 summary: &mut RunSummary,
3841) -> Result<()> {
3842 let mut reader = crate::record::RDWRecordReader::new(reader, options.strict_mode);
3843 let workers = effective_worker_count(options.threads);
3844 let batch_capacity = workers.saturating_mul(4).max(1);
3845 let mut pool = decode_worker_pool(schema, options);
3846 let mut record_index = 0_u64;
3847 let mut record_offset = 0_u64;
3848 let mut batch_len = 0_usize;
3849
3850 loop {
3851 let rdw_record = match reader.read_record() {
3852 Ok(record) => record,
3853 Err(error) => {
3854 let pending_result = if batch_len > 0 {
3855 process_decode_batch(&mut pool, batch_len, output, options, summary)
3856 } else {
3857 Ok(())
3858 };
3859 let _ = pool.shutdown();
3860 pending_result?;
3861 return Err(error);
3862 }
3863 };
3864 let Some(rdw_record) = rdw_record else { break };
3865
3866 record_index += 1;
3867 let record_bytes = rdw_record.header.len() + rdw_record.payload.len();
3868 let current_offset = record_offset;
3869 record_offset = record_offset.saturating_add(record_bytes as u64);
3870 summary.bytes_processed += record_bytes as u64;
3871 telemetry::record_read(record_bytes, options);
3872 if rdw_record.reserved() != 0 {
3873 increment_warning_counter();
3874 }
3875
3876 if let Some(schema_lrecl) = schema.lrecl_fixed
3877 && schema.tail_odo.is_none()
3878 && rdw_record.payload.len() < schema_lrecl as usize
3879 {
3880 let error = rdw_underflow_error(schema_lrecl, rdw_record.payload.len());
3881
3882 summary.note_failure(record_index, &error);
3883 telemetry::record_error(error.family_prefix());
3884 if options.strict_mode {
3885 let pending_result = if batch_len > 0 {
3886 process_decode_batch(&mut pool, batch_len, output, options, summary)
3887 } else {
3888 Ok(())
3889 };
3890 let _ = pool.shutdown();
3891 pending_result?;
3892 return Err(error);
3893 }
3894 continue;
3895 }
3896
3897 let full_raw_data = rdw_raw_data(&rdw_record, options.emit_raw);
3898
3899 if let Err(error) = pool.submit(DecodeWork {
3900 payload: rdw_record.payload,
3901 raw_data: full_raw_data,
3902 record_index,
3903 record_offset: current_offset,
3904 }) {
3905 let pending_result = if batch_len > 0 {
3906 process_decode_batch(&mut pool, batch_len, output, options, summary)
3907 } else {
3908 Ok(())
3909 };
3910 let _ = pool.shutdown();
3911 pending_result?;
3912 return Err(Error::new(
3913 ErrorCode::CBKI001_INVALID_STATE,
3914 error.to_string(),
3915 ));
3916 }
3917 batch_len += 1;
3918
3919 if batch_len == batch_capacity {
3920 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
3921 batch_len = 0;
3922 if let Err(error) = result {
3923 let _ = pool.shutdown();
3924 return Err(error);
3925 }
3926 }
3927 }
3928
3929 if batch_len > 0 {
3930 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
3931 if let Err(error) = result {
3932 let _ = pool.shutdown();
3933 return Err(error);
3934 }
3935 }
3936
3937 pool.shutdown().map_err(|error| {
3938 Error::new(
3939 ErrorCode::CBKI001_INVALID_STATE,
3940 format!("decode worker pool shutdown failed: {error}"),
3941 )
3942 })
3943}
3944
3945#[inline]
3946fn write_json_record<W: Write>(output: &mut W, value: &Value) -> Result<()> {
3947 if let Err(e) = serde_json::to_writer(&mut *output, value) {
3948 let error = Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string());
3949 telemetry::record_error(error.family_prefix());
3950 return Err(error);
3951 }
3952
3953 if let Err(e) = writeln!(output) {
3954 let error = Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string());
3955 telemetry::record_error(error.family_prefix());
3956 return Err(error);
3957 }
3958
3959 Ok(())
3960}
3961
3962fn encode_worker_pool(
3963 schema: &Schema,
3964 options: &EncodeOptions,
3965) -> crate::memory::WorkerPool<Value, Result<Vec<u8>>> {
3966 let workers = effective_worker_count(options.threads);
3967 let channel_capacity = workers.saturating_mul(4).max(1);
3968 let max_window_size = workers.saturating_mul(2).max(1);
3969 let schema = Arc::new(schema.clone());
3970 let options = Arc::new(options.clone());
3971
3972 crate::memory::WorkerPool::new(
3973 workers,
3974 channel_capacity,
3975 max_window_size,
3976 move |json_value: Value, _scratch: &mut crate::memory::ScratchBuffers| {
3977 encode_record(&schema, &json_value, &options)
3978 },
3979 )
3980}
3981
3982fn process_encode_batch<W: Write>(
3983 pool: &mut crate::memory::WorkerPool<Value, Result<Vec<u8>>>,
3984 batch_len: usize,
3985 records_before_batch: u64,
3986 output: &mut W,
3987 options: &EncodeOptions,
3988 summary: &mut RunSummary,
3989) -> Result<bool> {
3990 let mut stop_after_error = false;
3991
3992 for position in 0..batch_len {
3993 let result = pool
3994 .recv_ordered()
3995 .map_err(|error| Error::new(ErrorCode::CBKI001_INVALID_STATE, error.to_string()))?
3996 .ok_or_else(|| {
3997 Error::new(
3998 ErrorCode::CBKI001_INVALID_STATE,
3999 "encode worker pool ended before the submitted batch completed",
4000 )
4001 })?;
4002
4003 if stop_after_error {
4004 continue;
4005 }
4006
4007 match result {
4008 Ok(binary_data) => {
4009 output.write_all(&binary_data).map_err(|error| {
4010 Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, error.to_string())
4011 })?;
4012 summary.bytes_processed += binary_data.len() as u64;
4013 summary.records_processed += 1;
4014 }
4015 Err(error) => {
4016 summary.note_failure(records_before_batch + position as u64 + 1, &error);
4018 telemetry::record_error(error.family_prefix());
4019 if options.strict_mode {
4020 stop_after_error = true;
4021 }
4022 }
4023 }
4024 }
4025
4026 Ok(stop_after_error)
4027}
4028
4029fn shutdown_encode_pool(pool: crate::memory::WorkerPool<Value, Result<Vec<u8>>>) -> Result<()> {
4030 pool.shutdown().map_err(|error| {
4031 Error::new(
4032 ErrorCode::CBKI001_INVALID_STATE,
4033 format!("encode worker pool shutdown failed: {error}"),
4034 )
4035 })
4036}
4037
4038fn finish_encode_input_error<W: Write>(
4039 pool: crate::memory::WorkerPool<Value, Result<Vec<u8>>>,
4040 batch_len: usize,
4041 records_before_batch: u64,
4042 output: &mut W,
4043 options: &EncodeOptions,
4044 summary: &mut RunSummary,
4045 error: Error,
4046) -> Result<u64> {
4047 let mut pool = pool;
4048 let pending_result = if batch_len > 0 {
4049 process_encode_batch(
4050 &mut pool,
4051 batch_len,
4052 records_before_batch,
4053 output,
4054 options,
4055 summary,
4056 )
4057 } else {
4058 Ok(false)
4059 };
4060 let shutdown_result = shutdown_encode_pool(pool);
4061 let pending_stop = pending_result?;
4062 shutdown_result?;
4063 if pending_stop {
4064 Ok(summary.records_processed)
4065 } else {
4066 Err(error)
4067 }
4068}
4069
4070fn process_encode_jsonl_parallel<R: BufRead, W: Write>(
4071 schema: &Schema,
4072 reader: R,
4073 output: &mut W,
4074 options: &EncodeOptions,
4075 summary: &mut RunSummary,
4076) -> Result<u64> {
4077 let workers = effective_worker_count(options.threads);
4078 let batch_capacity = workers.saturating_mul(4).max(1);
4079 let mut pool = encode_worker_pool(schema, options);
4080 let mut records_seen = 0_u64;
4081 let mut records_before_batch = 0_u64;
4082 let mut batch_len = 0_usize;
4083
4084 for line in reader.lines() {
4085 let line = match line {
4086 Ok(line) => line,
4087 Err(error) => {
4088 return finish_encode_input_error(
4089 pool,
4090 batch_len,
4091 records_before_batch,
4092 output,
4093 options,
4094 summary,
4095 Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, error.to_string()),
4096 );
4097 }
4098 };
4099
4100 if line.trim().is_empty() {
4101 continue;
4102 }
4103
4104 let json_value: Value = match serde_json::from_str(&line) {
4105 Ok(json_value) => json_value,
4106 Err(error) => {
4107 return finish_encode_input_error(
4108 pool,
4109 batch_len,
4110 records_before_batch,
4111 output,
4112 options,
4113 summary,
4114 Error::new(ErrorCode::CBKE501_JSON_TYPE_MISMATCH, error.to_string()),
4115 );
4116 }
4117 };
4118
4119 records_seen += 1;
4120 if let Err(error) = pool.submit(json_value) {
4121 return finish_encode_input_error(
4122 pool,
4123 batch_len,
4124 records_before_batch,
4125 output,
4126 options,
4127 summary,
4128 Error::new(ErrorCode::CBKI001_INVALID_STATE, error.to_string()),
4129 );
4130 }
4131 batch_len += 1;
4132
4133 if batch_len == batch_capacity {
4134 let batch_result = process_encode_batch(
4135 &mut pool,
4136 batch_len,
4137 records_before_batch,
4138 output,
4139 options,
4140 summary,
4141 );
4142 let stop = match batch_result {
4143 Ok(stop) => stop,
4144 Err(error) => {
4145 let _ = shutdown_encode_pool(pool);
4146 return Err(error);
4147 }
4148 };
4149 batch_len = 0;
4150 records_before_batch = records_seen;
4151 if stop {
4152 shutdown_encode_pool(pool)?;
4153 return Ok(summary.records_processed);
4154 }
4155 }
4156 }
4157
4158 if batch_len > 0 {
4159 let batch_result = process_encode_batch(
4160 &mut pool,
4161 batch_len,
4162 records_before_batch,
4163 output,
4164 options,
4165 summary,
4166 );
4167 let stop = match batch_result {
4168 Ok(stop) => stop,
4169 Err(error) => {
4170 let _ = shutdown_encode_pool(pool);
4171 return Err(error);
4172 }
4173 };
4174 if stop {
4175 shutdown_encode_pool(pool)?;
4176 return Ok(summary.records_processed);
4177 }
4178 }
4179
4180 shutdown_encode_pool(pool)?;
4181 Ok(summary.records_processed)
4182}
4183
4184#[inline]
4221#[must_use = "Handle the Result or propagate the error"]
4222pub fn encode_jsonl_to_file(
4223 schema: &Schema,
4224 input: impl Read,
4225 mut output: impl Write,
4226 options: &EncodeOptions,
4227) -> Result<RunSummary> {
4228 let start_time = std::time::Instant::now();
4229 let mut summary = RunSummary::with_threads(effective_worker_count(options.threads));
4230 summary.set_schema_fingerprint(schema.fingerprint.clone());
4231
4232 let reader = BufReader::new(input);
4233 let record_count = if options.threads > 1 {
4234 process_encode_jsonl_parallel(schema, reader, &mut output, options, &mut summary)?
4235 } else {
4236 let mut records_seen = 0u64;
4237 let mut records_processed = 0u64;
4238
4239 for line in reader.lines() {
4240 let line =
4241 line.map_err(|e| Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string()))?;
4242
4243 if line.trim().is_empty() {
4244 continue;
4245 }
4246
4247 records_seen += 1;
4248
4249 let json_value: Value = serde_json::from_str(&line)
4251 .map_err(|e| Error::new(ErrorCode::CBKE501_JSON_TYPE_MISMATCH, e.to_string()))?;
4252
4253 match encode_record(schema, &json_value, options) {
4255 Ok(binary_data) => {
4256 output.write_all(&binary_data).map_err(|e| {
4257 Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string())
4258 })?;
4259 summary.bytes_processed += binary_data.len() as u64;
4260 records_processed += 1;
4261 }
4262 Err(error) => {
4263 summary.note_failure(records_seen, &error);
4264 telemetry::record_error(error.family_prefix());
4265 if options.strict_mode {
4266 break;
4267 }
4268 }
4269 }
4270 }
4271
4272 records_processed
4273 };
4274
4275 summary.records_processed = record_count;
4276 let elapsed_ms = start_time.elapsed().as_millis();
4277 summary.processing_time_ms = u64::try_from(elapsed_ms).unwrap_or(u64::MAX);
4278 summary.calculate_throughput();
4279
4280 Ok(summary)
4281}
4282
4283fn format_zoned_decimal_with_digits(
4285 decimal: &crate::numeric::SmallDecimal,
4286 digits: u16,
4287 blank_when_zero: bool,
4288) -> String {
4289 use std::fmt::Write;
4290
4291 if blank_when_zero {
4293 return decimal.to_string();
4294 }
4295
4296 if decimal.value == 0 {
4299 let natural_format = decimal.to_string();
4300 if natural_format == "0" {
4301 return "0".to_string();
4302 }
4303 }
4304
4305 let mut result = String::new();
4307 let value = decimal.value;
4308 let negative = decimal.negative && value != 0;
4309
4310 if negative {
4311 result.push('-');
4312 }
4313
4314 if decimal.scale <= 0 {
4316 let scaled_value = if decimal.scale < 0 {
4317 let exponent = u32::from(decimal.scale.unsigned_abs());
4318 value * 10_i64.pow(exponent)
4319 } else {
4320 value
4321 };
4322 if write!(result, "{:0width$}", scaled_value, width = digits as usize).is_err() {
4323 result.push('0');
4325 }
4326 } else {
4327 result.push_str(&decimal.to_string());
4329 }
4330
4331 result
4332}
4333
4334#[inline]
4335fn small_decimal_to_string(decimal: &crate::numeric::SmallDecimal) -> String {
4336 decimal.to_string()
4337}
4338
4339fn zoned_decimal_to_json_value(
4340 decimal: &crate::numeric::SmallDecimal,
4341 digits: u16,
4342 scale: i16,
4343 blank_when_zero: bool,
4344 options: &DecodeOptions,
4345) -> Value {
4346 let formatted = if scale == 0 {
4347 format_zoned_decimal_with_digits(decimal, digits, blank_when_zero)
4348 } else {
4349 small_decimal_to_string(decimal)
4350 };
4351 numeric_string_to_value(formatted, options)
4352}
4353
4354#[inline]
4355fn decimal_counter_to_u32(
4356 decimal: &crate::numeric::SmallDecimal,
4357 counter_path: &str,
4358) -> Result<u32> {
4359 let text = small_decimal_to_string(decimal);
4360 text.parse::<u32>().map_err(|_| {
4361 Error::new(
4362 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
4363 format!("ODO counter '{counter_path}' has invalid value: {text}"),
4364 )
4365 })
4366}
4367
4368#[cfg(test)]
4369#[allow(clippy::expect_used)]
4370#[allow(clippy::unwrap_used)]
4371mod tests {
4372 use super::*;
4373 use crate::Codepage;
4374 use crate::iterator::RecordIterator;
4375 use copybook_core::{Error, ErrorCode, Result, parse_copybook};
4376 use std::io::Cursor;
4377
4378 #[test]
4379 fn test_decode_record() -> Result<()> {
4380 let copybook_text = r"
4381 01 RECORD.
4382 05 ID PIC 9(3).
4383 05 NAME PIC X(5).
4384 ";
4385
4386 let schema = parse_copybook(copybook_text)?;
4387 let options = DecodeOptions {
4388 codepage: Codepage::ASCII, ..DecodeOptions::default()
4390 };
4391 let data = b"001ALICE";
4392
4393 let result = decode_record(&schema, data, &options)?;
4394 assert!(result.is_object());
4395 let object = result.as_object().ok_or_else(|| {
4396 Error::new(
4397 ErrorCode::CBKP001_SYNTAX,
4398 "decoded record should be an object".to_string(),
4399 )
4400 })?;
4401 assert!(object.len() > 1);
4402 Ok(())
4403 }
4404
4405 #[test]
4406 fn test_encode_record() -> Result<()> {
4407 let copybook_text = r"
4408 01 RECORD.
4409 05 ID PIC 9(3).
4410 05 NAME PIC X(5).
4411 ";
4412
4413 let schema = parse_copybook(copybook_text)?;
4414 let options = EncodeOptions::default();
4415
4416 let mut json_obj = serde_json::Map::new();
4417 json_obj.insert("ID".into(), Value::String("123".into()));
4418 json_obj.insert("NAME".into(), Value::String("HELLO".into()));
4419 let json = Value::Object(json_obj);
4420
4421 let result = encode_record(&schema, &json, &options)?;
4422 assert!(!result.is_empty());
4423 assert_eq!(result.len(), 8); Ok(())
4427 }
4428
4429 #[test]
4430 fn rdw_odo_length_preserves_storage_after_nested_group() -> Result<()> {
4431 let copybook_text = r"
4432 01 RECORD.
4433 05 HEADER PIC X(1).
4434 05 WRAP.
4435 10 CNT PIC 9(3).
4436 10 ITEMS OCCURS 1 TO 5 DEPENDING ON CNT PIC X(4).
4437 05 TRAILER PIC X(2).
4438 ";
4439 let schema = parse_copybook(copybook_text)?;
4440 let json = serde_json::json!({
4441 "HEADER": "H",
4442 "WRAP": {"CNT": "002", "ITEMS": ["ABCD", "WXYZ"]},
4443 "TRAILER": "TT"
4444 });
4445
4446 assert_eq!(rdw_record_length_for_json(&schema, &json), Some(14));
4447 Ok(())
4448 }
4449
4450 #[test]
4451 fn test_record_iterator() -> Result<()> {
4452 let copybook_text = r"
4453 01 RECORD.
4454 05 ID PIC 9(3).
4455 05 NAME PIC X(5).
4456 ";
4457
4458 let schema = parse_copybook(copybook_text)?;
4459 let options = DecodeOptions::default();
4460
4461 let test_data = vec![0u8; 16]; let cursor = Cursor::new(test_data);
4464
4465 let iterator = RecordIterator::new(cursor, &schema, &options)?;
4466 assert_eq!(iterator.current_record_index(), 0);
4467 assert!(!iterator.is_eof());
4468 Ok(())
4469 }
4470
4471 #[test]
4472 fn test_decode_file_to_jsonl() -> Result<()> {
4473 let copybook_text = r"
4474 01 RECORD.
4475 05 ID PIC 9(3).
4476 05 NAME PIC X(5).
4477 ";
4478
4479 let schema = parse_copybook(copybook_text)?;
4480 let options = DecodeOptions {
4481 codepage: Codepage::ASCII, ..DecodeOptions::default()
4483 };
4484
4485 let input_data = b"001ALICE002BOBBY".to_vec(); let input = Cursor::new(input_data);
4488
4489 let mut output = Vec::new();
4491
4492 let summary = decode_file_to_jsonl(&schema, input, &mut output, &options)?;
4493 assert!(summary.records_processed > 0);
4494 assert!(!output.is_empty());
4495 Ok(())
4496 }
4497
4498 #[test]
4499 fn test_encode_jsonl_to_file() -> Result<()> {
4500 let copybook_text = r"
4501 01 RECORD.
4502 05 ID PIC 9(3).
4503 05 NAME PIC X(5).
4504 ";
4505
4506 let schema = parse_copybook(copybook_text)?;
4507 let options = EncodeOptions::default();
4508
4509 let jsonl_data = "{\"__status\":\"test\"}\n{\"__status\":\"test2\"}";
4511 let input = Cursor::new(jsonl_data.as_bytes());
4512
4513 let mut output = Vec::new();
4515
4516 let summary = encode_jsonl_to_file(&schema, input, &mut output, &options)?;
4517 assert_eq!(summary.records_processed, 2);
4518 assert!(!output.is_empty());
4519 Ok(())
4520 }
4521}