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