diff --git a/CHANGELOG.md b/CHANGELOG.md index b6e8f8d1..7a8a5bf5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/Project.toml b/Project.toml index 3d60e3c6..7d2aa985 100644 --- a/Project.toml +++ b/Project.toml @@ -1,6 +1,6 @@ name = "DatabentoBinaryEncoding" uuid = "90689371-c8cb-40d1-831f-18033db90f74" -version = "0.1.6" +version = "0.1.7" authors = ["Tyler Beason "] [deps] diff --git a/src/export.jl b/src/export.jl index a135f5b6..5e1f82ce 100644 --- a/src/export.jl +++ b/src/export.jl @@ -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 @@ -542,4 +554,4 @@ Remove null bytes from a string. """ function strip_nulls(s::String) return replace(s, '\0' => "") -end \ No newline at end of file +end diff --git a/src/replay.jl b/src/replay.jl index 4bac762c..b05b5518 100644 --- a/src/replay.jl +++ b/src/replay.jl @@ -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, @@ -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 @@ -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`). @@ -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)) @@ -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)), @@ -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) diff --git a/test/test_phase10_complete.jl b/test/test_phase10_complete.jl index 72c3338c..e942b28e 100644 --- a/test/test_phase10_complete.jl +++ b/test/test_phase10_complete.jl @@ -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) @@ -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) @@ -260,4 +271,4 @@ using DataFrames # for nrow / ncol on dbn_to_csv / dbn_to_parquet / records_to end end end -end \ No newline at end of file +end diff --git a/test/test_phase9_working.jl b/test/test_phase9_working.jl index 4f7cd446..58079bc5 100644 --- a/test/test_phase9_working.jl +++ b/test/test_phase9_working.jl @@ -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 @@ -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 \ No newline at end of file +end diff --git a/test/test_replay.jl b/test/test_replay.jl index eadba7e2..312272aa 100644 --- a/test/test_replay.jl +++ b/test/test_replay.jl @@ -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