Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 33 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,39 @@ All notable changes to DatabentoBinaryEncoding.jl are documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [Unreleased]

### Fixed

- **`MBOMsg` wire layout.** `read_mbo_msg` and `write_record(::MBOMsg)` used a
field order that does not match the DBN spec: `ts_recv` was read from bytes
16-23, `order_id` from 24-31 and `price` from 40-47, whereas the official
`MboMsg` (identical in DBN v1/v2/v3) has `order_id` at 16, `price` at 24 and
`ts_recv` at 40. Because the encoder mirrored the decoder, files written by
this package round-tripped and every self-consistency test passed, but every
MBO file produced by Databento (historical downloads, live captures, the
official `test_data.mbo.*` fixtures) decoded with the three fields rotated:
`price` held the `ts_recv` nanoseconds, `ts_recv` held the `order_id`, and
`order_id` held the fixed-point price. Both paths now use the official layout.
New `test/test_mbo_wire_layout.jl` decodes Databento's fixtures against
reference values from the official decoder and checks a byte-exact re-encode.
**Migration:** MBO `.dbn` files that earlier versions *wrote from Julia-built
records* (`write_dbn`, `DBNStreamWriter`, `csv_to_dbn` / `json_to_dbn` /
`parquet_to_dbn`) carry the swapped layout on disk and will now decode
rotated; re-export them from the source data. Files captured by decoding and
re-encoding gateway bytes (e.g. DatabentoAPI.jl `stream_to_file`) are
byte-identical to the wire and decode correctly with this fix.
- **`StatMsg` undefined-quantity sentinel on write.** The encoder wrote an
undefined v3 `quantity` (`typemax(Int64)`) as `0xffffffffffffffff`, i.e. `-1`,
instead of the spec sentinel `typemax(Int64)` (`0x7fff…`). Other readers
(databento, duckdb-dbn) therefore showed `-1` where they should show
NULL/NaN, and a Julia round trip of an undefined quantity came back as `-1`.
The quantity is now written verbatim as a signed `Int64`, and the v3 decoder
reads it as a signed `Int64` (it used to read a `UInt64` and treat every
value >= `0x7fff…`, i.e. any negative quantity, as undefined, which masked
the bug). Files written by
earlier versions with undefined statistics quantities contain `-1` on disk.

## [0.1.6] - 2026-06-24

### Changed
Expand Down
20 changes: 13 additions & 7 deletions src/decode.jl
Original file line number Diff line number Diff line change
Expand Up @@ -502,15 +502,19 @@ end
# We've already read: length(1) + rtype(1) + publisher_id(2) + instrument_id(4) + ts_event(8) = 16 bytes
# Remaining to read: 56 - 16 = 40 bytes

# Based on Rust struct order and empirical evidence:
ts_recv = read(decoder.io, Int64) # 8 bytes (positions 16-23)
order_id = read(decoder.io, UInt64) # 8 bytes (positions 24-31)
# Official DBN MboMsg layout (dbn crate `MboMsg`; identical in v1/v2/v3):
# order_id u64 @16, price i64 @24, size u32 @32, flags u8 @36, channel_id u8 @37,
# action c_char @38, side c_char @39, ts_recv u64 @40, ts_in_delta i32 @48, sequence u32 @52.
# Struct field order == wire order, so no reordering is needed. (Through 0.1.6 this
# reader swapped order_id/ts_recv and price/ts_recv; see CHANGELOG.)
order_id = read(decoder.io, UInt64) # 8 bytes (positions 16-23)
price = read(decoder.io, Int64) # 8 bytes (positions 24-31)
size = read(decoder.io, UInt32) # 4 bytes (positions 32-35)
flags = read(decoder.io, UInt8) # 1 byte (position 36)
channel_id = read(decoder.io, UInt8) # 1 byte (position 37)
action = safe_action(read(decoder.io, UInt8)) # 1 byte (position 38)
side = safe_side(read(decoder.io, UInt8)) # 1 byte (position 39)
price = read(decoder.io, Int64) # 8 bytes (positions 40-47)
ts_recv = read(decoder.io, Int64) # 8 bytes (positions 40-47)
ts_in_delta = read(decoder.io, Int32) # 4 bytes (positions 48-51)
sequence = read(decoder.io, UInt32) # 4 bytes (positions 52-55)

