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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,14 @@

Notable changes are recorded from 0.5.0 onward.

## v0.7.1 1 Oct 2026

### Maintenance
- Compatibility with the 2026-09-30 Rust nightly, where `Vec::drain` on custom allocators requires the new `allocator_ext` feature (due to Rust planning to stabilise a lot of the allocator trait in the upcoming stable release).
Therefore, `Vec64` stream buffers now use `Vec64::delete_range` to avoid the feature gate.
- vec64 upgraded to 0.5.3.
- Minarrow upgraded to 0.18.3. The Python package pins minarrow and minarrow-pyo3 at 0.18.3.

## v0.7.0 19 Sep 2026

### Enhancements
Expand Down
16 changes: 8 additions & 8 deletions python/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 3 additions & 3 deletions python/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ cargo-features = ["trim-paths"]

[package]
name = "lightstream-py"
version = "0.7.0"
version = "0.7.1"
edition = "2024"
authors = ["Peter G. Bower"]
license = "MPL-2.0"
Expand All @@ -24,8 +24,8 @@ lightstream = { version = "0.7", path = "../rust", features = ["csv", "datetime"
# minarrow-pyo3's dictionary-index conversion needs the extended features,
# and they flow through lightstream's flags so its match arms gate in step
# with minarrow's variants.
minarrow = { version = "0.18.2", features = ["chunked"] }
minarrow-pyo3 = { version = "0.18.2", features = ["extended_categorical", "extended_numeric_types"] }
minarrow = { version = "0.18.3", features = ["chunked"] }
minarrow-pyo3 = { version = "0.18.3", features = ["extended_categorical", "extended_numeric_types"] }
futures-core = "0.3"
pyo3 = { version = "0.29", features = ["abi3-py39"] }
# The QUIC and WebTransport dependency pins mirror lightstream's, so
Expand Down
2 changes: 1 addition & 1 deletion python/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ build-backend = "maturin"

[project]
name = "lightstream-io"
version = "0.7.0"
version = "0.7.1"
description = "Streaming Arrow I/O for Python - files, sockets, and network transports with zero-copy minarrow interop."
readme = "README.md"
requires-python = ">=3.9"
Expand Down
2 changes: 1 addition & 1 deletion python/rust-toolchain.toml
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# lightstream (through minarrow and vec64) uses allocator_api and portable_simd,
# lightstream (through minarrow and vec64) uses allocator_ext and portable_simd,
# which are only available on the nightly toolchain. cargo selects it
# automatically here, so local builds, CI, and source installs all use the
# right channel.
Expand Down
10 changes: 6 additions & 4 deletions rust/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions rust/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "lightstream"
version = "0.7.0"
version = "0.7.1"
edition = "2024"
license = "MPL-2.0"
keywords = [
Expand Down Expand Up @@ -41,7 +41,7 @@ flatbuffers = "25.12.19"
libc = { version = "0.2.183", optional = true }
futures-sink = "0.3.32"
log = "0.4"
minarrow = { version = "0.18.2", features = ["chunked", "views", "select", "size"], default-features = false }
minarrow = { version = "0.18.3", features = ["chunked", "views", "select", "size"], default-features = false }
vec64 = { version = "0.5.2" }
snap = { version = "1.0", optional = true }
tokio-tungstenite = { version = "0.28", optional = true }
Expand Down
2 changes: 1 addition & 1 deletion rust/rust-toolchain.toml
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
# Lightstream builds on nightly: the minarrow and vec64 dependencies use
# `allocator_api` and `portable_simd`, both nightly-only features.
# `allocator_ext` and `portable_simd`, both nightly-only features.
#
# Note that Nightly is a temporary requirement and moving Lightstream off Nightly is a current priority
# that will be addressed in an upcoming release.
Expand Down
2 changes: 1 addition & 1 deletion rust/src/compression/ipc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
//! body is prefixed with an 8-byte i64 uncompressed length (-1 means the
//! buffer was stored uncompressed), followed by the compressed or raw data.
//!
//! [`decompress_ipc_body`](crate::compression::ipc::decompress_ipc_body) produces a new `Vec64<u8>` with all buffers placed
//! [`decompress_ipc_body`] produces a new `Vec64<u8>` with all buffers placed
//! at `B::ALIGN` offsets, consistent with the uncompressed wire layout. A
//! corrections Vec maps each buffer index to its (offset, length) within the
//! decompressed buffer, since the flatbuffer metadata still references the
Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/codecs/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@

//! Codecs for Arrow IPC and the Lightstream protocol.
//!
//! - [`ArrowIpcCodec`](crate::models::codecs::ipc::ArrowIpcCodec) - Arrow IPC streaming codec with zero-copy encode/decode
//! - [`LightstreamCodec`](crate::models::codecs::lightstream::LightstreamCodec) - Lightstream protocol codec with type registry and TLV framing
//! - [`ArrowIpcCodec`] - Arrow IPC streaming codec with zero-copy encode/decode
//! - [`LightstreamCodec`] - Lightstream protocol codec with type registry and TLV framing

/// Arrow IPC streaming codec with zero-copy encode and decode.
pub mod ipc;
Expand Down
6 changes: 3 additions & 3 deletions rust/src/models/decoders/csv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,11 @@

//! # CSV Decoder for Minarrow Tables
//!
//! - Accepts a CSV byte slice or any [`BufRead`](std::io::BufRead).
//! - Accepts a CSV byte slice or any [`BufRead`].
//! - Infers schema or uses a provided schema (optional).
//! - Supports: `Int32`, `Int64`, `UInt32`, `UInt64`, `Float32`, `Float64`, `Boolean`, `String32`, `Categorical32`, `Categorical8`.
//! - Custom delimiter, nulls, quoting, and dictionary mapping for categoricals.
//! - Produces a single [`Table`](minarrow::Table) via [`decode_csv`](crate::models::decoders::csv::decode_csv), or multiple batches via repeated calls to [`decode_csv_batch`](crate::models::decoders::csv::decode_csv_batch).
//! - Produces a single [`Table`] via [`decode_csv`], or multiple batches via repeated calls to [`decode_csv_batch`].
//!
//! ## Fast path
//!
Expand All @@ -23,7 +23,7 @@
//!
//! Notes:
//! - Input is treated as UTF-8; invalid byte sequences are lossily decoded via `String::from_utf8_lossy`.
//! - See [`CsvDecodeOptions`](crate::models::decoders::csv::CsvDecodeOptions) for configurable delimiter, quoting, header handling, and schema control.
//! - See [`CsvDecodeOptions`] for configurable delimiter, quoting, header handling, and schema control.

use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
Expand Down
8 changes: 4 additions & 4 deletions rust/src/models/decoders/json/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,17 +10,17 @@
//! - **Array-of-objects**: `[{"col_a": 1, "col_b": "x"}, ...]`
//! - **NDJSON** (newline-delimited): `{"col_a": 1, "col_b": "x"}\n...`
//!
//! Backed by `simd-json` via the [`simd`](crate::models::decoders::json::simd) module. Cells dispatch
//! directly into pre-allocated [`builder::ColumnBuilder`](crate::models::decoders::json::builder::ColumnBuilder) buffers - no
//! Backed by `simd-json` via the [`simd`] module. Cells dispatch
//! directly into pre-allocated [`builder::ColumnBuilder`] buffers - no
//! intermediate `Value` tree or per-cell allocations except when copying
//! string bytes into the column's data buffer.
//!
//! ## Schema
//! Schema is **required** - pass it via [`JsonDecodeOptions::schema`](crate::models::decoders::json::JsonDecodeOptions::schema). Schema
//! Schema is **required** - pass it via [`JsonDecodeOptions::schema`]. Schema
//! inference from sampled rows is a planned follow-up.
//!
//! ## Type mismatch handling
//! See [`builder::TypeMismatchPolicy`](crate::models::decoders::json::builder::TypeMismatchPolicy).
//! See [`builder::TypeMismatchPolicy`].

pub mod row_decoder;
pub mod value;
Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/decoders/limits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,11 @@
//!
//! Limits are checked at the point where a length or count is read from untrusted
//! bytes, before any allocation that scales with it. Each decoder accepts a
//! [`DecodeLimits`](crate::models::decoders::limits::DecodeLimits) either through its constructor or, for free functions, as a
//! [`DecodeLimits`] either through its constructor or, for free functions, as a
//! trailing parameter threaded from its host reader/codec.
//!
//! Callers that need to opt out (replay tools, internal test fixtures, trusted
//! pipelines feeding the decoder from disk) construct via [`DecodeLimits::unlimited`](crate::models::decoders::limits::DecodeLimits::unlimited).
//! pipelines feeding the decoder from disk) construct via [`DecodeLimits::unlimited`].

use std::io;

Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/decoders/tlv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@

//! # TLV Frame Decoder
//!
//! Provides a [`FrameDecoder`](crate::traits::frame_decoder::FrameDecoder) implementation for Type-Length-Value (TLV) encoded frames.
//! Provides a [`FrameDecoder`] implementation for Type-Length-Value (TLV) encoded frames.
//!
//! ## Format
//! - **Type**: 1 byte (`u8`)
Expand All @@ -15,7 +15,7 @@
//!
//! Example frame layout: `[type][length][value...]`
//!
//! Produces [`TLVDecodedFrame`](crate::models::frames::tlv_frame::TLVDecodedFrame) instances for downstream consumers.
//! Produces [`TLVDecodedFrame`] instances for downstream consumers.

use std::convert::TryInto;
use std::io;
Expand Down
2 changes: 1 addition & 1 deletion rust/src/models/encoders/tlv/tlv_stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
//! Asynchronous, pull-based producer that encodes Type-Length-Value frames and yields them
//! as a `futures_core::Stream` of buffers. Uses the generic [`StreamBuffer`] so you can emit
//! either standard `Vec<u8>` or SIMD-aligned `Vec64<u8>` buffers. Push frames with
//! [`TLVStreamWriter::write_frame`](crate::models::encoders::tlv::tlv_stream::TLVStreamWriter::write_frame), then call [`TLVStreamWriter::finish`](crate::models::encoders::tlv::tlv_stream::TLVStreamWriter::finish) to signal end of stream.
//! [`TLVStreamWriter::write_frame`], then call [`TLVStreamWriter::finish`] to signal end of stream.
//!
//! No alignment padding is applied.

Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/frames/ipc_message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@
//!
//! Core data structures for Arrow IPC frame encoding.
//!
//! - [`ArrowIPCMessage`](crate::models::frames::ipc_message::ArrowIPCMessage) wraps a FlatBuffers metadata message and its associated body buffer.
//! - [`IPCFrameMetadata`](crate::models::frames::ipc_message::IPCFrameMetadata) tracks byte lengths and padding for all frame sections, used to compute
//! - [`ArrowIPCMessage`] wraps a FlatBuffers metadata message and its associated body buffer.
//! - [`IPCFrameMetadata`] tracks byte lengths and padding for all frame sections, used to compute
//! total frame size and to ensure compliance with Arrow IPC alignment rules.
//!
//! These are internal, low-level components used by the IPC encoders and writers.
Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/frames/json.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@
//!
//! A json frame can hold many records under one envelope, so the
//! [`JsonInterface`](crate::models::interfaces::json::JsonInterface) parses
//! it once and returns a [`JsonFrame`](crate::models::frames::json::JsonFrame) to read the records. Each
//! [`JsonRecord`](crate::models::frames::json::JsonRecord) gives its columns as [`JsonValueRef`](crate::models::decoders::json::value::JsonValueRef)s borrowed
//! it once and returns a [`JsonFrame`] to read the records. Each
//! [`JsonRecord`] gives its columns as [`JsonValueRef`]s borrowed
//! from the parsed buffer for the destination to handle and/or convert.

use std::io;
Expand Down
2 changes: 1 addition & 1 deletion rust/src/models/frames/lightstream_message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
//! [type_tag: u8][payload_len: u32 LE][payload: N bytes]
//! ```
//!
//! After decoding, frames are represented as [`LightstreamMessage`](crate::models::frames::lightstream_message::LightstreamMessage) variants
//! After decoding, frames are represented as [`LightstreamMessage`] variants
//! - either an opaque message or a decoded Arrow table.
//!
//! With the `protobuf` feature enabled, message variants gain typed decode
Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/frames/tlv_frame.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@
//! Type-Length-Value (TLV) frame definitions.
//!
//! Provides lightweight frame structs for TLV-based protocols:
//! - [`TLVFrame`](crate::models::frames::tlv_frame::TLVFrame) for encoding, which borrows value slices.
//! - [`TLVDecodedFrame`](crate::models::frames::tlv_frame::TLVDecodedFrame) for decoding, and owns the buffer via [`StreamBuffer`](crate::traits::stream_buffer::StreamBuffer)).
//! - [`TLVFrame`] for encoding, which borrows value slices.
//! - [`TLVDecodedFrame`] for decoding, and owns the buffer via [`StreamBuffer`]).

use crate::traits::stream_buffer::StreamBuffer;

Expand Down
6 changes: 3 additions & 3 deletions rust/src/models/interfaces/json/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,19 +17,19 @@
//! takes a wire shape that already matches the table - JSON keys as column names,
//! and one flat object per row. In contrast, this interface takes a connector's shape
//! (that may differ by vendor, etc.) and a [`JsonSchema`]'s per-column
//! [`ValueSource`](crate::models::interfaces::json::ValueSource) to map it onto the same columns.
//! [`ValueSource`] to map it onto the same columns.
//!
//! ## Behaviour
//!
//! Flow: The mapping compiles once at construction. Each frame is parsed through
//! simd-json's tape with reusable [`Buffers`](simd_json::Buffers). The record envelope is
//! simd-json's tape with reusable [`Buffers`]. The record envelope is
//! then resolved, and every record yields one value per column.
//!
//! Timing: A wall-clock column fills from the caller's receive `now`,
//! passed in epoch nanoseconds and scaled to the column's time unit when
//! the source compiles.
//!
//! Size Limits: [`DecodeLimits`](crate::models::decoders::limits::DecodeLimits) caps frame bytes, records per frame, and string value
//! Size Limits: [`DecodeLimits`] caps frame bytes, records per frame, and string value
//! length before the corresponding work happens. Downstream connector input is considered
//! untrusted.

Expand Down
4 changes: 2 additions & 2 deletions rust/src/models/interfaces/json/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ pub enum ReadAs {
}

/// A target schema that maps each upstream JSON record into typed
/// columns. Built column by column with [`column`](Self::column), then
/// columns. Built column by column with [`column`], then
/// executed by
/// [`JsonInterface`](crate::models::interfaces::json::JsonInterface).
#[derive(Debug, Clone, Default)]
Expand All @@ -63,7 +63,7 @@ pub struct JsonSchema {
}

impl JsonSchema {
/// An empty schema. Add columns with [`column`](Self::column).
/// An empty schema. Add columns with [`column`].
pub fn new() -> Self {
Self::default()
}
Expand Down
2 changes: 1 addition & 1 deletion rust/src/models/io_uring/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
//! no ring thread, channels, or cross-thread overhead. The io_uring
//! driver is integrated into the tokio event loop.
//!
//! Generic over any [`UringStream`](crate::models::io_uring::UringStream) implementor (UDS, TCP, etc.).
//! Generic over any [`UringStream`] implementor (UDS, TCP, etc.).
//! Monomorphised at compile time for no overhead dispatch.
//!
//! Requires the `io_uring` feature and Linux. Connections must be
Expand Down
8 changes: 4 additions & 4 deletions rust/src/models/protocol/ipc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,20 +12,20 @@
//!
//! ## Codec
//!
//! [`ArrowIpcCodec`](crate::models::codecs::ipc::ArrowIpcCodec) is the central codec for Arrow IPC encode and decode.
//! [`ArrowIpcCodec`] is the central codec for Arrow IPC encode and decode.
//! It owns the encoder state machine, decoded schema, dictionary registry,
//! and SharedBuffer cache for zero-copy buffer recycling.
//!
//! ## Readers
//!
//! [`TableReader`](crate::models::readers::ipc::table::TableReader) wraps the streaming decoder and reads Arrow IPC tables
//! [`TableReader`] wraps the streaming decoder and reads Arrow IPC tables
//! from any `AsyncRead` source. Transport-specific readers (TCP, UDS, etc.)
//! delegate to it internally.
//!
//! ## Writers
//!
//! - [`TableWriter`](crate::models::writers::ipc::table::TableWriter) - async writer to files or streams via the Sink trait
//! - [`TableStreamWriter`](crate::models::writers::ipc::table_stream::TableStreamWriter) - synchronous frame-by-frame writer for pipes
//! - [`TableWriter`] - async writer to files or streams via the Sink trait
//! - [`TableStreamWriter`] - synchronous frame-by-frame writer for pipes
//! and custom protocols
//!
//! ## Sinks
Expand Down
2 changes: 1 addition & 1 deletion rust/src/models/protocol/lightstream/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@

//! Bidirectional Lightstream protocol connection.
//!
//! Wraps a [`LightstreamReader`](crate::models::readers::lightstream::LightstreamReader) and [`LightstreamWriter`](crate::models::writers::lightstream::LightstreamWriter) together with
//! Wraps a [`LightstreamReader`] and [`LightstreamWriter`] together with
//! transport-specific constructors for TCP, UDS, WebSocket, QUIC, stdio,
//! and WebTransport.
//!
Expand Down
Loading
Loading