Skip to main content

datui_lib/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::{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
33/// What datui does with a FIX log: see [`crate::readers`].
34pub(crate) const READER: crate::readers::Reader = crate::readers::Reader {
35    // The dictionaries are read before the file, so one that does not parse says so at
36    // once.
37    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
55/// The longest message read; a longer one is cut there.
56pub const MAX_MESSAGE: usize = 1 << 20;
57/// The most fields one message may hold; the rest are left out.
58pub const MAX_FIELDS: usize = 4096;
59/// The most tag columns; tags past this are left out.
60pub const MAX_TAGS: usize = 4096;
61/// The longest prefix kept.
62pub const MAX_PREFIX: usize = 4096;
63/// The most distinct names one tag is shown with in the Info panel.
64const MAX_NAMES: usize = 8;
65/// The most distinct counterparties whose dictionaries are remembered.
66const MAX_VIEWS: usize = 1024;
67/// Rows held before a batch is handed over.
68pub const BATCH_ROWS: usize = 32_768;
69/// Text held before a batch is handed over, whatever its rows. Each cell counts
70/// too, empty or not: every tag seen is a column in every row, so a log of
71/// thousands of tags would otherwise hold gigabytes of empty cells in one batch.
72pub const BATCH_TEXT: usize = 32 << 20;
73
74/// The columns that are not tags, in front of them and after them.
75pub const FRONT: [&str; 3] = ["prefix", "direction", "session"];
76pub const BACK: [&str; 2] = ["body_length_ok", "checksum_ok"];
77
78/// Where `8=FIX` starts in `line`, where it starts a tag: not after a digit.
79fn 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
94/// Whether a value byte can be part of a BeginString: `FIX.4.4`, `FIXT.1.1`.
95fn begin_byte(b: u8) -> bool {
96    b.is_ascii_alphanumeric() || b == b'.'
97}
98
99/// The delimiter after a BeginString's value: SOH, `|`, `^A`, or another byte.
100fn 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
113/// Whether the first bytes hold a FIX message: `8=FIX...`, a delimiter and `9=`.
114pub 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/// One message as read.
127#[derive(Debug, Default)]
128struct Message {
129    fields: Vec<(u32, String)>,
130    /// Bytes read, from `8=` on.
131    len: usize,
132    /// Whether it ended with tag 10.
133    complete: bool,
134    body_length_ok: Option<bool>,
135    checksum_ok: Option<bool>,
136    /// Fields past [`MAX_FIELDS`].
137    dropped: u64,
138}
139
140/// Read one message from `bytes`, which start at `8=`. `None` when more is needed.
141fn 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
256/// The lines of `text` (newline-separated) that are not blank.
257fn 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
263/// `in` or `out`, from the words of a log line's prefix.
264fn 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
276/// A session a prefix names: `FIX.4.4:SENDER->TARGET`, or any `A->B`.
277fn 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/// What reading noticed.
289#[derive(Debug, Clone, Default, PartialEq)]
290pub struct Stats {
291    pub messages: u64,
292    /// Lines with no FIX message, passed over.
293    pub skipped_lines: u64,
294    /// Messages that end before tag 10.
295    pub incomplete: u64,
296    pub body_length_failed: u64,
297    pub checksum_failed: u64,
298    /// Values of tags past [`MAX_TAGS`].
299    pub tags_dropped: u64,
300    /// Fields past [`MAX_FIELDS`] in a message.
301    pub fields_dropped: u64,
302    /// Messages cut at [`MAX_MESSAGE`].
303    pub cut: u64,
304    /// BeginStrings and how many messages each.
305    pub versions: Vec<(String, u64)>,
306    /// Per dictionary, the messages it applied to.
307    pub applied: Vec<u64>,
308}
309
310/// One tag's columns and what its values were.
311#[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    /// The names the dictionaries gave it, each with the dictionary.
318    names: Vec<(String, usize)>,
319    ty: Option<FixType>,
320    /// The dictionaries gave it two types.
321    types_differ: bool,
322    /// Every value read as `ty`.
323    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    /// The type its column is cast to, if any: a tag with enums keeps its names.
341    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/// The dictionaries one counterparty's messages are read with, and what they said.
351#[derive(Debug, Default)]
352struct View {
353    applying: Vec<usize>,
354    resolved: HashMap<u32, Arc<Resolved>>,
355}
356
357/// Reads a FIX log a piece at a time.
358#[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    /// Read `bytes`, the next piece of the file.
415    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    /// A full batch, once one is held: [`BATCH_ROWS`] rows, or [`BATCH_TEXT`] of text.
425    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    /// The end of the file: what is left, and the rows not yet taken.
433    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            // The next message, and the whole lines before it that hold none. Found
444            // before any newline is looked for, so a capture of messages back to back,
445            // with no newlines at all, is read in one pass.
446            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                    // A line with no `8=FIX` in its first MiB: passed over, keeping
458                    // what could be the start of one.
459                    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        // Each tag's first value, and the rest when it repeats.
553        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                // A batch of nulls only is a list of nulls; every batch's type must match.
692                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    /// Each tag and the names the dictionaries gave it, each with the dictionary's index
718    /// in [`Self::layers`].
719    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    /// Each tag's column name: the name every dictionary that applied gave it, or its
727    /// number when none did, when they disagree, or when the name is taken.
728    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            // A name of digits or with `#` could meet another column's raw name.
737            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    /// `batch`, the last, which has every column, as the dataset shows it: for checks.
756    pub fn finished(&self, batch: DataFrame) -> PolarsResult<DataFrame> {
757        self.finish_frame(batch.lazy()).collect()
758    }
759
760    /// The segments' columns renamed, typed and, for those that never held a value,
761    /// dropped.
762    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
812/// A string column as `ty`.
813fn 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
835/// The element type a repeated tag's list is cast to: numbers only.
836fn 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
925/// The FIX tab of the Info panel: versions, dictionaries and each tag's number.
926pub 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
988/// The dictionaries a log is read with: the search path's, then each `--dict`.
989pub 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                // A TOML dictionary's error may be in the file it names.
1004                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
1012/// Read a FIX log, given a piece at a time by `pieces`, into segments written through
1013/// `writer`.
1014pub(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    /// A log with no messages names itself; a dictionary that is not one names the
1053    /// dictionary and the flag that took it.
1054    #[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    /// A message of `body` (fields after 9, before 10) with a correct BodyLength and
1087    /// CheckSum, joined by `delim`.
1088    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    /// Empty cells count toward a batch: one message of 2,000 tags makes every later
1124    /// row 2,000 cells wide, so the batch is handed over long before `BATCH_ROWS`.
1125    #[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}