Expand Down Expand Up @@ -642,9 +646,11 @@ end
quantity_raw = read(decoder.io, Int32)
quantity_raw == typemax(Int32) ? typemax(Int64) : Int64(quantity_raw)
else
# v3: 64-bit quantity, UNDEF_STAT_QUANTITY = typemax(Int64)
quantity_raw = read(decoder.io, UInt64)
quantity_raw >= 0x7fffffffffffffff ? typemax(Int64) : Int64(quantity_raw)
# v3: signed 64-bit quantity; UNDEF_STAT_QUANTITY = typemax(Int64) (0x7fff...), which
# needs no mapping. (Through 0.1.6 this read a UInt64 and treated every value
# >= 0x7fff... - i.e. any negative quantity - as UNDEF, masking that the encoder
# wrote -1 for UNDEF.)
read(decoder.io, Int64)
end
sequence = read(decoder.io, UInt32)
ts_in_delta = read(decoder.io, Int32)
Expand Down
20 changes: 10 additions & 10 deletions src/encode.jl
Original file line number Diff line number Diff line change
Expand Up @@ -280,17 +280,17 @@ Each record type is serialized according to its specific binary layout.
unsafe_write(encoder.io, Ref(record), sizeof(record))
end

# Specialized optimized write for MBOMsg with field reordering
# Specialized optimized write for MBOMsg
@inline function write_record(encoder::DBNEncoder, record::MBOMsg)
# MBOMsg requires custom serialization: struct field order != binary order
# Binary order: hd → ts_recv → order_id → size → flags → channel_id → action → side → price → ts_in_delta → sequence
# Struct order: hd → order_id → price → size → flags → channel_id → action → side → ts_recv → ts_in_delta → sequence
# Wire order == struct order (official DBN MboMsg): hd -> order_id -> price -> size -> flags
# -> channel_id -> action -> side -> ts_recv -> ts_in_delta -> sequence. (Through 0.1.6 this
# writer mirrored the decoder's swapped order_id/ts_recv/price layout; see CHANGELOG.)
#
# Performance: This IOBuffer approach achieves 1.4x speedup (40% faster) compared to field-by-field write()
# by batching all fields into a buffer and performing a single write operation (2.1M vs 1.5M records/sec).
# The reduction in IO syscalls more than compensates for the temporary 56-byte allocation per record.

# Use IOBuffer to reorder fields, then write in one operation
# Batch all fields into a buffer, then write in one operation
buffer = IOBuffer()

# Write header (16 bytes)
Expand All @@ -301,14 +301,14 @@ end
write(buffer, record.hd.ts_event)

# Write body in binary order (40 bytes)
write(buffer, record.ts_recv)
write(buffer, record.order_id)
write(buffer, record.price)
write(buffer, record.size)
write(buffer, record.flags)
write(buffer, record.channel_id)
write(buffer, UInt8(record.action))
write(buffer, UInt8(record.side))
write(buffer, record.price)
write(buffer, record.ts_recv)
write(buffer, record.ts_in_delta)
write(buffer, record.sequence)

Expand Down Expand Up @@ -623,9 +623,9 @@ function write_record_complex(encoder::DBNEncoder, record)
write(io, record.ts_recv)
write(io, record.ts_ref)
write(io, record.price)
# Write quantity as UInt64, converting back if needed
quantity_to_write = record.quantity == typemax(Int64) ? 0xffffffffffffffff : UInt64(record.quantity)
write(io, quantity_to_write)
# quantity is a signed Int64 on the wire; the v3 UNDEF sentinel is typemax(Int64)
# (0x7fff...). Through 0.1.6 this wrote 0xffffffffffffffff (-1) for UNDEF instead.
write(io, record.quantity)
write(io, record.sequence)
write(io, record.ts_in_delta)
write(io, record.stat_type)
Expand Down
2 changes: 2 additions & 0 deletions test/runtests.jl
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ include("test_utils.jl")
include("test_show.jl") # compact one-line Base.show for record types
include("test_symbols.jl") # symbol_map / symbol_for / add_symbol_column! / records_to_dataframe(records, metadata)
include("test_issue40_mbp_dataframe.jl") # regression: MBP-1/MBP-10/BBO records_to_dataframe via nested levels (issue #40)
include("test_mbo_wire_layout.jl") # regression: MBOMsg wire layout (order_id/price/ts_recv offsets) vs Databento's own fixtures
include("test_stat_quantity_sentinel.jl") # regression: v3 StatMsg UNDEF quantity written as typemax(Int64), not -1

