1pub mod dict;
16
17use std::collections::HashMap;
18use std::path::{Path, PathBuf};
19use std::sync::Arc;
20
21use crate::error_display::FileError;
22use color_eyre::Result;
23use polars::prelude::*;
24
25use crate::OpenOptions;
26use crate::model_files::MetaValue;
27use crate::notes::Note;
28use crate::segments::{Converted, Segments};
29use crate::text_formats::{Detail, Pieces, capped_list, count, note};
30use crate::unfinished::Writer;
31use dict::{FixType, Layers, Resolved};
32
33pub(crate) const READER: crate::readers::Reader = crate::readers::Reader {
35 convert: Some(|input| {
38 let layers = layers(input.formats, &input.options.dicts)?;
39 crate::text_formats::read_one(input, |pieces| {
40 convert(input.display, input.options, layers, input.writer, pieces)
41 })
42 }),
43 scan: crate::readers::read_into,
44 signatures: &[crate::readers::Signature {
45 says: |head, _| looks_like(head),
46 kind: crate::readers::Kind::Text,
47 trusted: crate::readers::Trusted {
48 listing: false,
49 ..crate::readers::EVERYWHERE
50 },
51 }],
52 ..crate::readers::BASE
53};
54
55pub const MAX_MESSAGE: usize = 1 << 20;
57pub const MAX_FIELDS: usize = 4096;
59pub const MAX_TAGS: usize = 4096;
61pub const MAX_PREFIX: usize = 4096;
63const MAX_NAMES: usize = 8;
65const MAX_VIEWS: usize = 1024;
67pub const BATCH_ROWS: usize = 32_768;
69pub const BATCH_TEXT: usize = 32 << 20;
73
74pub const FRONT: [&str; 3] = ["prefix", "direction", "session"];
76pub const BACK: [&str; 2] = ["body_length_ok", "checksum_ok"];
77
78fn find_begin(line: &[u8]) -> Option<usize> {
80 let mut from = 0;
81 while let Some(at) = line[from..]
82 .windows(5)
83 .position(|w| w == b"8=FIX")
84 .map(|i| from + i)
85 {
86 if at == 0 || !line[at - 1].is_ascii_digit() {
87 return Some(at);
88 }
89 from = at + 1;
90 }
91 None
92}
93
94fn begin_byte(b: u8) -> bool {
96 b.is_ascii_alphanumeric() || b == b'.'
97}
98
99fn delimiter_at(rest: &[u8]) -> Option<&'static [u8]> {
101 match rest {
102 [b'^', b'A', ..] => Some(b"^A"),
103 [0x01, ..] => Some(b"\x01"),
104 [b'|', ..] => Some(b"|"),
105 [b';', ..] => Some(b";"),
106 [b',', ..] => Some(b","),
107 [b'\t', ..] => Some(b"\t"),
108 [b' ', ..] => Some(b" "),
109 _ => None,
110 }
111}
112
113pub fn looks_like(head: &[u8]) -> bool {
115 head.split(|&b| b == b'\n').any(|line| {
116 let Some(at) = find_begin(line) else {
117 return false;
118 };
119 let value = &line[at + 2..];
120 let len = value.iter().take_while(|&&b| begin_byte(b)).count();
121 let after = &value[len..];
122 delimiter_at(after).is_some_and(|d| after[d.len()..].starts_with(b"9="))
123 })
124}
125
126#[derive(Debug, Default)]
128struct Message {
129 fields: Vec<(u32, String)>,
130 len: usize,
132 complete: bool,
134 body_length_ok: Option<bool>,
135 checksum_ok: Option<bool>,
136 dropped: u64,
138}
139
140fn parse_message(bytes: &[u8], eof: bool, layers: &Layers) -> Option<Message> {
142 let mut m = Message::default();
143 let mut i = 0;
144 let mut delim: Option<&'static [u8]> = None;
145 let mut sized: Option<(u32, usize)> = None;
146 let (mut sum, mut body) = (0u64, 0usize);
147 let (mut declared_len, mut declared_sum) = (None::<usize>, None::<u64>);
148 let mut after_length = false;
149 loop {
150 let start = i;
151 while i < bytes.len() && bytes[i].is_ascii_digit() && i - start < 10 {
152 i += 1;
153 }
154 if i == bytes.len() {
155 if !eof {
156 return None;
157 }
158 i = start;
159 break;
160 }
161 if i == start || bytes[i] != b'=' {
162 i = start;
163 break;
164 }
165 let tag = std::str::from_utf8(&bytes[start..i])
166 .ok()
167 .and_then(|t| t.parse::<u32>().ok())
168 .filter(|t| (1..=dict::MAX_TAG).contains(t));
169 let Some(tag) = tag else {
170 i = start;
171 break;
172 };
173 let value_start = i + 1;
174 let value_end = match sized.take().filter(|(data, _)| *data == tag) {
175 Some((_, n)) => match value_start.checked_add(n) {
176 Some(end) if end <= bytes.len() => end,
177 _ if !eof => return None,
178 _ => bytes.len(),
179 },
180 None => {
181 if delim.is_none() {
182 let len = bytes[value_start..]
183 .iter()
184 .take_while(|&&b| begin_byte(b))
185 .count();
186 let after = &bytes[value_start + len..];
187 if after.len() < 2 && !eof {
188 return None;
189 }
190 delim = delimiter_at(after);
191 }
192 let stop = bytes[value_start..].iter().enumerate().position(|(k, &b)| {
193 b == b'\n'
194 || b == b'\r'
195 || delim.is_some_and(|d| bytes[value_start + k..].starts_with(d))
196 });
197 match stop {
198 Some(k) => value_start + k,
199 None if !eof => return None,
200 None => bytes.len(),
201 }
202 }
203 };
204 let value = &bytes[value_start..value_end];
205 if tag == 10 {
206 declared_sum = std::str::from_utf8(value).ok().and_then(|v| v.parse().ok());
207 } else {
208 sum += bytes[start..value_end]
209 .iter()
210 .map(|&b| u64::from(b))
211 .sum::<u64>()
212 + 1;
213 if after_length {
214 body += value_end - start + 1;
215 }
216 }
217 if tag == 9 && declared_len.is_none() {
218 declared_len = std::str::from_utf8(value).ok().and_then(|v| v.parse().ok());
219 after_length = true;
220 }
221 if let Some(data) = layers.data_tag(tag)
222 && let Some(n) = std::str::from_utf8(value)
223 .ok()
224 .and_then(|v| v.parse::<usize>().ok())
225 .filter(|&n| n <= MAX_MESSAGE)
226 {
227 sized = Some((data, n));
228 }
229 if m.fields.len() < MAX_FIELDS {
230 m.fields
231 .push((tag, String::from_utf8_lossy(value).into_owned()));
232 } else {
233 m.dropped += 1;
234 }
235 i = value_end;
236 let delimited = delim.is_some_and(|d| bytes[i..].starts_with(d));
237 if delimited {
238 i += delim.map_or(0, <[u8]>::len);
239 }
240 if tag == 10 {
241 m.complete = true;
242 break;
243 }
244 if !delimited {
245 break;
246 }
247 }
248 m.len = i;
249 if m.complete {
250 m.body_length_ok = declared_len.map(|d| d == body);
251 m.checksum_ok = declared_sum.map(|d| d == sum % 256);
252 }
253 Some(m)
254}
255
256fn blank_free_lines(text: &[u8]) -> u64 {
258 text.split(|&b| b == b'\n')
259 .filter(|line| !line.trim_ascii().is_empty())
260 .count() as u64
261}
262
263fn direction(prefix: &str) -> Option<&'static str> {
265 prefix
266 .split(|c: char| c.is_whitespace() || matches!(c, '[' | ']' | '(' | ')' | ':' | ','))
267 .find_map(|word| match word.to_ascii_uppercase().as_str() {
268 "IN" | "INCOMING" | "INBOUND" | "RECV" | "RECEIVED" | "RCV" | "RX" | "<" | "<<"
269 | "<-" | "<<<" => Some("in"),
270 "OUT" | "OUTGOING" | "OUTBOUND" | "SENT" | "SEND" | "SND" | "TX" | ">" | ">>"
271 | "->" | ">>>" => Some("out"),
272 _ => None,
273 })
274}
275
276fn session(prefix: &str) -> Option<String> {
278 prefix
279 .split(|c: char| c.is_whitespace() || matches!(c, '[' | ']' | '(' | ')' | ','))
280 .map(|w| w.trim_end_matches(':'))
281 .find(|w| {
282 w.split_once("->")
283 .is_some_and(|(a, b)| !a.is_empty() && !b.is_empty() && !b.contains("->"))
284 })
285 .map(str::to_string)
286}
287
288#[derive(Debug, Clone, Default, PartialEq)]
290pub struct Stats {
291 pub messages: u64,
292 pub skipped_lines: u64,
294 pub incomplete: u64,
296 pub body_length_failed: u64,
297 pub checksum_failed: u64,
298 pub tags_dropped: u64,
300 pub fields_dropped: u64,
302 pub cut: u64,
304 pub versions: Vec<(String, u64)>,
306 pub applied: Vec<u64>,
308}
309
310#[derive(Debug)]
312struct Tag {
313 tag: u32,
314 values: Vec<Option<String>>,
315 codes: Option<Vec<Option<String>>>,
316 rest: Option<Vec<Option<Vec<String>>>>,
317 names: Vec<(String, usize)>,
319 ty: Option<FixType>,
320 types_differ: bool,
322 all_read: bool,
324}
325
326impl Tag {
327 fn new(tag: u32, rows: usize) -> Self {
328 Self {
329 tag,
330 values: vec![None; rows],
331 codes: None,
332 rest: None,
333 names: Vec::new(),
334 ty: None,
335 types_differ: false,
336 all_read: true,
337 }
338 }
339
340 fn cast(&self) -> Option<FixType> {
342 if self.codes.is_some() || self.types_differ || !self.all_read {
343 return None;
344 }
345 self.ty
346 .filter(|t| !matches!(t, FixType::Text | FixType::Data))
347 }
348}
349
350#[derive(Debug, Default)]
352struct View {
353 applying: Vec<usize>,
354 resolved: HashMap<u32, Arc<Resolved>>,
355}
356
357#[derive(Debug)]
359pub struct FixReader {
360 layers: Layers,
361 buf: Vec<u8>,
362 pos: usize,
363 tags: Vec<Tag>,
364 by_tag: HashMap<u32, usize>,
365 prefix: Vec<Option<String>>,
366 direction: Vec<Option<&'static str>>,
367 session: Vec<Option<String>>,
368 body_ok: Vec<Option<bool>>,
369 checksum_ok: Vec<Option<bool>>,
370 rows: usize,
371 held: usize,
372 views: HashMap<[Option<String>; 3], View>,
373 stats: Stats,
374 seen_prefix: bool,
375 seen_direction: bool,
376 seen_session: bool,
377}
378
379impl FixReader {
380 pub fn new(layers: Layers) -> Self {
381 let applied = vec![0; layers.dicts.len()];
382 Self {
383 layers,
384 buf: Vec::new(),
385 pos: 0,
386 tags: Vec::new(),
387 by_tag: HashMap::new(),
388 prefix: Vec::new(),
389 direction: Vec::new(),
390 session: Vec::new(),
391 body_ok: Vec::new(),
392 checksum_ok: Vec::new(),
393 rows: 0,
394 held: 0,
395 views: HashMap::new(),
396 stats: Stats {
397 applied,
398 ..Stats::default()
399 },
400 seen_prefix: false,
401 seen_direction: false,
402 seen_session: false,
403 }
404 }
405
406 pub fn stats(&self) -> &Stats {
407 &self.stats
408 }
409
410 pub fn layers(&self) -> &Layers {
411 &self.layers
412 }
413
414 pub fn push(&mut self, bytes: &[u8]) {
416 if self.pos > 0 {
417 self.buf.drain(..self.pos);
418 self.pos = 0;
419 }
420 self.buf.extend_from_slice(bytes);
421 self.scan(false);
422 }
423
424 pub fn take_batch(&mut self) -> PolarsResult<Option<DataFrame>> {
426 if self.rows < BATCH_ROWS && self.held < BATCH_TEXT {
427 return Ok(None);
428 }
429 self.batch().map(Some)
430 }
431
432 pub fn finish(&mut self) -> PolarsResult<DataFrame> {
434 self.scan(true);
435 self.buf.clear();
436 self.pos = 0;
437 self.batch()
438 }
439
440 fn scan(&mut self, eof: bool) {
441 while self.pos < self.buf.len() {
442 let rest = &self.buf[self.pos..];
443 let Some(at) = find_begin(rest) else {
447 let last = rest.iter().rposition(|&b| b == b'\n');
448 match last {
449 Some(end) => {
450 self.stats.skipped_lines += blank_free_lines(&rest[..end]);
451 self.pos += end + 1;
452 }
453 None if eof => {
454 self.stats.skipped_lines += blank_free_lines(rest);
455 self.pos = self.buf.len();
456 }
457 None if rest.len() > MAX_MESSAGE => {
460 self.stats.skipped_lines += 1;
461 self.pos += rest.len() - 4;
462 }
463 None => return,
464 }
465 continue;
466 };
467 let line_start = rest[..at]
468 .iter()
469 .rposition(|&b| b == b'\n')
470 .map_or(0, |i| i + 1);
471 if line_start > 0 {
472 self.stats.skipped_lines += blank_free_lines(&rest[..line_start - 1]);
473 self.pos += line_start;
474 continue;
475 }
476 let long = rest.len() - at > MAX_MESSAGE;
477 let window = &rest[at..rest.len().min(at + MAX_MESSAGE)];
478 let Some(message) = parse_message(window, eof || long, &self.layers) else {
479 return;
480 };
481 if long && message.len == window.len() {
482 self.stats.cut += 1;
483 }
484 let prefix = String::from_utf8_lossy(&rest[..at]).into_owned();
485 let len = at + message.len.max(2);
486 self.message(prefix, message);
487 self.pos += len;
488 }
489 }
490
491 fn message(&mut self, prefix: String, m: Message) {
492 let prefix = prefix.trim().trim_end_matches(':').trim_end();
493 let mut cut = prefix.len().min(MAX_PREFIX);
494 while !prefix.is_char_boundary(cut) {
495 cut -= 1;
496 }
497 let prefix = &prefix[..cut];
498 let dir = direction(prefix);
499 let sess = session(prefix);
500 self.seen_prefix |= !prefix.is_empty();
501 self.seen_direction |= dir.is_some();
502 self.seen_session |= sess.is_some();
503 self.held += prefix.len();
504 self.prefix
505 .push(Some(prefix.to_string()).filter(|p| !p.is_empty()));
506 self.direction.push(dir);
507 self.session.push(sess);
508 self.body_ok.push(m.body_length_ok);
509 self.checksum_ok.push(m.checksum_ok);
510 self.stats.messages += 1;
511 self.stats.fields_dropped += m.dropped;
512 if !m.complete {
513 self.stats.incomplete += 1;
514 }
515 if m.body_length_ok == Some(false) {
516 self.stats.body_length_failed += 1;
517 }
518 if m.checksum_ok == Some(false) {
519 self.stats.checksum_failed += 1;
520 }
521
522 let first = |tag: u32| {
523 m.fields
524 .iter()
525 .find(|(t, _)| *t == tag)
526 .map(|(_, v)| v.clone())
527 };
528 let key = [first(49), first(56), first(8)];
529 if let Some(begin) = &key[2] {
530 let versions = &mut self.stats.versions;
531 match versions.iter().position(|(b, _)| b == begin) {
532 Some(i) => versions[i].1 += 1,
533 None if versions.len() < 16 => versions.push((begin.clone(), 1)),
534 None => {}
535 }
536 }
537 let mut view = match self.views.remove(&key) {
538 Some(view) => view,
539 None => View {
540 applying: self.layers.applying(
541 key[0].as_deref(),
542 key[1].as_deref(),
543 key[2].as_deref(),
544 ),
545 resolved: HashMap::new(),
546 },
547 };
548 for &i in &view.applying {
549 self.stats.applied[i] += 1;
550 }
551
552 let mut order: Vec<(u32, String, Vec<String>)> = Vec::new();
554 let mut at: HashMap<u32, usize> = HashMap::new();
555 for (tag, value) in m.fields {
556 match at.get(&tag) {
557 Some(&i) => order[i].2.push(value),
558 None => {
559 at.insert(tag, order.len());
560 order.push((tag, value, Vec::new()));
561 }
562 }
563 }
564 let row = self.rows;
565 for (tag, value, rest) in order {
566 let index = match self.by_tag.get(&tag) {
567 Some(&i) => i,
568 None if self.tags.len() >= MAX_TAGS => {
569 self.stats.tags_dropped += 1;
570 continue;
571 }
572 None => {
573 self.by_tag.insert(tag, self.tags.len());
574 self.tags.push(Tag::new(tag, row));
575 self.tags.len() - 1
576 }
577 };
578 let resolved = view
579 .resolved
580 .entry(tag)
581 .or_insert_with(|| Arc::new(self.layers.resolve(&view.applying, tag)))
582 .clone();
583 let t = &mut self.tags[index];
584 if let (Some(name), Some(by)) = (&resolved.name, resolved.named_by)
585 && !t.names.iter().any(|(n, _)| n == name)
586 && t.names.len() < MAX_NAMES
587 {
588 t.names.push((name.clone(), by));
589 }
590 if let Some(ty) = resolved.ty {
591 match t.ty {
592 None => t.ty = Some(ty),
593 Some(had) if had != ty => t.types_differ = true,
594 _ => {}
595 }
596 }
597 if let Some(ty) = t.ty {
598 t.all_read &= ty.reads(&value) && rest.iter().all(|v| ty.reads(v));
599 }
600 let named = |v: &String| resolved.enums.get(v).cloned().unwrap_or_else(|| v.clone());
601 if !resolved.enums.is_empty() && t.codes.is_none() {
602 t.codes = Some(vec![None; row]);
603 }
604 let shown = if resolved.enums.is_empty() {
605 value.clone()
606 } else {
607 named(&value)
608 };
609 self.held += shown.len() + rest.iter().map(String::len).sum::<usize>();
610 if let Some(codes) = t.codes.as_mut() {
611 codes.push(Some(value));
612 }
613 if !rest.is_empty() {
614 let rest: Vec<String> = if resolved.enums.is_empty() {
615 rest
616 } else {
617 rest.iter().map(named).collect()
618 };
619 t.rest
620 .get_or_insert_with(|| vec![None; row])
621 .push(Some(rest));
622 }
623 t.values.push(Some(shown));
624 }
625 self.rows += 1;
626 for t in &mut self.tags {
627 self.held += size_of::<Option<String>>()
628 * (1 + usize::from(t.codes.is_some()) + usize::from(t.rest.is_some()));
629 if t.values.len() < self.rows {
630 t.values.push(None);
631 }
632 if let Some(codes) = t.codes.as_mut()
633 && codes.len() < self.rows
634 {
635 codes.push(None);
636 }
637 if let Some(rest) = t.rest.as_mut()
638 && rest.len() < self.rows
639 {
640 rest.push(None);
641 }
642 }
643 if self.views.len() < MAX_VIEWS {
644 self.views.insert(key, view);
645 }
646 }
647
648 fn batch(&mut self) -> PolarsResult<DataFrame> {
649 self.held = 0;
650 let height = std::mem::take(&mut self.rows);
651 let mut columns = vec![
652 StringChunked::from_iter_options(
653 "prefix".into(),
654 std::mem::take(&mut self.prefix).into_iter(),
655 )
656 .into_column(),
657 StringChunked::from_iter_options(
658 "direction".into(),
659 std::mem::take(&mut self.direction).into_iter(),
660 )
661 .into_column(),
662 StringChunked::from_iter_options(
663 "session".into(),
664 std::mem::take(&mut self.session).into_iter(),
665 )
666 .into_column(),
667 ];
668 for t in &mut self.tags {
669 let name = t.tag.to_string();
670 columns.push(
671 StringChunked::from_iter_options(
672 name.as_str().into(),
673 std::mem::take(&mut t.values).into_iter(),
674 )
675 .into_column(),
676 );
677 if let Some(codes) = t.codes.as_mut() {
678 columns.push(
679 StringChunked::from_iter_options(
680 format!("{name}#code").into(),
681 std::mem::take(codes).into_iter(),
682 )
683 .into_column(),
684 );
685 }
686 if let Some(rest) = t.rest.as_mut() {
687 let list: ListChunked = std::mem::take(rest)
688 .into_iter()
689 .map(|items| items.map(|items| Series::new(PlSmallStr::EMPTY, items)))
690 .collect();
691 columns.push(
693 list.with_name(format!("{name}#rest").into())
694 .into_series()
695 .cast(&DataType::List(Box::new(DataType::String)))?
696 .into_column(),
697 );
698 }
699 }
700 columns.push(
701 BooleanChunked::from_iter_options(
702 "body_length_ok".into(),
703 std::mem::take(&mut self.body_ok).into_iter(),
704 )
705 .into_column(),
706 );
707 columns.push(
708 BooleanChunked::from_iter_options(
709 "checksum_ok".into(),
710 std::mem::take(&mut self.checksum_ok).into_iter(),
711 )
712 .into_column(),
713 );
714 DataFrame::new(height, columns)
715 }
716
717 pub fn tag_names(&self) -> Vec<(u32, String, usize)> {
720 self.tags
721 .iter()
722 .flat_map(|t| t.names.iter().map(|(n, by)| (t.tag, n.clone(), *by)))
723 .collect()
724 }
725
726 pub fn column_names(&self) -> Vec<(u32, String)> {
729 let mut taken: Vec<String> = FRONT
730 .iter()
731 .chain(BACK.iter())
732 .map(|s| s.to_string())
733 .collect();
734 let mut out = Vec::with_capacity(self.tags.len());
735 for t in &self.tags {
736 let usable = |name: &String| {
738 !taken.contains(name)
739 && !taken.contains(&format!("{name}_code"))
740 && !name.contains('#')
741 && !name.bytes().all(|b| b.is_ascii_digit())
742 };
743 let name = match t.names.as_slice() {
744 [(name, _)] if usable(name) => name.clone(),
745 _ => t.tag.to_string(),
746 };
747 taken.push(name.clone());
748 taken.push(format!("{name}_code"));
749 taken.push(format!("{name}_rest"));
750 out.push((t.tag, name));
751 }
752 out
753 }
754
755 pub fn finished(&self, batch: DataFrame) -> PolarsResult<DataFrame> {
757 self.finish_frame(batch.lazy()).collect()
758 }
759
760 pub fn finish_frame(&self, lf: LazyFrame) -> LazyFrame {
763 let mut drop: Vec<&str> = Vec::new();
764 if !self.seen_prefix {
765 drop.push("prefix");
766 }
767 if !self.seen_direction {
768 drop.push("direction");
769 }
770 if !self.seen_session {
771 drop.push("session");
772 }
773 let mut casts = Vec::new();
774 let (mut from, mut to) = (Vec::new(), Vec::new());
775 for (t, (_, name)) in self.tags.iter().zip(self.column_names()) {
776 let raw = t.tag.to_string();
777 if let Some(ty) = t.cast() {
778 casts.push(cast(col(raw.as_str()), ty).alias(raw.as_str()));
779 if t.rest.is_some()
780 && let Some(inner) = list_type(ty)
781 {
782 let rest = format!("{raw}#rest");
783 casts.push(
784 col(rest.as_str())
785 .cast(DataType::List(Box::new(inner)))
786 .alias(rest.as_str()),
787 );
788 }
789 }
790 from.push(raw.clone());
791 to.push(name.clone());
792 if t.codes.is_some() {
793 from.push(format!("{raw}#code"));
794 to.push(format!("{name}_code"));
795 }
796 if t.rest.is_some() {
797 from.push(format!("{raw}#rest"));
798 to.push(format!("{name}_rest"));
799 }
800 }
801 let mut lf = lf;
802 if !drop.is_empty() {
803 lf = lf.drop(cols(drop));
804 }
805 if !casts.is_empty() {
806 lf = lf.with_columns(casts);
807 }
808 lf.rename(from, to, true)
809 }
810}
811
812fn cast(expr: Expr, ty: FixType) -> Expr {
814 match ty {
815 FixType::Int | FixType::Length => expr.cast(DataType::Int64),
816 FixType::Float => expr.cast(DataType::Float64),
817 FixType::Bool => expr.eq(lit("Y")),
818 FixType::Timestamp => expr.str().to_datetime(
819 Some(TimeUnit::Nanoseconds),
820 Some(TimeZone::UTC),
821 StrptimeOptions {
822 format: Some("%Y%m%d-%H:%M:%S%.f".into()),
823 ..Default::default()
824 },
825 lit("raise"),
826 ),
827 FixType::Date => expr.str().to_date(StrptimeOptions {
828 format: Some("%Y%m%d".into()),
829 ..Default::default()
830 }),
831 FixType::Data | FixType::Text => expr,
832 }
833}
834
835fn list_type(ty: FixType) -> Option<DataType> {
837 match ty {
838 FixType::Int | FixType::Length => Some(DataType::Int64),
839 FixType::Float => Some(DataType::Float64),
840 _ => None,
841 }
842}
843
844fn notes(reader: &FixReader) -> Vec<Note> {
845 let stats = reader.stats();
846 let mut notes = Vec::new();
847 let of_messages = format!("of {}", count(stats.messages, "message", "messages"));
848 if stats.skipped_lines > 0 {
849 notes.push(note(
850 format!(
851 "{} left out: no FIX message",
852 count(stats.skipped_lines, "line", "lines")
853 ),
854 "in the whole file".to_string(),
855 ));
856 }
857 if stats.incomplete > 0 {
858 notes.push(note(
859 format!(
860 "ended before tag 10: {} {} checks null",
861 count(stats.incomplete, "message", "messages"),
862 crate::glyphs::get().middot
863 ),
864 of_messages.clone(),
865 ));
866 }
867 if stats.body_length_failed > 0 {
868 notes.push(note(
869 format!(
870 "BodyLength (9) failed: {} {} body_length_ok false",
871 count(stats.body_length_failed, "message", "messages"),
872 crate::glyphs::get().middot
873 ),
874 of_messages.clone(),
875 ));
876 }
877 if stats.checksum_failed > 0 {
878 notes.push(note(
879 format!(
880 "CheckSum (10) failed: {} {} checksum_ok false",
881 count(stats.checksum_failed, "message", "messages"),
882 crate::glyphs::get().middot
883 ),
884 of_messages.clone(),
885 ));
886 }
887 if stats.tags_dropped > 0 || stats.fields_dropped > 0 {
888 notes.push(note(
889 format!(
890 "{} left out: past limits ({MAX_TAGS} tags, {MAX_FIELDS} fields a message)",
891 count(stats.tags_dropped + stats.fields_dropped, "value", "values")
892 ),
893 of_messages.clone(),
894 ));
895 }
896 if stats.cut > 0 {
897 notes.push(note(
898 format!(
899 "{} cut at {} MiB",
900 count(stats.cut, "message", "messages"),
901 MAX_MESSAGE >> 20
902 ),
903 of_messages.clone(),
904 ));
905 }
906 let differ: Vec<String> = reader
907 .tags
908 .iter()
909 .filter(|t| t.names.len() > 1)
910 .map(|t| t.tag.to_string())
911 .collect();
912 if !differ.is_empty() {
913 notes.push(note(
914 format!(
915 "dictionaries disagree on {} {} named by number (FIX tab has the names)",
916 count(differ.len() as u64, "tag", "tags"),
917 crate::glyphs::get().middot
918 ),
919 format!("tags {}", differ.join(", ")),
920 ));
921 }
922 notes
923}
924
925pub fn detail(reader: &FixReader) -> Detail {
927 let stats = reader.stats();
928 let sep = format!(" {} ", crate::glyphs::get().middot);
929 let mut head = format!("FIX{sep}{}", count(stats.messages, "message", "messages"));
930 for (begin, n) in &stats.versions {
931 head.push_str(&sep);
932 head.push_str(&format!(
933 "{begin} ({})",
934 crate::numfmt::group_chrome(*n as usize)
935 ));
936 }
937 let mut lines = vec![head];
938 let dicts = &reader.layers().dicts;
939 let mut used = vec!["built-in FIX 4.2, 4.4, 5.0 SP2".to_string()];
940 for (i, d) in dicts.iter().enumerate().skip(1) {
941 let mut said = d.name.clone();
942 let summary = d.matcher.summary();
943 if !summary.is_empty() {
944 said.push_str(&format!(" ({summary})"));
945 }
946 said.push_str(&format!(
947 ", {}",
948 count(
949 stats.applied.get(i).copied().unwrap_or(0),
950 "message",
951 "messages"
952 )
953 ));
954 used.push(said);
955 }
956 lines.push(format!("Dictionaries: {}", used.join(&sep)));
957 let names = reader.column_names();
958 let list = capped_list(
959 reader.tags.iter().zip(names).map(|(t, (_, column))| {
960 let said = match t.names.as_slice() {
961 [] => format!("tag {}{sep}in no dictionary", t.tag),
962 [(name, _)] => format!("tag {}{sep}{name}", t.tag),
963 many => {
964 let each: Vec<String> = many
965 .iter()
966 .map(|(name, by)| {
967 let dict = dicts.get(*by).map_or("?", |d| d.name.as_str());
968 format!("{name} ({dict})")
969 })
970 .collect();
971 format!("tag {}{sep}{}", t.tag, each.join(&sep))
972 }
973 };
974 (column, MetaValue::Text(said))
975 }),
976 reader.tags.len(),
977 );
978 Detail {
979 tab: crate::text_formats::tab(crate::FileFormat::Fix),
980 lines,
981 list_title: "Tags",
982 list,
983 first: false,
984 ..Default::default()
985 }
986}
987
988pub fn layers(registry: &crate::formats::Registry, dicts: &[PathBuf]) -> Result<Layers> {
990 let mut custom: Vec<Arc<dict::Dictionary>> =
991 registry.fix.iter().map(|f| f.dict.clone()).collect();
992 for path in dicts {
993 match dict::Dictionary::load(path) {
994 Ok(Some(d)) => custom.push(Arc::new(d)),
995 Ok(None) => {
996 return Err(FileError::new(
997 path,
998 "not a QuickFIX dictionary. --dict takes a QuickFIX XML file, or TOML with kind = \"fix\".",
999 )
1000 .into());
1001 }
1002 Err(e) => {
1003 let at = e.path.as_deref().unwrap_or(path);
1005 return Err(FileError::at(at, e.line, e.column, e.message).into());
1006 }
1007 }
1008 }
1009 Ok(Layers::new(custom))
1010}
1011
1012pub(crate) fn convert(
1015 display: &Path,
1016 options: &OpenOptions,
1017 layers: Layers,
1018 writer: &Writer,
1019 pieces: &mut Pieces<'_>,
1020) -> Result<(Converted, Detail)> {
1021 let mut reader = FixReader::new(layers);
1022 let mut segments = Segments::new(options, writer);
1023 pieces(&mut |piece| {
1024 reader.push(piece);
1025 if let Some(df) = reader.take_batch()? {
1026 segments.write(&df)?;
1027 }
1028 Ok(())
1029 })?;
1030 let last = reader.finish()?;
1031 if reader.stats().messages == 0 {
1032 return Err(FileError::new(display, "no FIX messages: no line holds 8=FIX.").into());
1033 }
1034 segments.write(&last)?;
1035 let (lf, files) = segments.finish()?;
1036 let lf = reader.finish_frame(lf);
1037 Ok((
1038 Converted {
1039 lf,
1040 files,
1041 notes: notes(&reader),
1042 other_tables: Vec::new(),
1043 },
1044 detail(&reader),
1045 ))
1046}
1047
1048#[cfg(test)]
1049mod tests {
1050 use super::*;
1051
1052 #[test]
1055 fn errors_name_the_file() {
1056 use crate::readers::bad_input::{assert_shape, each_names_its_file, opening};
1057 each_names_its_file(
1058 crate::FileFormat::Fix,
1059 &[("words.fix", b"hello there\n", "No FIX messages")],
1060 );
1061 let dir = tempfile::tempdir().unwrap();
1062 for (name, text, says) in [
1063 ("plain.toml", "a = 1\n", "--dict takes"),
1064 ("broken.xml", "<fix><fields><field", "broken.xml\":1:"),
1065 ] {
1066 let dict = dir.path().join(name);
1067 std::fs::write(&dict, text).unwrap();
1068 let options = OpenOptions {
1069 dicts: vec![dict.clone()],
1070 ..Default::default()
1071 };
1072 let message = opening(
1073 dir.path(),
1074 "a.fix",
1075 b"8=FIX.4.4\x019=5\x0135=0\x0110=000\x01\n",
1076 crate::FileFormat::Fix,
1077 &options,
1078 )
1079 .expect("the dictionary is refused");
1080 eprintln!("{message}");
1081 assert_shape(&message, &dict);
1082 assert!(message.contains(says), "{message}");
1083 }
1084 }
1085
1086 pub(crate) fn message(begin: &str, body: &[(u32, &str)], delim: &str) -> String {
1089 let body: String = body.iter().map(|(t, v)| format!("{t}={v}\x01")).collect();
1090 let head = format!("8={begin}\x019={}\x01", body.len());
1091 let sum: u32 = head.bytes().chain(body.bytes()).map(u32::from).sum();
1092 format!("{head}{body}10={:03}\x01", sum % 256).replace('\x01', delim)
1093 }
1094
1095 fn read(text: &[u8], piece: usize, layers: Layers) -> (DataFrame, FixReader) {
1096 let mut reader = FixReader::new(layers);
1097 let mut frames = Vec::new();
1098 for chunk in text.chunks(piece) {
1099 reader.push(chunk);
1100 if let Some(df) = reader.take_batch().unwrap() {
1101 frames.push(df);
1102 }
1103 }
1104 frames.push(reader.finish().unwrap());
1105 let mut df = frames.remove(0);
1106 for f in frames {
1107 df.vstack_mut(&f).unwrap();
1108 }
1109 let df = reader.finish_frame(df.lazy()).collect().unwrap();
1110 (df, reader)
1111 }
1112
1113 fn strings(df: &DataFrame, name: &str) -> Vec<Option<String>> {
1114 df.column(name)
1115 .unwrap_or_else(|_| panic!("{name} in {:?}", df.get_column_names()))
1116 .str()
1117 .unwrap()
1118 .iter()
1119 .map(|s| s.map(String::from))
1120 .collect()
1121 }
1122
1123 #[test]
1126 fn a_wide_log_hands_over_small_batches() {
1127 let mut reader = FixReader::new(Layers::default());
1128 let tags: Vec<(u32, String)> = (20_000..22_000).map(|t| (t, "x".to_string())).collect();
1129 let tags: Vec<(u32, &str)> = tags.iter().map(|(t, v)| (*t, v.as_str())).collect();
1130 reader.push(format!("{}\n", message("FIX.4.4", &tags, "\x01")).as_bytes());
1131 let narrow = format!("{}\n", message("FIX.4.4", &[(35, "0")], "\x01"));
1132 let mut rows_at_batch = None;
1133 for row in 1..BATCH_ROWS {
1134 reader.push(narrow.as_bytes());
1135 if let Some(df) = reader.take_batch().unwrap() {
1136 rows_at_batch = Some((row, df.height()));
1137 break;
1138 }
1139 }
1140 let (row, height) = rows_at_batch.expect("a batch is handed over");
1141 assert!(height < BATCH_ROWS / 8, "{height} rows of 2,000 tags");
1142 assert_eq!(height, row + 1);
1143 }
1144
1145 fn log() -> String {
1146 let order = message(
1147 "FIX.4.4",
1148 &[
1149 (35, "D"),
1150 (49, "BUYSIDE"),
1151 (56, "BROKERX"),
1152 (52, "20260102-03:04:05.123"),
1153 (55, "ACME"),
1154 (54, "1"),
1155 (44, "101.25"),
1156 (38, "500"),
1157 ],
1158 "\x01",
1159 );
1160 let fill = message(
1161 "FIX.4.4",
1162 &[
1163 (35, "8"),
1164 (49, "BROKERX"),
1165 (56, "BUYSIDE"),
1166 (52, "20260102-03:04:06"),
1167 (55, "ACME"),
1168 (54, "2"),
1169 (44, "101.5"),
1170 (453, "2"),
1171 (448, "PARTY1"),
1172 (448, "PARTY2"),
1173 ],
1174 "|",
1175 );
1176 format!("{order}\n2026-01-02 03:04:06 FIX.4.4:BROKERX->BUYSIDE IN : {fill}\nnot fix\n")
1177 }
1178
1179 #[test]
1180 fn messages_read_as_rows_named_by_the_dictionary() {
1181 for piece in [1, 7, 4096] {
1182 let (df, reader) = read(log().as_bytes(), piece, Layers::default());
1183 assert_eq!(df.height(), 2, "piece {piece}");
1184 assert_eq!(reader.stats().skipped_lines, 1);
1185 assert_eq!(
1186 strings(&df, "MsgType"),
1187 [
1188 Some("NewOrderSingle".into()),
1189 Some("ExecutionReport".into())
1190 ]
1191 );
1192 assert_eq!(
1193 strings(&df, "MsgType_code"),
1194 [Some("D".into()), Some("8".into())]
1195 );
1196 assert_eq!(
1197 strings(&df, "Side"),
1198 [Some("Buy".into()), Some("Sell".into())]
1199 );
1200 let price: Vec<_> = df.column("Price").unwrap().f64().unwrap().iter().collect();
1201 assert_eq!(price, [Some(101.25), Some(101.5)]);
1202 let qty = df.column("OrderQty").unwrap();
1203 assert_eq!(qty.dtype(), &DataType::Float64);
1204 let sent = df.column("SendingTime").unwrap();
1205 assert_eq!(
1206 sent.dtype(),
1207 &DataType::Datetime(TimeUnit::Nanoseconds, Some(TimeZone::UTC))
1208 );
1209 let ns: Vec<_> = sent.datetime().unwrap().physical().iter().collect();
1210 assert_eq!(ns[0], dict::parse_timestamp("20260102-03:04:05.123"));
1211 assert_eq!(ns[1], dict::parse_timestamp("20260102-03:04:06"));
1212 assert_eq!(strings(&df, "PartyID"), [None, Some("PARTY1".into())]);
1213 let rest = df.column("PartyID_rest").unwrap().list().unwrap();
1214 assert!(rest.get_as_series(0).is_none());
1215 assert_eq!(rest.get_as_series(1).unwrap().len(), 1);
1216 assert_eq!(
1217 strings(&df, "prefix"),
1218 [
1219 None,
1220 Some("2026-01-02 03:04:06 FIX.4.4:BROKERX->BUYSIDE IN".into())
1221 ]
1222 );
1223 assert_eq!(strings(&df, "direction"), [None, Some("in".into())]);
1224 assert_eq!(
1225 strings(&df, "session"),
1226 [None, Some("FIX.4.4:BROKERX->BUYSIDE".into())]
1227 );
1228 for check in ["body_length_ok", "checksum_ok"] {
1229 let ok: Vec<_> = df.column(check).unwrap().bool().unwrap().iter().collect();
1230 assert_eq!(ok, [Some(true), Some(true)], "{check}");
1231 }
1232 }
1233 }
1234
1235 #[test]
1236 fn bad_checks_are_false_and_a_cut_message_null() {
1237 let good = message("FIX.4.2", &[(35, "0")], "|");
1238 let bad_sum = good.replace("10=", "10=9");
1239 let bad_len = good.replace("9=5", "9=6");
1240 let text = format!("{bad_sum}\n{bad_len}\n8=FIX.4.2|9=5|35=0|\n");
1241 let (df, reader) = read(text.as_bytes(), 3, Layers::default());
1242 let sum: Vec<_> = df
1243 .column("checksum_ok")
1244 .unwrap()
1245 .bool()
1246 .unwrap()
1247 .iter()
1248 .collect();
1249 let len: Vec<_> = df
1250 .column("body_length_ok")
1251 .unwrap()
1252 .bool()
1253 .unwrap()
1254 .iter()
1255 .collect();
1256 assert_eq!(sum, [Some(false), Some(false), None]);
1257 assert_eq!(len, [Some(true), Some(false), None]);
1258 assert_eq!(reader.stats().incomplete, 1);
1259 assert!(df.get_column_names().iter().all(|c| *c != "prefix"));
1260 }
1261
1262 #[test]
1263 fn a_length_tagged_value_may_hold_the_delimiter() {
1264 let text = message(
1265 "FIX.4.4",
1266 &[(35, "A"), (95, "7"), (96, "a\x01b|c\nd"), (58, "hi")],
1267 "\x01",
1268 );
1269 let (df, _) = read(text.as_bytes(), 2, Layers::default());
1270 assert_eq!(df.height(), 1);
1271 assert_eq!(strings(&df, "RawData"), [Some("a\x01b|c\nd".into())]);
1272 assert_eq!(strings(&df, "Text"), [Some("hi".into())]);
1273 let ok: Vec<_> = df
1274 .column("checksum_ok")
1275 .unwrap()
1276 .bool()
1277 .unwrap()
1278 .iter()
1279 .collect();
1280 assert_eq!(ok, [Some(true)]);
1281 }
1282
1283 #[test]
1284 fn counterparties_name_custom_tags_their_own_way() {
1285 let x = dict::Dictionary::parse_toml(
1286 "name = \"acme.x\"\nkind = \"fix\"\nmatch = { sender = \"X\" }\ntags = { 9001 = \"AlgoName\", 9002 = { name = \"Urgency\", type = \"int\", enum = { 1 = \"Low\" } } }",
1287 None,
1288 )
1289 .unwrap();
1290 let y = dict::Dictionary::parse_toml(
1291 "name = \"acme.y\"\nkind = \"fix\"\nmatch = { sender = \"Y\" }\ntags = { 9001 = \"Strategy\" }",
1292 None,
1293 )
1294 .unwrap();
1295 let layers = Layers::new(vec![Arc::new(x), Arc::new(y)]);
1296 let text = [
1297 message(
1298 "FIX.4.4",
1299 &[(35, "D"), (49, "X"), (9001, "VWAP"), (9002, "1")],
1300 "|",
1301 ),
1302 message(
1303 "FIX.4.4",
1304 &[(35, "D"), (49, "Y"), (9001, "TWAP"), (9003, "z")],
1305 "|",
1306 ),
1307 ]
1308 .join("\n");
1309 let (df, reader) = read(text.as_bytes(), 5, layers);
1310 assert_eq!(
1311 strings(&df, "9001"),
1312 [Some("VWAP".into()), Some("TWAP".into())]
1313 );
1314 assert_eq!(strings(&df, "Urgency"), [Some("Low".into()), None]);
1315 assert_eq!(strings(&df, "Urgency_code"), [Some("1".into()), None]);
1316 assert_eq!(strings(&df, "9003"), [None, Some("z".into())]);
1317 assert_eq!(reader.stats().applied, [2, 1, 1]);
1318 let detail = detail(&reader);
1319 let row = detail.list.iter().find(|(k, _)| k == "9001").unwrap();
1320 let MetaValue::Text(said) = &row.1 else {
1321 panic!()
1322 };
1323 assert!(
1324 said.contains("AlgoName (acme.x)") && said.contains("Strategy (acme.y)"),
1325 "{said}"
1326 );
1327 assert!(
1328 notes(&reader)
1329 .iter()
1330 .any(|n| n.summary.contains("disagree"))
1331 );
1332 }
1333
1334 #[test]
1335 fn sniffing() {
1336 assert!(looks_like(log().as_bytes()));
1337 assert!(looks_like(b"in: 8=FIX.4.2^A9=12^A35=0^A10=000^A"));
1338 assert!(!looks_like(b"the 8=FIX tag starts a message"));
1339 assert!(!looks_like(b"a,b\n18=FIX.4.4|9=1\n"));
1340 assert_eq!(direction("12:00 OUT >"), Some("out"));
1341 assert_eq!(direction("received"), Some("in"));
1342 assert_eq!(session("[FIX.4.4:A->B]"), Some("FIX.4.4:A->B".into()));
1343 assert_eq!(session("a -> b"), None);
1344 }
1345}