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
10 changes: 9 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,18 @@ 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]
## [0.1.7] - 2026-09-10

### Changed

- Replay pacing now uses a monotonic clock by default, so NTP or manual
wall-clock adjustments cannot change replay timing.

### Fixed

- Parquet compression names are validated before being interpolated into the
DuckDB `COPY` statement; unsupported codecs now throw `ArgumentError`.

- **`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
Expand Down
2 changes: 1 addition & 1 deletion Project.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
name = "DatabentoBinaryEncoding"
uuid = "90689371-c8cb-40d1-831f-18033db90f74"
version = "0.1.6"
version = "0.1.7"
authors = ["Tyler Beason <tbeas12@gmail.com>"]

[deps]
Expand Down
22 changes: 17 additions & 5 deletions src/export.jl
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,22 @@ using DataFrames

# Escape a path for a DuckDB SQL string literal: double single quotes, and on
# Windows use forward slashes (accepted by DuckDB, avoids backslash ambiguity).
_duckdb_sql_path(p::AbstractString) = replace(replace(p, '\\' => '/'), "'" => "''")

function _write_parquet(df::DataFrame, output_file::AbstractString; compression="zstd")
comp = uppercase(String(compression)) # ZSTD | SNAPPY | GZIP | UNCOMPRESSED
_duckdb_sql_path(p::AbstractString) = replace(replace(p, '\\' => '/'), "'" => "''")

const _PARQUET_COMPRESSION_CODECS = ("ZSTD", "SNAPPY", "GZIP", "UNCOMPRESSED")

function _parquet_compression(compression)
comp = uppercase(String(compression))
comp in _PARQUET_COMPRESSION_CODECS || throw(ArgumentError(
"unsupported Parquet compression $(repr(compression)); expected one of " *
join(lowercase.(_PARQUET_COMPRESSION_CODECS), ", ")))
return comp
end

function _write_parquet(df::DataFrame, output_file::AbstractString; compression="zstd")
# Validate before interpolating the codec into DuckDB SQL. Unlike the path,
# this token cannot be parameter-bound or quoted in COPY options.
comp = _parquet_compression(compression)
# DuckDB cannot COPY a zero-column frame, which is what `records_to_dataframe`
# returns for an empty record set (e.g. a schema with no rows in the window).
# Fall back to a minimal header-only schema so we still write a valid, empty
Expand Down Expand Up @@ -542,4 +554,4 @@ Remove null bytes from a string.
"""
function strip_nulls(s::String)
return replace(s, '\0' => "")
end
end
13 changes: 8 additions & 5 deletions src/replay.jl
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,8 @@ means time spent inside `f` is absorbed instead of accumulating as drift — a
slow callback simply shortens the next wait rather than pushing every
subsequent record later.
"""
_monotonic_seconds() = time_ns() / 1.0e9

function _replay_loop(f::Function, produce::Function;
speed::Real,
timestamp::Symbol,
Expand Down Expand Up @@ -160,7 +162,7 @@ end
timestamp::Symbol = :ts_event,
max_sleep::Union{Real,Nothing} = nothing,
precise::Bool = false,
clock::Function = time,
clock::Function = _monotonic_seconds,
sleep_fn::Union{Function,Nothing} = nothing) -> Int

Replay a DBN file, invoking `f(record)` for each record paced in real time
Expand All @@ -183,7 +185,8 @@ like [`DBNStream`](@ref). Records are delivered in file order.
- `precise`: when `true`, busy-wait sub-millisecond gaps instead of using
`Base.sleep`. See the timing-resolution note below.
- `clock` / `sleep_fn`: injection points for the wall clock and sleep function,
primarily for testing. `clock` defaults to `Base.time`. An explicit `sleep_fn`
primarily for testing. `clock` defaults to a monotonic clock derived from
`Base.time_ns`, so NTP or manual wall-clock adjustments cannot distort pacing. An explicit `sleep_fn`
overrides `precise`; when left as `nothing` the sleeper is `Base.sleep`
(or the precise busy-wait sleeper when `precise = true`).

Expand Down Expand Up @@ -224,7 +227,7 @@ function replay_dbn(f::Function, filename::AbstractString;
timestamp::Symbol = :ts_event,
max_sleep::Union{Real,Nothing} = nothing,
precise::Bool = false,
clock::Function = time,
clock::Function = _monotonic_seconds,
sleep_fn::Union{Function,Nothing} = nothing)
sleeper = _resolve_sleep_fn(sleep_fn, precise)
decoder = DBNDecoder(String(filename))
Expand Down Expand Up @@ -256,7 +259,7 @@ end
timestamp::Symbol = :ts_event,
max_sleep::Union{Real,Nothing} = nothing,
precise::Bool = false,
clock::Function = time,
clock::Function = _monotonic_seconds,
sleep_fn::Union{Function,Nothing} = nothing) -> Int

Replay an in-memory collection of records (e.g. the result of [`read_dbn`](@ref)),
Expand All @@ -280,7 +283,7 @@ function replay_records(f::Function, records;
timestamp::Symbol = :ts_event,
max_sleep::Union{Real,Nothing} = nothing,
precise::Bool = false,
clock::Function = time,
clock::Function = _monotonic_seconds,
sleep_fn::Union{Function,Nothing} = nothing)
sleeper = _resolve_sleep_fn(sleep_fn, precise)
next = iterate(records)
Expand Down
17 changes: 14 additions & 3 deletions test/test_phase10_complete.jl
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,7 @@ using DataFrames # for nrow / ncol on dbn_to_csv / dbn_to_parquet / records_to
end
end

@testset "Parquet Export" begin
@testset "Parquet Export" begin
temp_parquet = tempname() * ".parquet"
try
df = dbn_to_parquet(test_file, temp_parquet)
Expand All @@ -249,7 +249,18 @@ using DataFrames # for nrow / ncol on dbn_to_csv / dbn_to_parquet / records_to
finally
safe_rm(temp_parquet)
end
end
end

@testset "Parquet compression validation" begin
temp_parquet = tempname() * ".parquet"
try
@test_throws ArgumentError dbn_to_parquet(
test_file, temp_parquet; compression = "zstd); DROP TABLE x; --")
@test !isfile(temp_parquet)
finally
safe_rm(temp_parquet)
end
end

@testset "DataFrame Conversion" begin
metadata, records = read_dbn_with_metadata(test_file)
Expand All @@ -260,4 +271,4 @@ using DataFrames # for nrow / ncol on dbn_to_csv / dbn_to_parquet / records_to
end
end
end
end
end
33 changes: 19 additions & 14 deletions test/test_phase9_working.jl
Original file line number Diff line number Diff line change
Expand Up @@ -414,14 +414,18 @@ using Dates
end

@testset "Write Permission Errors" begin
@testset "Read-only directory" begin
# This test might not work in all environments
# Try to write to a system directory
readonly_paths = ["/", "/etc", "/usr"]

for path in readonly_paths
if isdir(path) && !Sys.iswindows() # Skip on Windows
readonly_file = joinpath(path, "test_dbn_readonly.dbn")
@testset "Read-only directory" begin
# Root can write to these paths despite their permissions, so the
# assertion is not meaningful in root-run containers.
if Sys.iswindows() || Base.Libc.geteuid() == 0
@test_skip false
else
# Try to write to a system directory
readonly_paths = ["/", "/etc", "/usr"]

for path in readonly_paths
if isdir(path)
readonly_file = joinpath(path, "test_dbn_readonly.dbn")

metadata = Metadata(
UInt8(3), # version
Expand All @@ -439,11 +443,12 @@ using Dates
Tuple{String,String,Int64,Int64}[] # mappings
)

# Should throw an error when trying to write
@test_throws Exception write_dbn(readonly_file, metadata, TradeMsg[])
break # Only need one successful test
end
end
# Should throw an error when trying to write
@test_throws Exception write_dbn(readonly_file, metadata, TradeMsg[])
break # Only need one successful test
end
end
end
end
end
end
end
8 changes: 8 additions & 0 deletions test/test_replay.jl
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,14 @@ import DatabentoBinaryEncoding as DBN
@test DBN._resolve_sleep_fn(custom, true) === custom
end

@testset "default pacing clock is monotonic" begin
t1 = DBN._monotonic_seconds()
yield()
t2 = DBN._monotonic_seconds()
@test isfinite(t1)
@test t2 >= t1
end

@testset "_precise_sleep does not return early" begin
# Zero / negative is a no-op.
@test DBN._precise_sleep(0) === nothing
Expand Down
Loading