# Run compatibility tests if the Rust CLI is available
dbn_cli_path = if Sys.iswindows()
Expand Down
123 changes: 123 additions & 0 deletions test/test_mbo_wire_layout.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
# Regression: MBOMsg wire layout must match the official DBN `MboMsg` struct.
#
# Before 0.1.7, read_mbo_msg / write_record(::MBOMsg) used a swapped layout
# (ts_recv read from bytes 16-23, order_id from 24-31, price from 40-47). The
# encoder mirrored the decoder, so files written by this package round-tripped
# and every self-consistency test passed — but any file produced by Databento
# (historical downloads, live captures written by other tools, the official
# test fixtures) decoded with order_id / price / ts_recv rotated:
# order_id <- price bytes, price <- ts_recv bytes, ts_recv <- order_id bytes.
# Official layout (identical for DBN v1/v2/v3, 56 bytes):
# hd(16) order_id u64 @16, price i64 @24, size u32 @32, flags u8 @36,
# channel_id u8 @37, action @38, side @39, ts_recv u64 @40, ts_in_delta i32 @48,
# sequence u32 @52.

@testset "MBOMsg wire layout matches the official DBN struct" begin
data_dir = joinpath(@__DIR__, "data")

# Databento's own MBO fixture (same two records in every DBN version).
# Reference values from the official decoder:
# order_id 647784973705, price 3722.75, ts_recv 2020-12-28T13:00:00.000704060Z
expected = (
order_id = UInt64(647784973705),
price = Int64(3_722_750_000_000),
size = UInt32(1),
flags = 0x80,
channel_id = 0x00,
action = Action.CANCEL,
side = Side.ASK,
ts_event = Int64(1609160400000429831),
ts_recv = Int64(1609160400000704060),
ts_in_delta = Int32(22993),
sequence = UInt32(1170352),
)

fixtures = filter(f -> occursin(r"^test_data\.mbo(\.v[123])?\.dbn(\.zst)?$", f), readdir(data_dir))
@test !isempty(fixtures)

for f in fixtures
@testset "$f" begin
recs = read_dbn(joinpath(data_dir, f))
@test length(recs) == 2
r = recs[1]
@test r isa MBOMsg
@test r.hd.instrument_id == 5482
@test r.hd.ts_event == expected.ts_event
@test r.order_id == expected.order_id
@test r.price == expected.price
@test price_to_float(r.price) == 3722.75
@test r.size == expected.size
@test r.flags == expected.flags
@test r.channel_id == expected.channel_id
@test r.action == expected.action
@test r.side == expected.side
@test r.ts_recv == expected.ts_recv
@test r.ts_in_delta == expected.ts_in_delta
@test r.sequence == expected.sequence
# Sanity invariants that the swapped layout violated: ts_recv is a
# nanosecond timestamp at or after ts_event; price is a sane fixed-point value.
for x in recs
@test x.ts_recv >= x.hd.ts_event
@test 0 < price_to_float(x.price) < 1_000_000
end
end
end

# Byte-exact re-encode of the official (uncompressed) fixtures: decoding then
# re-encoding must reproduce the original record bytes, which pins the writer
# to the wire layout independently of the reader.
for f in filter(f -> endswith(f, ".dbn"), fixtures)
@testset "re-encode $f" begin
path = joinpath(data_dir, f)
bytes = read(path)
meta, recs = read_dbn_with_metadata(path)
io = IOBuffer()
enc = DBNEncoder(io, meta)
for r in recs
write_record(enc, r)
end
out = take!(io)
@test length(out) == 2 * sizeof(MBOMsg)
@test bytes[end-length(out)+1:end] == out
end
end

# Hand-encoded record at the official offsets decodes to the right fields.
@testset "hand-encoded official layout" begin
ts0 = Int64(1_700_000_000_000_000_000)
meta = Metadata(UInt8(3), "GLBX.MDP3", Schema.MBO, ts0, ts0 + 1, nothing,
SType.RAW_SYMBOL, SType.INSTRUMENT_ID, false,
String[], String[], String[], Tuple{String,String,Int64,Int64}[])
io = IOBuffer()
enc = DBNEncoder(io, meta)
write_header(enc)
write(io, UInt8(14)); write(io, UInt8(RType.MBO_MSG)); write(io, UInt16(1)); write(io, UInt32(42)); write(io, ts0)
write(io, UInt64(777)) # order_id @16
write(io, Int64(1_234_000_000_000)) # price @24 (1234.0)
write(io, UInt32(5)) # size @32
write(io, UInt8(0x80)) # flags @36
write(io, UInt8(3)) # channel @37
write(io, UInt8('A')) # action @38
write(io, UInt8('B')) # side @39
write(io, Int64(ts0 + 99)) # ts_recv @40
write(io, Int32(-7)) # ts_in_delta @48
write(io, UInt32(9)) # sequence @52
tmp = tempname() * ".dbn"
try
write(tmp, take!(io))
r = only(read_dbn(tmp))
@test r.order_id == 777
@test r.price == 1_234_000_000_000
@test r.size == 5
@test r.flags == 0x80
@test r.channel_id == 3
@test r.action == Action.ADD
@test r.side == Side.BID
@test r.ts_recv == ts0 + 99
@test r.ts_in_delta == -7
@test r.sequence == 9
finally
safe_rm(tmp)
end
end
end
11 changes: 8 additions & 3 deletions test/test_phase5.jl
Original file line number Diff line number Diff line change
Expand Up @@ -650,7 +650,9 @@ using Dates

