Skip to main content

datui_lib/formats/fix/
mod.rs

1//! FIX logs: `tag=value` messages, SOH or `|` delimited, read into a table of one row
2//! per message.
3//!
4//! Each tag in the file is a column, named from the FIX dictionaries ([`dict`]): `35`
5//! is `MsgType`, `55` is `Symbol`. An enumerated tag shows its names (`54=1` is `Buy`)
6//! with the code beside it in `<name>_code`; a tag repeated within a message (a
7//! repeating group) keeps its first value and lists the rest in `<name>_rest`. Prices
8//! and quantities become numbers and `52`/`60` timestamps datetimes when every value
9//! reads as one. A line's text before `8=FIX` is kept as `prefix`, with the direction
10//! and session it names. `body_length_ok` and `checksum_ok` check tags 9 and 10.
11//!
12//! The log is read once, a piece at a time, into temporary Arrow IPC segments; a
13//! message, a value and the number of columns are bounded.
14
15pub 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
30/// What datui does with a FIX log: see [`crate::formats::readers`].
31pub(crate) const READER: crate::formats::readers::Reader = crate::formats::readers::Reader {
32    // The dictionaries are read before the file, so one that does not parse says so at
33    // once.
34    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
57/// The longest message read; a longer one is cut there.
58pub const MAX_MESSAGE: usize = 1 << 20;
59/// The longest prefix kept.
60pub const MAX_PREFIX: usize = 4096;
61/// The most distinct names one tag is shown with in the Info panel.
62const MAX_NAMES: usize = 8;
63/// The most distinct counterparties whose dictionaries are remembered.
64const MAX_VIEWS: usize = 1024;
65/// Rows held before a batch is handed over.
66pub const BATCH_ROWS: usize = 32_768;
67/// Text held before a batch is handed over, whatever its rows. Each cell counts
68/// too, empty or not: every tag seen is a column in every row, so a log of
69/// thousands of tags would otherwise hold gigabytes of empty cells in one batch.
70pub const BATCH_TEXT: usize = 32 << 20;
71
72/// The columns that are not tags, in front of them and after them.
73pub const FRONT: [&str; 3] = ["prefix", "direction", "session"];
74pub const BACK: [&str; 2] = ["body_length_ok", "checksum_ok"];
75
76/// Where `8=FIX` starts in `line`, where it starts a tag: not after a digit.
77fn 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
92/// Whether a value byte can be part of a BeginString: `FIX.4.4`, `FIXT.1.1`.
93fn begin_byte(b: u8) -> bool {
94    b.is_ascii_alphanumeric() || b == b'.'
95}
96
97/// The delimiter after a BeginString's value: SOH, `|`, `^A`, or another byte.
98fn 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
111/// Whether the first bytes hold a FIX message: `8=FIX...`, a delimiter and `9=`.
112pub 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/// One message as read.
125#[derive(Debug, Default)]
126struct Message {
127    fields: Vec<(u32, String)>,
128    /// Bytes read, from `8=` on.
129    len: usize,
130    /// Whether it ended with tag 10.
131    complete: bool,
132    body_length_ok: Option<bool>,
133    checksum_ok: Option<bool>,
134    /// Fields past `limits.fix_fields`.
135    dropped: u64,
136}
137
138/// Read one message from `bytes`, which start at `8=`. `None` when more is needed.
139/// Fields past `max_fields` are left out.
140fn 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
255/// The lines of `text` (newline-separated) that are not blank.
256fn 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
262/// `in` or `out`, from the words of a log line's prefix.
263fn 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
275/// A session a prefix names: `FIX.4.4:SENDER->TARGET`, or any `A->B`.
276fn 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/// What reading noticed.
288#[derive(Debug, Clone, Default, PartialEq)]
289pub struct Stats {
290    pub messages: u64,
291    /// Lines with no FIX message, passed over.
292    pub skipped_lines: u64,
293    /// Messages that end before tag 10.
294    pub incomplete: u64,
295    pub body_length_failed: u64,
296    pub checksum_failed: u64,
297    /// Values of tags past `limits.fix_tags`.
298    pub tags_dropped: u64,
299    /// Fields past `limits.fix_fields` in a message.
300    pub fields_dropped: u64,
301    /// Messages cut at [`MAX_MESSAGE`].
302    pub cut: u64,
303    /// BeginStrings and how many messages each.
304    pub versions: Vec<(String, u64)>,
305    /// Per dictionary, the messages it applied to.
306    pub applied: Vec<u64>,
307}
308
309/// One tag's columns and what its values were.
310#[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    /// The names the dictionaries gave it, each with the dictionary.
317    names: Vec<(String, usize)>,
318    ty: Option<FixType>,
319    /// The dictionaries gave it two types.
320    types_differ: bool,
321    /// Every value read as `ty`.
322    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    /// The type its column is cast to, if any: a tag with enums keeps its names.
340    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/// The dictionaries one counterparty's messages are read with, and what they said.
350#[derive(Debug, Default)]
351struct View {
352    applying: Vec<usize>,
353    resolved: HashMap<u32, Arc<Resolved>>,
354}
355
356/// Reads a FIX log a piece at a time.
357#[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    /// `limits.fix_tags` and `limits.fix_fields`, read once for the file.
377    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    /// Read `bytes`, the next piece of the file.
419    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    /// A full batch, once one is held: [`BATCH_ROWS`] rows, or [`BATCH_TEXT`] of text.
429    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    /// The end of the file: what is left, and the rows not yet taken.
437    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            // The next message, and the whole lines before it that hold none. Found
448            // before any newline is looked for, so a capture of messages back to back,
449            // with no newlines at all, is read in one pass.
450            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                    // A line with no `8=FIX` in its first MiB: passed over, keeping
462                    // what could be the start of one.
463                    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        // Each tag's first value, and the rest when it repeats.
558        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                // A batch of nulls only is a list of nulls; every batch's type must match.
697                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    /// Each tag and the names the dictionaries gave it, each with the dictionary's index
723    /// in [`Self::layers`].
724    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    /// Each tag's column name: the name every dictionary that applied gave it, or its
732    /// number when none did, when they disagree, or when the name is taken.
733    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            // A name of digits or with `#` could meet another column's raw name.
742            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    /// `batch`, the last, which has every column, as the dataset shows it: for checks.
761    pub fn finished(&self, batch: DataFrame) -> PolarsResult<DataFrame> {
762        self.finish_frame(batch.lazy()).collect()
763    }
764
765    /// The segments' columns renamed, typed and, for those that never held a value,
766    /// dropped.
767    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
817/// A string column as `ty`.
818fn 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
840/// The element type a repeated tag's list is cast to: numbers only.
841fn 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
941/// The FIX tab of the Info panel: versions, dictionaries and each tag's number.
942pub 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
1004/// The dictionaries a log is read with: the search path's, then each `--dict`.
1005pub 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                // A TOML dictionary's error may be in the file it names.
1020                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    /// A log with no messages names itself; a dictionary that is not one names the
1048    /// dictionary and the flag that took it.
1049    #[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    /// A message of `body` (fields after 9, before 10) with a correct BodyLength and
1082    /// CheckSum, joined by `delim`.
1083    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    /// Empty cells count toward a batch: one message of 2,000 tags makes every later
1119    /// row 2,000 cells wide, so the batch is handed over long before `BATCH_ROWS`.
1120    #[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}