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