# First record
r1 = records[1]
@test r1.order_id == 3722750000000
@test r1.order_id == 647784973705 # official value; pre-0.1.7 this held the price bytes
@test r1.price == 3722750000000
@test r1.ts_recv == 1609160400000704060
@test r1.action == Action.CANCEL
@test r1.side == Side.ASK
@test r1.size == 1
Expand All @@ -659,7 +661,8 @@ using Dates

# Second record
r2 = records[2]
@test r2.order_id == 3723000000000
@test r2.order_id == 647784973631
@test r2.price == 3723000000000
@test r2.action == Action.CANCEL
@test r2.side == Side.ASK
@test r2.sequence == r1.sequence + 1
Expand All @@ -675,7 +678,9 @@ using Dates

# Should match uncompressed data
r1 = records[1]
@test r1.order_id == 3722750000000
@test r1.order_id == 647784973705 # official value; pre-0.1.7 this held the price bytes
@test r1.price == 3722750000000
@test r1.ts_recv == 1609160400000704060
@test r1.action == Action.CANCEL
@test r1.side == Side.ASK
end
Expand Down
43 changes: 43 additions & 0 deletions test/test_stat_quantity_sentinel.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
# Regression: the v3 StatMsg "undefined quantity" sentinel on the wire is
# typemax(Int64) (0x7fff_ffff_ffff_ffff, i64::MAX in the dbn crate). Through
# 0.1.6 the encoder wrote 0xffff_ffff_ffff_ffff (-1) for an undefined quantity,
# so other readers (databento, duckdb-dbn) showed -1 instead of NULL/NaN and a
# Julia round trip of an undefined quantity came back as -1.

@testset "StatMsg undefined quantity is written as typemax(Int64)" begin
ts0 = Int64(1_700_000_000_000_000_000)
meta = Metadata(UInt8(3), "GLBX.MDP3", Schema.STATISTICS, ts0, ts0 + 1, nothing,
SType.RAW_SYMBOL, SType.INSTRUMENT_ID, false,
String[], String[], String[], Tuple{String,String,Int64,Int64}[])
hd = RecordHeader(UInt8(20), RType.STAT_MSG, UInt16(1), UInt32(7), ts0)
mk(qty, seq) = StatMsg(hd, UInt64(ts0 + 1), UInt64(ts0), Int64(5_100_000_000_000), qty,
UInt32(seq), Int32(0), UInt16(17), UInt16(0), UInt8(1), UInt8(0))
undef_q = mk(typemax(Int64), 1)
neg_one = mk(Int64(-1), 2) # a real -1 must stay -1 (not be confused with UNDEF)
normal = mk(Int64(12345), 3)

# Raw bytes: quantity is the i64 at offset 40 of the 80-byte v3 record.
io = IOBuffer()
enc = DBNEncoder(io, meta)
write_record(enc, undef_q)
bytes = take!(io)
@test length(bytes) == 80
@test reinterpret(Int64, bytes[41:48])[1] == typemax(Int64)
@test bytes[41:48] == UInt8[0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0x7f]

# Round trip preserves the sentinel, a genuine -1, and a normal value.
tmp = tempname() * ".dbn"
try
write_dbn(tmp, meta, [undef_q, neg_one, normal])
recs = read_dbn(tmp)
@test length(recs) == 3
@test recs[1].quantity == typemax(Int64)
@test recs[2].quantity == -1
@test recs[3].quantity == 12345
# DataFrame export: the sentinel is the only one that is not a real quantity
df = records_to_dataframe(StatMsg[r for r in recs])
@test df.quantity[2] == -1 && df.quantity[3] == 12345
finally
safe_rm(tmp)
end
end
Loading