diff --git a/CHANGELOG.md b/CHANGELOG.md index 44ad409..af13a6b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/python/Cargo.lock b/python/Cargo.lock index c538bef..eb871fb 100644 --- a/python/Cargo.lock +++ b/python/Cargo.lock @@ -662,7 +662,7 @@ checksum = "68ab91017fe16c622486840e4c83c9a37afeff978bd239b5293d61ece587de66" [[package]] name = "lightstream" -version = "0.7.0" +version = "0.7.1" dependencies = [ "bytes", "fast-float2", @@ -692,7 +692,7 @@ dependencies = [ [[package]] name = "lightstream-py" -version = "0.7.0" +version = "0.7.1" dependencies = [ "futures-core", "lightstream", @@ -734,9 +734,9 @@ checksum = "f8ca58f447f06ed17d5fc4043ce1b10dd205e060fb3ce5b979b8ed8e59ff3f79" [[package]] name = "minarrow" -version = "0.18.1" +version = "0.18.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e53f5031b7d16aa76bc8671171814419694c687f5160917a3668f010dab2b609" +checksum = "4b272b619741fbb23ed37c0079a7b8dde4bd66239672b201752cef93fe13c05b" dependencies = [ "log", "num-traits", @@ -745,9 +745,9 @@ dependencies = [ [[package]] name = "minarrow-pyo3" -version = "0.18.1" +version = "0.18.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "94e91845f51b0aab761625144ab1c71b989d5237d8cc6159245db07b4b78bc26" +checksum = "25949301756ebc55abdcd6eaa987b9821404e77c8de0df3bfa1b445f33368650" dependencies = [ "minarrow", "pyo3", @@ -1673,9 +1673,9 @@ dependencies = [ [[package]] name = "vec64" -version = "0.5.1" +version = "0.5.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1aef6bbef159f21ac387220ebc71141524423f7886afd390bc7b140f12630425" +checksum = "662d4c3c6fe8d14a8b4a0b97155a0b22dd6d9b6626f42224be6987b0e0f1ca2b" [[package]] name = "version_check" diff --git a/python/Cargo.toml b/python/Cargo.toml index 0f2e433..b46f69c 100644 --- a/python/Cargo.toml +++ b/python/Cargo.toml @@ -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" @@ -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 diff --git a/python/pyproject.toml b/python/pyproject.toml index 5c6f7a1..8129100 100644 --- a/python/pyproject.toml +++ b/python/pyproject.toml @@ -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" diff --git a/python/rust-toolchain.toml b/python/rust-toolchain.toml index 4c8b29b..cf94bbe 100644 --- a/python/rust-toolchain.toml +++ b/python/rust-toolchain.toml @@ -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. diff --git a/rust/Cargo.lock b/rust/Cargo.lock index fcc1a62..4aa3056 100644 --- a/rust/Cargo.lock +++ b/rust/Cargo.lock @@ -1942,7 +1942,7 @@ checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981" [[package]] name = "lightstream" -version = "0.7.0" +version = "0.7.1" dependencies = [ "arrow", "arrow-flight", @@ -2066,7 +2066,9 @@ checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" [[package]] name = "minarrow" -version = "0.18.2" +version = "0.18.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4b272b619741fbb23ed37c0079a7b8dde4bd66239672b201752cef93fe13c05b" dependencies = [ "arrow", "arrow-schema", @@ -4444,9 +4446,9 @@ dependencies = [ [[package]] name = "vec64" -version = "0.5.2" +version = "0.5.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "54db55dbf75205d7a52f11a60aa6742d0ce6a9eb3ed97346556f26aec246f139" +checksum = "662d4c3c6fe8d14a8b4a0b97155a0b22dd6d9b6626f42224be6987b0e0f1ca2b" dependencies = [ "libc", ] diff --git a/rust/Cargo.toml b/rust/Cargo.toml index cc021d4..83b4a38 100644 --- a/rust/Cargo.toml +++ b/rust/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "lightstream" -version = "0.7.0" +version = "0.7.1" edition = "2024" license = "MPL-2.0" keywords = [ @@ -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 } diff --git a/rust/rust-toolchain.toml b/rust/rust-toolchain.toml index 02e8c48..b8de14d 100644 --- a/rust/rust-toolchain.toml +++ b/rust/rust-toolchain.toml @@ -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. diff --git a/rust/src/compression/ipc.rs b/rust/src/compression/ipc.rs index 1136dc6..c8c60b6 100644 --- a/rust/src/compression/ipc.rs +++ b/rust/src/compression/ipc.rs @@ -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` with all buffers placed +//! [`decompress_ipc_body`] produces a new `Vec64` 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 diff --git a/rust/src/models/codecs/mod.rs b/rust/src/models/codecs/mod.rs index ca4334e..a69a16e 100644 --- a/rust/src/models/codecs/mod.rs +++ b/rust/src/models/codecs/mod.rs @@ -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; diff --git a/rust/src/models/decoders/csv.rs b/rust/src/models/decoders/csv.rs index bfed083..f381d06 100644 --- a/rust/src/models/decoders/csv.rs +++ b/rust/src/models/decoders/csv.rs @@ -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 //! @@ -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}; diff --git a/rust/src/models/decoders/json/mod.rs b/rust/src/models/decoders/json/mod.rs index 522e03e..321b61e 100644 --- a/rust/src/models/decoders/json/mod.rs +++ b/rust/src/models/decoders/json/mod.rs @@ -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; diff --git a/rust/src/models/decoders/limits.rs b/rust/src/models/decoders/limits.rs index a96bcf4..75fde16 100644 --- a/rust/src/models/decoders/limits.rs +++ b/rust/src/models/decoders/limits.rs @@ -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; diff --git a/rust/src/models/decoders/tlv.rs b/rust/src/models/decoders/tlv.rs index cba8975..f53ea97 100644 --- a/rust/src/models/decoders/tlv.rs +++ b/rust/src/models/decoders/tlv.rs @@ -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`) @@ -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; diff --git a/rust/src/models/encoders/tlv/tlv_stream.rs b/rust/src/models/encoders/tlv/tlv_stream.rs index 9e962f9..c9b9aa8 100644 --- a/rust/src/models/encoders/tlv/tlv_stream.rs +++ b/rust/src/models/encoders/tlv/tlv_stream.rs @@ -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` or SIMD-aligned `Vec64` 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. diff --git a/rust/src/models/frames/ipc_message.rs b/rust/src/models/frames/ipc_message.rs index 9fadb48..a336e32 100644 --- a/rust/src/models/frames/ipc_message.rs +++ b/rust/src/models/frames/ipc_message.rs @@ -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. diff --git a/rust/src/models/frames/json.rs b/rust/src/models/frames/json.rs index 9a1db17..37f112c 100644 --- a/rust/src/models/frames/json.rs +++ b/rust/src/models/frames/json.rs @@ -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; diff --git a/rust/src/models/frames/lightstream_message.rs b/rust/src/models/frames/lightstream_message.rs index c42278d..4e8bde2 100644 --- a/rust/src/models/frames/lightstream_message.rs +++ b/rust/src/models/frames/lightstream_message.rs @@ -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 diff --git a/rust/src/models/frames/tlv_frame.rs b/rust/src/models/frames/tlv_frame.rs index 2d9911d..856a8cf 100644 --- a/rust/src/models/frames/tlv_frame.rs +++ b/rust/src/models/frames/tlv_frame.rs @@ -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; diff --git a/rust/src/models/interfaces/json/mod.rs b/rust/src/models/interfaces/json/mod.rs index 80f35d7..e096c5b 100644 --- a/rust/src/models/interfaces/json/mod.rs +++ b/rust/src/models/interfaces/json/mod.rs @@ -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. diff --git a/rust/src/models/interfaces/json/schema.rs b/rust/src/models/interfaces/json/schema.rs index b43eb9c..c7dd5dd 100644 --- a/rust/src/models/interfaces/json/schema.rs +++ b/rust/src/models/interfaces/json/schema.rs @@ -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)] @@ -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() } diff --git a/rust/src/models/io_uring/mod.rs b/rust/src/models/io_uring/mod.rs index 3beb30e..ee1a738 100644 --- a/rust/src/models/io_uring/mod.rs +++ b/rust/src/models/io_uring/mod.rs @@ -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 diff --git a/rust/src/models/protocol/ipc.rs b/rust/src/models/protocol/ipc.rs index 7f9f1c0..26e3990 100644 --- a/rust/src/models/protocol/ipc.rs +++ b/rust/src/models/protocol/ipc.rs @@ -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 diff --git a/rust/src/models/protocol/lightstream/connection.rs b/rust/src/models/protocol/lightstream/connection.rs index e9e903c..c420266 100644 --- a/rust/src/models/protocol/lightstream/connection.rs +++ b/rust/src/models/protocol/lightstream/connection.rs @@ -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. //! diff --git a/rust/src/models/protocol/mod.rs b/rust/src/models/protocol/mod.rs index 436830c..fbc8e44 100644 --- a/rust/src/models/protocol/mod.rs +++ b/rust/src/models/protocol/mod.rs @@ -8,12 +8,12 @@ //! //! ## Arrow IPC //! -//! The [`ipc`](crate::models::protocol::ipc) module re-exports the Arrow IPC codec, readers, writers, and +//! The [`ipc`] module re-exports the Arrow IPC codec, readers, writers, and //! sinks. Always available - no feature gate required. //! //! ## Lightstream //! -//! The [`lightstream`](crate::models::protocol::lightstream) module provides multiplexed typed messages and Arrow +//! The [`lightstream`] module provides multiplexed typed messages and Arrow //! tables over a single async stream using TLV framing on top of Arrow IPC. //! Requires the `protocol` feature. //! diff --git a/rust/src/models/readers/chunked/arrow.rs b/rust/src/models/readers/chunked/arrow.rs index 694bc55..fd28f9a 100644 --- a/rust/src/models/readers/chunked/arrow.rs +++ b/rust/src/models/readers/chunked/arrow.rs @@ -23,7 +23,7 @@ //! //! Across files the picture changes: chunk files are independent, so //! parallel reads scale roughly with disk queue depth and core count. -//! [`ChunkedTableReader::par_load_batched`](crate::traits::chunked_table_reader::ChunkedTableReader::par_load_batched) (inherited from the trait) uses +//! [`ChunkedTableReader::par_load_batched`] (inherited from the trait) uses //! `std::thread::scope` to fan per-chunk work (open + footer + body + //! decode) across worker threads and returns a `SuperTable` with batches //! in write order. diff --git a/rust/src/models/readers/chunked/csv.rs b/rust/src/models/readers/chunked/csv.rs index 6ccc5a1..3caa2b9 100644 --- a/rust/src/models/readers/chunked/csv.rs +++ b/rust/src/models/readers/chunked/csv.rs @@ -12,7 +12,7 @@ //! rows). The reader sorts files by their numeric index so consumers see //! batches in write order. //! -//! Inherits [`ChunkedTableReader::par_load_batched`](crate::traits::chunked_table_reader::ChunkedTableReader::par_load_batched) for sync parallel +//! Inherits [`ChunkedTableReader::par_load_batched`] for sync parallel //! decode across chunk files. use std::fs::{self, File}; diff --git a/rust/src/models/readers/chunked/parquet.rs b/rust/src/models/readers/chunked/parquet.rs index 270fdfd..1fc0f14 100644 --- a/rust/src/models/readers/chunked/parquet.rs +++ b/rust/src/models/readers/chunked/parquet.rs @@ -17,7 +17,7 @@ //! The default [`Iterator`] path is sync and serial. Per-file Parquet //! decode is CPU-heavy (decompression, page parsing, dictionary //! resolution), so across files the gain from parallel reads is larger -//! than for raw IPC. [`ChunkedTableReader::par_load_batched`](crate::traits::chunked_table_reader::ChunkedTableReader::par_load_batched) (inherited from +//! than for raw IPC. [`ChunkedTableReader::par_load_batched`] (inherited from //! the trait) uses `std::thread::scope` to fan per-chunk work end-to-end //! across worker threads and returns a `SuperTable` with batches in //! write order. diff --git a/rust/src/models/readers/http.rs b/rust/src/models/readers/http.rs index d127b18..16d4d3f 100644 --- a/rust/src/models/readers/http.rs +++ b/rust/src/models/readers/http.rs @@ -11,11 +11,11 @@ //! Service layer is not in the dep tree. //! //! Plug-and-play one-liner over `http://` is [`HttpTableReader::get`](crate::models::readers::http::HttpTableReader::get); -//! over `https://` is [`HttpTableReader::get_tls`](crate::models::readers::http::HttpTableReader::get_tls) (requires the `tls` +//! over `https://` is [`HttpTableReader::get_tls`] (requires the `tls` //! feature and a caller-supplied `rustls::ClientConfig` with ALPN set //! to `h2`). Callers that need custom headers (auth tokens, API keys) //! pass a fully-built `http::Request<()>` via -//! [`HttpTableReader::from_request`](crate::models::readers::http::HttpTableReader::from_request). +//! [`HttpTableReader::from_request`]. //! //! ## Continuous streaming //! diff --git a/rust/src/models/readers/json.rs b/rust/src/models/readers/json.rs index 867dc03..653049c 100644 --- a/rust/src/models/readers/json.rs +++ b/rust/src/models/readers/json.rs @@ -8,12 +8,12 @@ //! //! High-level API for reading JSON files or streams into Minarrow Tables. //! Supports either a single JSON array-of-objects or streaming NDJSON -//! (newline-delimited), chosen via [`JsonFormat`](crate::models::encoders::json::JsonFormat). +//! (newline-delimited), chosen via [`JsonFormat`]. //! //! For NDJSON, records can be pulled in fixed-size batches via -//! [`JsonReader::next_batch`](crate::models::readers::json::JsonReader::next_batch); array-of-objects must be fully read in one pass. +//! [`JsonReader::next_batch`]; array-of-objects must be fully read in one pass. //! -//! See [`JsonDecodeOptions`](crate::models::decoders::json::JsonDecodeOptions) for schema handling and type control. +//! See [`JsonDecodeOptions`] for schema handling and type control. use std::fs::File; use std::io::{self, BufRead, BufReader}; diff --git a/rust/src/models/readers/lightstream.rs b/rust/src/models/readers/lightstream.rs index 2cf9149..62aa1a8 100644 --- a/rust/src/models/readers/lightstream.rs +++ b/rust/src/models/readers/lightstream.rs @@ -6,8 +6,8 @@ //! Async Lightstream protocol reader. //! -//! Reads TLV frames from an [`AsyncRead`](tokio::io::AsyncRead) source, decoding each into a -//! [`LightstreamMessage`](crate::models::frames::lightstream_message::LightstreamMessage) via the codec's type registry. +//! Reads TLV frames from an [`AsyncRead`] source, decoding each into a +//! [`LightstreamMessage`] via the codec's type registry. //! //! The reader accumulates the 5-byte TLV header, then reads the payload //! into a Vec64 for zero-copy decode. Column data is mapped in place diff --git a/rust/src/models/readers/parallel/http.rs b/rust/src/models/readers/parallel/http.rs index a8d433d..63f7ca2 100644 --- a/rust/src/models/readers/parallel/http.rs +++ b/rust/src/models/readers/parallel/http.rs @@ -14,9 +14,9 @@ //! //! The h2 server connection is the I/O driver for the in-flight request //! bodies, so a background task keeps it polled while the accepted streams -//! decode. Under [`SortBehaviour::None`](crate::traits::parallel_transport_reader::SortBehaviour::None) and [`SortBehaviour::RequestKeys`](crate::traits::parallel_transport_reader::SortBehaviour::RequestKeys) +//! decode. Under [`SortBehaviour::None`] and [`SortBehaviour::RequestKeys`] //! tables surface in the order the streams produce them. Under -//! [`SortBehaviour::Ordered`](crate::traits::parallel_transport_reader::SortBehaviour::Ordered) the reader pulls the streams in the writer's +//! [`SortBehaviour::Ordered`] the reader pulls the streams in the writer's //! round-robin rotation, so tables surface in global write order. use std::io; diff --git a/rust/src/models/readers/parallel/lightstream.rs b/rust/src/models/readers/parallel/lightstream.rs index b3d68ad..f6fa474 100644 --- a/rust/src/models/readers/parallel/lightstream.rs +++ b/rust/src/models/readers/parallel/lightstream.rs @@ -7,15 +7,15 @@ //! # Parallel Lightstream protocol reader //! //! Accepts several concurrent Lightstream protocol connections on a -//! [`TcpListener`](tokio::net::TcpListener) and decodes them across cores, one task per connection. +//! [`TcpListener`] and decodes them across cores, one task per connection. //! Each task feeds its own channel, and the reader merges the channels into a -//! single frame stream of [`LightstreamMessage`](crate::models::frames::lightstream_message::LightstreamMessage) values, so protobuf messages +//! single frame stream of [`LightstreamMessage`] values, so protobuf messages //! and Arrow tables share the wire. //! //! Message and table types are registered on every connection at -//! [`accept`](crate::models::readers::parallel::lightstream::LightstreamParallelReader::accept). Under [`SortBehaviour::None`](crate::traits::parallel_transport_reader::SortBehaviour::None) -//! and [`SortBehaviour::RequestKeys`](crate::traits::parallel_transport_reader::SortBehaviour::RequestKeys) frames surface in the order the -//! connections produce them. Under [`SortBehaviour::Ordered`](crate::traits::parallel_transport_reader::SortBehaviour::Ordered) the reader pulls +//! [`accept`](crate::models::readers::parallel::lightstream::LightstreamParallelReader::accept). Under [`SortBehaviour::None`] +//! and [`SortBehaviour::RequestKeys`] frames surface in the order the +//! connections produce them. Under [`SortBehaviour::Ordered`] the reader pulls //! the connections in the writer's round-robin rotation, so frames surface in //! global send order. Each connection announces its index before any frames, //! so it is placed by that index rather than by accept order - the global order diff --git a/rust/src/models/readers/parallel/quic.rs b/rust/src/models/readers/parallel/quic.rs index 0499ae9..d8092b0 100644 --- a/rust/src/models/readers/parallel/quic.rs +++ b/rust/src/models/readers/parallel/quic.rs @@ -13,9 +13,9 @@ //! sequence key - `Some` when the peer used an ordered writer, `None` //! otherwise. //! -//! Under [`SortBehaviour::None`](crate::traits::parallel_transport_reader::SortBehaviour::None) and [`SortBehaviour::RequestKeys`](crate::traits::parallel_transport_reader::SortBehaviour::RequestKeys) tables +//! Under [`SortBehaviour::None`] and [`SortBehaviour::RequestKeys`] tables //! surface in the order the streams produce them. Under -//! [`SortBehaviour::Ordered`](crate::traits::parallel_transport_reader::SortBehaviour::Ordered) the reader pulls the streams in the writer's +//! [`SortBehaviour::Ordered`] the reader pulls the streams in the writer's //! round-robin rotation, so tables surface in global write order. use std::io; diff --git a/rust/src/models/readers/parallel/tcp.rs b/rust/src/models/readers/parallel/tcp.rs index f3567dd..c082955 100644 --- a/rust/src/models/readers/parallel/tcp.rs +++ b/rust/src/models/readers/parallel/tcp.rs @@ -6,15 +6,15 @@ //! # Parallel TCP table reader //! -//! Accepts several concurrent TCP connections on a [`TcpListener`](tokio::net::TcpListener) and decodes +//! Accepts several concurrent TCP connections on a [`TcpListener`] and decodes //! them across cores, one task per connection. Each task feeds its own //! channel, and the reader merges the channels into a single table stream. //! Each table is paired with its sequence key - `Some` when the peer used an //! ordered writer, `None` otherwise. //! -//! Under [`SortBehaviour::None`](crate::traits::parallel_transport_reader::SortBehaviour::None) and [`SortBehaviour::RequestKeys`](crate::traits::parallel_transport_reader::SortBehaviour::RequestKeys) tables +//! Under [`SortBehaviour::None`] and [`SortBehaviour::RequestKeys`] tables //! surface in the order the connections produce them. Under -//! [`SortBehaviour::Ordered`](crate::traits::parallel_transport_reader::SortBehaviour::Ordered) the reader pulls the connections in the writer's +//! [`SortBehaviour::Ordered`] the reader pulls the connections in the writer's //! round-robin rotation, so tables surface in global write order. Connections //! are accepted in order, so the `i`-th accepted connection pairs with the //! writer's `i`-th connection. diff --git a/rust/src/models/readers/quic.rs b/rust/src/models/readers/quic.rs index c321601..fcab712 100644 --- a/rust/src/models/readers/quic.rs +++ b/rust/src/models/readers/quic.rs @@ -9,7 +9,7 @@ //! High-level async reader that wraps a QUIC receive stream and decodes //! Arrow IPC data into MinArrow tables. //! -//! Wraps [`TableReader`](crate::models::readers::ipc::table::TableReader) over a [`QuicByteStream`](crate::models::streams::quic::QuicByteStream), hiding the wiring +//! Wraps [`TableReader`] over a [`QuicByteStream`], hiding the wiring //! so callers get a one-liner API. //! //! ## Continuous streaming diff --git a/rust/src/models/readers/stdio.rs b/rust/src/models/readers/stdio.rs index 2e933e9..8c31234 100644 --- a/rust/src/models/readers/stdio.rs +++ b/rust/src/models/readers/stdio.rs @@ -9,7 +9,7 @@ //! High-level async reader that reads Arrow IPC data from stdin //! and decodes it into MinArrow tables. //! -//! Wraps [`TableReader`](crate::models::readers::ipc::table::TableReader) over a [`StdinByteStream`](crate::models::streams::stdio::StdinByteStream), hiding the wiring +//! Wraps [`TableReader`] over a [`StdinByteStream`], hiding the wiring //! so callers get a simple API for CLI tools. //! //! ## Continuous streaming diff --git a/rust/src/models/readers/tcp.rs b/rust/src/models/readers/tcp.rs index f49363c..9329f86 100644 --- a/rust/src/models/readers/tcp.rs +++ b/rust/src/models/readers/tcp.rs @@ -9,7 +9,7 @@ //! High-level async reader that connects to a TCP endpoint streaming //! Arrow IPC data and decodes it into MinArrow tables. //! -//! Wraps [`TableReader`](crate::models::readers::ipc::table::TableReader) over a [`TcpByteStream`](crate::models::streams::tcp::TcpByteStream), hiding the wiring +//! Wraps [`TableReader`] over a [`TcpByteStream`], hiding the wiring //! so callers get a one-liner API. //! //! ## Continuous streaming diff --git a/rust/src/models/readers/uds.rs b/rust/src/models/readers/uds.rs index 4593d08..0f8d1b3 100644 --- a/rust/src/models/readers/uds.rs +++ b/rust/src/models/readers/uds.rs @@ -9,7 +9,7 @@ //! High-level async reader that connects to a UDS endpoint streaming //! Arrow IPC data and decodes it into MinArrow tables. //! -//! Wraps [`TableReader`](crate::models::readers::ipc::table::TableReader) over a [`UdsByteStream`](crate::models::streams::uds::UdsByteStream), hiding the wiring +//! Wraps [`TableReader`] over a [`UdsByteStream`], hiding the wiring //! so callers get a one-liner API. //! //! ## Continuous streaming diff --git a/rust/src/models/readers/websocket.rs b/rust/src/models/readers/websocket.rs index ad6dbc1..921af9a 100644 --- a/rust/src/models/readers/websocket.rs +++ b/rust/src/models/readers/websocket.rs @@ -10,7 +10,7 @@ //! Arrow IPC data and decodes it into MinArrow tables. //! //! Extracts the raw TCP stream after the tungstenite handshake and uses -//! [`WsRead`](crate::models::streams::websocket::WsRead) for zero-copy WebSocket frame parsing on the data path. +//! [`WsRead`] for zero-copy WebSocket frame parsing on the data path. //! //! ## Security //! @@ -20,7 +20,7 @@ //! that integration is compiled in. //! //! For pinned roots, a custom verifier, or client-auth keys, use -//! [`WebSocketTableReader::connect_tls`](crate::models::readers::websocket::WebSocketTableReader::connect_tls) - it takes an +//! [`WebSocketTableReader::connect_tls`] - it takes an //! `Arc` directly and bypasses the bundled //! verifier. The library does not enforce a transport policy; if a //! deployment requires TLS, that is the caller's deployment decision. diff --git a/rust/src/models/readers/webtransport.rs b/rust/src/models/readers/webtransport.rs index f5979db..74f0299 100644 --- a/rust/src/models/readers/webtransport.rs +++ b/rust/src/models/readers/webtransport.rs @@ -9,7 +9,7 @@ //! High-level async reader that wraps a WebTransport receive stream and decodes //! Arrow IPC data into MinArrow tables. //! -//! Wraps [`TableReader`](crate::models::readers::ipc::table::TableReader) over a [`WebTransportByteStream`](crate::models::streams::webtransport::WebTransportByteStream), hiding the wiring +//! Wraps [`TableReader`] over a [`WebTransportByteStream`], hiding the wiring //! so callers get a one-liner API. //! //! ## Stability: unstable diff --git a/rust/src/models/serialise/ipc.rs b/rust/src/models/serialise/ipc.rs index 25e4cb5..9997464 100644 --- a/rust/src/models/serialise/ipc.rs +++ b/rust/src/models/serialise/ipc.rs @@ -6,14 +6,14 @@ //! Arrow IPC serialisation for minarrow value types. //! -//! [`Table`](minarrow::Table) provides the base implementation. Other value types are converted +//! [`Table`] provides the base implementation. Other value types are converted //! to a table-compatible representation before encoding and converted back //! after decoding. //! -//! [`SuperTable`](minarrow::SuperTable) and [`SuperArray`](minarrow::SuperArray) use the streaming codec API to encode and +//! [`SuperTable`] and [`SuperArray`] use the streaming codec API to encode and //! decode multiple batches in one IPC stream. //! -//! [`IpcSerialise`](crate::models::serialise::ipc::IpcSerialise) can be used as a generic bound for types that support this +//! [`IpcSerialise`] can be used as a generic bound for types that support this //! Arrow IPC round trip. use std::sync::Arc; diff --git a/rust/src/models/sinks/table_sink.rs b/rust/src/models/sinks/table_sink.rs index a80d273..8851ddb 100644 --- a/rust/src/models/sinks/table_sink.rs +++ b/rust/src/models/sinks/table_sink.rs @@ -11,7 +11,7 @@ //! ## Overview: //! - Supports Stream or File protocol //! - Handles schema emission, optional compression, dictionary batches, record batches, and end-of-stream/footer generation. -//! - Supports both 8-byte (`Vec`) and 64-byte SIMD-aligned (`Vec64`) buffers via [`TableSink`](crate::models::sinks::table_sink::TableSink) and [`TableSink64`](crate::models::sinks::table_sink::TableSink64) type aliases. +//! - Supports both 8-byte (`Vec`) and 64-byte SIMD-aligned (`Vec64`) buffers via [`TableSink`] and [`TableSink64`] type aliases. //! - Supports backpressure-friendly, chunked writes with partial-write handling in async runtimes (e.g. Tokio). use crate::compression::Compression; diff --git a/rust/src/models/sinks/tlv_sink.rs b/rust/src/models/sinks/tlv_sink.rs index 0f62cfa..b14c1ac 100644 --- a/rust/src/models/sinks/tlv_sink.rs +++ b/rust/src/models/sinks/tlv_sink.rs @@ -7,7 +7,7 @@ //! # Asynchronous TLV sink //! //! Wraps any `AsyncWrite` and streams Type-Length-Value (TLV) frames produced -//! by [`TLVStreamWriter`](crate::models::encoders::tlv::tlv_stream::TLVStreamWriter). +//! by [`TLVStreamWriter`]. //! //! ## Overview: //! - Encodes frames as `[type: u8][length: u32 LE][value: bytes]`. diff --git a/rust/src/models/streams/async_read.rs b/rust/src/models/streams/async_read.rs index 62e9ae4..ec02b6a 100644 --- a/rust/src/models/streams/async_read.rs +++ b/rust/src/models/streams/async_read.rs @@ -6,13 +6,13 @@ //! # Generic asynchronous byte stream adapter //! -//! Wraps any [`AsyncRead`](tokio::io::AsyncRead) source as both [`AsyncRead`](tokio::io::AsyncRead) and [`Stream`](futures_core::Stream). +//! Wraps any [`AsyncRead`] source as both [`AsyncRead`] and [`Stream`]. //! //! ## AsyncRead //! Passthrough to the inner source for the direct decode path. //! //! ## Stream -//! Yields [`SharedBuffer`](minarrow::structs::shared_buffer::SharedBuffer) windows from a [`StreamArena`](crate::models::streams::stream_arena::StreamArena) for zero-allocation +//! Yields [`SharedBuffer`] windows from a [`StreamArena`] for zero-allocation //! streaming. This is the generic building block behind transport-specific //! byte streams - QUIC, WebTransport, and Stdin are type aliases over this. diff --git a/rust/src/models/streams/disk.rs b/rust/src/models/streams/disk.rs index 2e7343d..3190e6a 100644 --- a/rust/src/models/streams/disk.rs +++ b/rust/src/models/streams/disk.rs @@ -6,13 +6,13 @@ //! # Asynchronous disk byte stream //! -//! Wraps a file in a [`Stream`](futures_core::Stream) that yields fixed-size byte chunks. +//! Wraps a file in a [`Stream`] that yields fixed-size byte chunks. //! //! ## Overview -//! - Uses Tokio [`File`](tokio::fs::File). +//! - Uses Tokio [`File`]. //! - Supports async backpressure via `poll_next`. //! - One copy into a `Vec64` output buffer per chunk. -//! - Chunk size controlled by [`BufferChunkSize`](crate::enums::BufferChunkSize). +//! - Chunk size controlled by [`BufferChunkSize`]. //! //! ## Use cases //! - Ingest large files without loading them fully into memory. diff --git a/rust/src/models/streams/framed_byte_stream.rs b/rust/src/models/streams/framed_byte_stream.rs index 298977b..75535d2 100644 --- a/rust/src/models/streams/framed_byte_stream.rs +++ b/rust/src/models/streams/framed_byte_stream.rs @@ -7,7 +7,7 @@ //! # Generic async framed byte stream //! //! Adapts any chunked byte source into a stream of protocol frames using a -//! user-supplied [`FrameDecoder`](crate::traits::frame_decoder::FrameDecoder). +//! user-supplied [`FrameDecoder`]. //! //! - Works with any `GenByteStream` (e.g., network/file sources). //! - Buffers partial input and yields complete frames as soon as available. diff --git a/rust/src/models/streams/http.rs b/rust/src/models/streams/http.rs index 53ad72f..ec8c41c 100644 --- a/rust/src/models/streams/http.rs +++ b/rust/src/models/streams/http.rs @@ -6,10 +6,10 @@ //! # HTTP/2 byte stream adapters //! -//! Receive side: [`H2RecvRead`](crate::models::streams::http::H2RecvRead) wraps [`h2::RecvStream`] as `AsyncRead` -//! so it flows through [`AsyncReadByteStream`](crate::models::streams::async_read::AsyncReadByteStream) like every other transport. +//! Receive side: [`H2RecvRead`] wraps [`h2::RecvStream`] as `AsyncRead` +//! so it flows through [`AsyncReadByteStream`] like every other transport. //! -//! Send side: [`H2SendWrite`](crate::models::streams::http::H2SendWrite) wraps [`h2::SendStream`] as +//! Send side: [`H2SendWrite`] wraps [`h2::SendStream`] as //! `AsyncWrite` so `HttpTableWriter` can hold it inside `TableSink64` //! exactly the same way `QuicTableWriter` holds a `quinn::SendStream`. diff --git a/rust/src/models/streams/quic.rs b/rust/src/models/streams/quic.rs index c71abd9..abdb0aa 100644 --- a/rust/src/models/streams/quic.rs +++ b/rust/src/models/streams/quic.rs @@ -6,7 +6,7 @@ //! # Asynchronous QUIC byte stream //! -//! Type alias over [`AsyncReadByteStream`](crate::models::streams::async_read::AsyncReadByteStream) for QUIC receive streams. +//! Type alias over [`AsyncReadByteStream`] for QUIC receive streams. //! //! ## Use cases //! - Receive Arrow IPC streams over QUIC without loading them fully into memory. diff --git a/rust/src/models/streams/stdio.rs b/rust/src/models/streams/stdio.rs index f4f1e46..54bc3d3 100644 --- a/rust/src/models/streams/stdio.rs +++ b/rust/src/models/streams/stdio.rs @@ -6,7 +6,7 @@ //! # Asynchronous stdin byte stream //! -//! Type alias over [`AsyncReadByteStream`](crate::models::streams::async_read::AsyncReadByteStream) for standard input. +//! Type alias over [`AsyncReadByteStream`] for standard input. //! //! ## Use cases //! - Receive Arrow IPC streams from Unix pipes without loading them fully into memory. diff --git a/rust/src/models/streams/tcp.rs b/rust/src/models/streams/tcp.rs index 7818868..3214cf6 100644 --- a/rust/src/models/streams/tcp.rs +++ b/rust/src/models/streams/tcp.rs @@ -6,14 +6,14 @@ //! # Asynchronous TCP byte stream //! -//! Wraps a TCP connection's read half as both [`AsyncRead`](tokio::io::AsyncRead) and [`Stream`](futures_core::Stream). +//! Wraps a TCP connection's read half as both [`AsyncRead`] and [`Stream`]. //! //! ## AsyncRead //! The direct decode path uses `AsyncRead` for zero-copy reads into the //! decoder's managed buffers. This is the internal fast path. //! //! ## Stream -//! Yields [`SharedBuffer`](minarrow::structs::shared_buffer::SharedBuffer) windows from a [`StreamArena`](crate::models::streams::stream_arena::StreamArena) for zero-allocation +//! Yields [`SharedBuffer`] windows from a [`StreamArena`] for zero-allocation //! streaming. Each poll reads into the arena's spare capacity and yields an //! immutable view of the filled region. In steady state, one arena allocation //! is reused forever. diff --git a/rust/src/models/streams/uds.rs b/rust/src/models/streams/uds.rs index f52206b..8b34c9e 100644 --- a/rust/src/models/streams/uds.rs +++ b/rust/src/models/streams/uds.rs @@ -6,14 +6,14 @@ //! # Asynchronous Unix domain socket byte stream //! -//! Wraps a UDS connection's read half as both [`AsyncRead`](tokio::io::AsyncRead) and [`Stream`](futures_core::Stream). +//! Wraps a UDS connection's read half as both [`AsyncRead`] and [`Stream`]. //! //! ## AsyncRead //! The direct decode path uses `AsyncRead` for zero-copy reads into the //! decoder's managed buffers. This is the internal fast path. //! //! ## Stream -//! Yields [`SharedBuffer`](minarrow::structs::shared_buffer::SharedBuffer) windows from a [`StreamArena`](crate::models::streams::stream_arena::StreamArena) for zero-allocation +//! Yields [`SharedBuffer`] windows from a [`StreamArena`] for zero-allocation //! streaming. Each poll reads into the arena's spare capacity and yields an //! immutable view of the filled region. diff --git a/rust/src/models/streams/websocket.rs b/rust/src/models/streams/websocket.rs index 11185ad..d72e8bc 100644 --- a/rust/src/models/streams/websocket.rs +++ b/rust/src/models/streams/websocket.rs @@ -6,14 +6,14 @@ //! # WebSocket byte stream adapters //! -//! Provides [`WsRead`](crate::models::streams::websocket::WsRead) and [`WsWrite`](crate::models::streams::websocket::WsWrite) for WebSocket I/O over a raw TCP +//! Provides [`WsRead`] and [`WsWrite`] for WebSocket I/O over a raw TCP //! stream extracted after the tungstenite handshake. //! //! WS frame parsing and construction happens inline with zero intermediate //! allocations. Payload bytes flow between the TCP socket and the caller's -//! buffer via the standard [`AsyncRead`](tokio::io::AsyncRead) / [`AsyncWrite`](tokio::io::AsyncWrite) traits. +//! buffer via the standard [`AsyncRead`] / [`AsyncWrite`] traits. //! -//! [`WsRead`](crate::models::streams::websocket::WsRead) also implements [`Stream`](futures_core::Stream) yielding arena-backed [`SharedBuffer`](minarrow::structs::shared_buffer::SharedBuffer) +//! [`WsRead`] also implements [`Stream`] yielding arena-backed [`SharedBuffer`] //! windows for consumers that prefer the StreamExt API. use std::io; diff --git a/rust/src/models/streams/webtransport.rs b/rust/src/models/streams/webtransport.rs index 1e9e63c..b65b56c 100644 --- a/rust/src/models/streams/webtransport.rs +++ b/rust/src/models/streams/webtransport.rs @@ -6,7 +6,7 @@ //! # Asynchronous WebTransport byte stream //! -//! Type alias over [`AsyncReadByteStream`](crate::models::streams::async_read::AsyncReadByteStream) for WebTransport receive streams. +//! Type alias over [`AsyncReadByteStream`] for WebTransport receive streams. //! //! ## Use cases //! - Receive Arrow IPC streams over WebTransport without loading them fully into memory. diff --git a/rust/src/models/transports/quic.rs b/rust/src/models/transports/quic.rs index d1f842c..1257e30 100644 --- a/rust/src/models/transports/quic.rs +++ b/rust/src/models/transports/quic.rs @@ -8,7 +8,7 @@ //! //! Connection establishment for QUIC in both peer roles. QUIC requires //! TLS, so each role carries its own quinn configuration through -//! [`QuicEndpoint`](crate::models::transports::quic::QuicEndpoint). +//! [`QuicEndpoint`]. //! Each established connection opens one //! bidirectional stream whose receive and send sides are the returned //! halves. diff --git a/rust/src/models/writers/csv.rs b/rust/src/models/writers/csv.rs index 5ecf307..05cf9dd 100644 --- a/rust/src/models/writers/csv.rs +++ b/rust/src/models/writers/csv.rs @@ -11,9 +11,9 @@ //! //! ## Features //! - Pluggable destination: in-memory (`Vec`) or files (`std::fs::File`) -//! - Configurable delimiter, header emission, and null representation via [`CsvEncodeOptions`](crate::models::encoders::csv::CsvEncodeOptions) +//! - Configurable delimiter, header emission, and null representation via [`CsvEncodeOptions`] //! - RFC 4180-style quoting/escaping handled by the encoders -//! - Write single tables or concatenate multi-batch [`SuperTable`](minarrow::SuperTable)s - header on first batch only +//! - Write single tables or concatenate multi-batch [`SuperTable`]s - header on first batch only //! //! ## Quick start //! ```no_run diff --git a/rust/src/models/writers/http.rs b/rust/src/models/writers/http.rs index f0a7a06..58f27fb 100644 --- a/rust/src/models/writers/http.rs +++ b/rust/src/models/writers/http.rs @@ -9,17 +9,17 @@ //! POSTs an Arrow IPC stream to an HTTP/2 endpoint. Transport is `h2` //! directly - hyper's Body / Service layer is not in the dep tree. //! -//! Wraps a [`TableSink64`](crate::models::sinks::table_sink::TableSink64) over an [`H2SendWrite`](crate::models::streams::http::H2SendWrite) adapter, matching +//! Wraps a [`TableSink64`] over an [`H2SendWrite`] adapter, matching //! the structural shape of every other lightstream transport writer //! (`TcpTableWriter`, `QuicTableWriter`, ...). The encoder's //! `encode_buf` is reused across frames as in `TableSink64`, and the //! `H2SendWrite::poll_write` adapter feeds h2 with one chunk per //! flow-control grant. //! -//! Plug-and-play one-liner over `http://` is [`HttpTableWriter::post`](crate::models::writers::http::HttpTableWriter::post); -//! over `https://` is [`HttpTableWriter::post_tls`](crate::models::writers::http::HttpTableWriter::post_tls). Callers that need +//! Plug-and-play one-liner over `http://` is [`HttpTableWriter::post`]; +//! over `https://` is [`HttpTableWriter::post_tls`]. Callers that need //! custom headers pass a fully-built `http::Request<()>` via -//! [`HttpTableWriter::from_request`](crate::models::writers::http::HttpTableWriter::from_request). +//! [`HttpTableWriter::from_request`]. use std::io; use std::pin::Pin; diff --git a/rust/src/models/writers/ipc/sync_table.rs b/rust/src/models/writers/ipc/sync_table.rs index 867fda8..e3f467e 100644 --- a/rust/src/models/writers/ipc/sync_table.rs +++ b/rust/src/models/writers/ipc/sync_table.rs @@ -16,7 +16,7 @@ //! `next_frame` and is responsible for getting bytes onto the wire. //! - [`TableWriter`] - **async end-to-end writer**. Wraps `TableStreamWriter` //! and drives the queued frames into a `tokio::io::AsyncWrite` sink. -//! - [`SyncTableWriter`](crate::models::writers::ipc::sync_table::SyncTableWriter) (this struct) - **sync end-to-end writer**. Wraps +//! - [`SyncTableWriter`] (this struct) - **sync end-to-end writer**. Wraps //! `TableStreamWriter` and drives the queued frames into a //! `std::io::Write` sink. Use this for one-shot file writes from a sync //! caller (e.g. `ChunkedArrowWriter`'s per-batch path) so the caller does diff --git a/rust/src/models/writers/ipc/table.rs b/rust/src/models/writers/ipc/table.rs index 6fdcfbf..a3cc12e 100644 --- a/rust/src/models/writers/ipc/table.rs +++ b/rust/src/models/writers/ipc/table.rs @@ -10,16 +10,16 @@ //! in Arrow IPC format using its `File` or `Stream` protocols. //! //! ## Features -//! - Wraps [`GTableSink`](crate::models::sinks::table_sink::GTableSink) for Arrow IPC framing and schema management +//! - Wraps [`GTableSink`] for Arrow IPC framing and schema management //! - Supports optional compression //! - Automatically registers categorical dictionaries //! - Provides helpers for writing one or many tables to disk //! //! ## Typical usage -//! - Create with [`TableWriter::new`](crate::models::writers::ipc::table::TableWriter::new) or `TableWriter::new_with_compression` -//! - Optionally register dictionaries with [`TableWriter::register_dictionary`](crate::models::writers::ipc::table::TableWriter::register_dictionary) -//! - Write tables using [`TableWriter::write_table`](crate::models::writers::ipc::table::TableWriter::write_table) or [`TableWriter::write_all_tables`](crate::models::writers::ipc::table::TableWriter::write_all_tables) -//! - Finalise with [`TableWriter::finish`](crate::models::writers::ipc::table::TableWriter::finish) +//! - Create with [`TableWriter::new`] or `TableWriter::new_with_compression` +//! - Optionally register dictionaries with [`TableWriter::register_dictionary`] +//! - Write tables using [`TableWriter::write_table`] or [`TableWriter::write_all_tables`] +//! - Finalise with [`TableWriter::finish`] use std::io; diff --git a/rust/src/models/writers/ipc/table_stream.rs b/rust/src/models/writers/ipc/table_stream.rs index 640d082..cfd8e2f 100644 --- a/rust/src/models/writers/ipc/table_stream.rs +++ b/rust/src/models/writers/ipc/table_stream.rs @@ -16,8 +16,8 @@ //! - Frames can be pulled incrementally (`next_frame`) or drained all at once //! //! ## Async Helpers -//! - [`write_tables_to_stream`](crate::models::writers::ipc::table_stream::write_tables_to_stream) - write a sequence of tables to an async sink. -//! - [`write_table_to_stream`](crate::models::writers::ipc::table_stream::write_table_to_stream) - write a single table to an async sink. +//! - [`write_tables_to_stream`] - write a sequence of tables to an async sink. +//! - [`write_table_to_stream`] - write a single table to an async sink. //! //! ## Usage //! ```ignore diff --git a/rust/src/models/writers/json.rs b/rust/src/models/writers/json.rs index 14f1bbb..e230697 100644 --- a/rust/src/models/writers/json.rs +++ b/rust/src/models/writers/json.rs @@ -12,9 +12,9 @@ //! //! ## Features //! - Pluggable destination: in-memory `Vec` or [`minarrow::Vec64`] (64-byte aligned), -//! files, or any user-supplied [`Write`](std::io::Write) impl -//! - Configurable output shape and formatting via [`JsonEncodeOptions`](crate::models::encoders::json::JsonEncodeOptions) -//! - Writes single tables or multi-batch [`SuperTable`](minarrow::SuperTable)s +//! files, or any user-supplied [`Write`] impl +//! - Configurable output shape and formatting via [`JsonEncodeOptions`] +//! - Writes single tables or multi-batch [`SuperTable`]s //! //! JSON is a text format; output bytes are not consumed as SIMD payloads, so //! a standard `Vec` is the ergonomic default. A 64-byte aligned `Vec64` diff --git a/rust/src/models/writers/lightstream.rs b/rust/src/models/writers/lightstream.rs index 186aae3..972cf2e 100644 --- a/rust/src/models/writers/lightstream.rs +++ b/rust/src/models/writers/lightstream.rs @@ -6,7 +6,7 @@ //! Async Lightstream protocol writer. //! -//! Wraps a [`LightstreamCodec`](crate::models::codecs::lightstream::LightstreamCodec) and an [`AsyncWrite`](tokio::io::AsyncWrite) destination, providing +//! Wraps a [`LightstreamCodec`] and an [`AsyncWrite`] destination, providing //! methods to send messages and Arrow tables over a single connection. //! //! Tables are encoded using the Arrow IPC streaming protocol - schema is diff --git a/rust/src/models/writers/parallel/http.rs b/rust/src/models/writers/parallel/http.rs index 8c36b0c..a24dd47 100644 --- a/rust/src/models/writers/parallel/http.rs +++ b/rust/src/models/writers/parallel/http.rs @@ -8,7 +8,7 @@ //! //! Fans one table sequence across several concurrent HTTP/2 request //! streams on a single h2 client connection. Each stream runs its own -//! [`HttpTableWriter`](crate::models::writers::http::HttpTableWriter) driven by a dedicated task, so the streams upload +//! [`HttpTableWriter`] driven by a dedicated task, so the streams upload //! in parallel and aggregate throughput is the sum across them. //! //! Each stream carries an independent ordered sequence of batches; diff --git a/rust/src/models/writers/parallel/lightstream.rs b/rust/src/models/writers/parallel/lightstream.rs index 7b75141..5368068 100644 --- a/rust/src/models/writers/parallel/lightstream.rs +++ b/rust/src/models/writers/parallel/lightstream.rs @@ -8,7 +8,7 @@ //! //! Fans one frame sequence across several concurrent Lightstream protocol //! connections to a single endpoint. Each connection runs its own -//! [`LightstreamWriter`](crate::models::writers::lightstream::LightstreamWriter) driven by a dedicated task, so the connections send +//! [`LightstreamWriter`] driven by a dedicated task, so the connections send //! in parallel and aggregate throughput is the sum across them. //! //! Message types and table types are registered on every connection at diff --git a/rust/src/models/writers/parallel/quic.rs b/rust/src/models/writers/parallel/quic.rs index 22be895..a0459de 100644 --- a/rust/src/models/writers/parallel/quic.rs +++ b/rust/src/models/writers/parallel/quic.rs @@ -8,7 +8,7 @@ //! //! Fans one table sequence across several concurrent QUIC streams on a //! single [`quinn::Connection`]. Each stream runs its own -//! [`QuicTableWriter`](crate::models::writers::quic::QuicTableWriter) driven by a dedicated task, so the streams send +//! [`QuicTableWriter`] driven by a dedicated task, so the streams send //! in parallel and aggregate throughput is the sum across them. //! //! Each stream carries an independent ordered sequence of batches; diff --git a/rust/src/models/writers/parallel/tcp.rs b/rust/src/models/writers/parallel/tcp.rs index 48f1a48..dc9c724 100644 --- a/rust/src/models/writers/parallel/tcp.rs +++ b/rust/src/models/writers/parallel/tcp.rs @@ -8,7 +8,7 @@ //! //! Fans one table sequence across several concurrent TCP connections to a //! single endpoint. TCP has no in-band stream multiplexing, so each "stream" -//! is its own connection running its own [`TcpTableWriter`](crate::models::writers::tcp::TcpTableWriter) driven by a +//! is its own connection running its own [`TcpTableWriter`] driven by a //! dedicated task. The connections send in parallel and aggregate throughput //! is the sum across them. //! diff --git a/rust/src/models/writers/quic.rs b/rust/src/models/writers/quic.rs index 7fd0f02..3ef0349 100644 --- a/rust/src/models/writers/quic.rs +++ b/rust/src/models/writers/quic.rs @@ -9,7 +9,7 @@ //! High-level async writer that sends Arrow IPC encoded tables over a //! QUIC send stream. //! -//! Wraps a [`TableSink64`](crate::models::sinks::table_sink::TableSink64) over a [`quinn::SendStream`], hiding the wiring +//! Wraps a [`TableSink64`] over a [`quinn::SendStream`], hiding the wiring //! so callers get a one-liner API. //! //! Uses `Vec64` for 64-byte SIMD aligned encoding, matching the diff --git a/rust/src/models/writers/stdio.rs b/rust/src/models/writers/stdio.rs index e4a1795..5b00183 100644 --- a/rust/src/models/writers/stdio.rs +++ b/rust/src/models/writers/stdio.rs @@ -8,7 +8,7 @@ //! //! High-level async writer that writes Arrow IPC encoded tables to stdout. //! -//! Wraps a [`TableSink64`](crate::models::sinks::table_sink::TableSink64) over `tokio::io::Stdout`, hiding the wiring +//! Wraps a [`TableSink64`] over `tokio::io::Stdout`, hiding the wiring //! so callers get a simple API for CLI tools. //! //! Uses `Vec64` for 64-byte SIMD aligned encoding, matching the diff --git a/rust/src/models/writers/uds.rs b/rust/src/models/writers/uds.rs index b8b20e2..0422e7a 100644 --- a/rust/src/models/writers/uds.rs +++ b/rust/src/models/writers/uds.rs @@ -9,7 +9,7 @@ //! High-level async writer that connects to a UDS endpoint and sends //! Arrow IPC encoded tables over the wire. //! -//! Wraps a [`TableSink64`](crate::models::sinks::table_sink::TableSink64) over a UDS write half, hiding the wiring +//! Wraps a [`TableSink64`] over a UDS write half, hiding the wiring //! so callers get a one-liner API. //! //! Uses `Vec64` for 64-byte SIMD aligned encoding, matching the diff --git a/rust/src/models/writers/websocket.rs b/rust/src/models/writers/websocket.rs index fcbe9e3..1d51ae6 100644 --- a/rust/src/models/writers/websocket.rs +++ b/rust/src/models/writers/websocket.rs @@ -10,7 +10,7 @@ //! Arrow IPC encoded tables as binary WebSocket messages. //! //! Extracts the raw TCP stream after the tungstenite handshake and uses -//! [`WsWrite`](crate::models::streams::websocket::WsWrite) for WebSocket binary frame encoding on the data path. +//! [`WsWrite`] for WebSocket binary frame encoding on the data path. //! //! Uses `Vec64` for 64-byte SIMD aligned encoding. //! @@ -22,7 +22,7 @@ //! that integration is compiled in. //! //! For pinned roots, a custom verifier, or client-auth keys, use -//! [`WebSocketTableWriter::connect_tls`](crate::models::writers::websocket::WebSocketTableWriter::connect_tls) - it takes an +//! [`WebSocketTableWriter::connect_tls`] - it takes an //! `Arc` directly and bypasses the bundled //! verifier. The library does not enforce a transport policy; if a //! deployment requires TLS, that is the caller's deployment decision. diff --git a/rust/src/models/writers/webtransport.rs b/rust/src/models/writers/webtransport.rs index 9f0b7a6..1a6a1f2 100644 --- a/rust/src/models/writers/webtransport.rs +++ b/rust/src/models/writers/webtransport.rs @@ -9,7 +9,7 @@ //! High-level async writer that sends Arrow IPC encoded tables over a //! WebTransport send stream. //! -//! Wraps a [`TableSink64`](crate::models::sinks::table_sink::TableSink64) over a [`wtransport::SendStream`], hiding the wiring +//! Wraps a [`TableSink64`] over a [`wtransport::SendStream`], hiding the wiring //! so callers get a one-liner API. //! //! Uses `Vec64` for 64-byte SIMD aligned encoding, matching the diff --git a/rust/src/traits/parallel_transport_reader.rs b/rust/src/traits/parallel_transport_reader.rs index 41f7dd6..dac0f09 100644 --- a/rust/src/traits/parallel_transport_reader.rs +++ b/rust/src/traits/parallel_transport_reader.rs @@ -8,7 +8,7 @@ //! //! Tables from all streams are exposed through a single merged stream. Order //! within a source stream is always preserved. Global write order across -//! streams is recovered with [`SortBehaviour::Ordered`](crate::traits::parallel_transport_reader::SortBehaviour::Ordered), which pulls the +//! streams is recovered with [`SortBehaviour::Ordered`], which pulls the //! streams in the writer's round-robin rotation. use std::future::Future; diff --git a/rust/src/traits/serialise.rs b/rust/src/traits/serialise.rs index 271bd0a..0066c2c 100644 --- a/rust/src/traits/serialise.rs +++ b/rust/src/traits/serialise.rs @@ -6,7 +6,7 @@ //! Codec-based serialisation and deserialisation. //! -//! [`Serialise`](crate::traits::serialise::Serialise) is implemented for Minarrow value types such as `Table`, +//! [`Serialise`] is implemented for Minarrow value types such as `Table`, //! `Array` and `FieldArray`. The codec type `C` defines the encoded format and //! provides the corresponding encoder and decoder implementations. //! diff --git a/rust/src/traits/stream_buffer.rs b/rust/src/traits/stream_buffer.rs index 571d9ae..db85ad5 100644 --- a/rust/src/traits/stream_buffer.rs +++ b/rust/src/traits/stream_buffer.rs @@ -181,7 +181,7 @@ impl StreamBuffer for Vec64 { #[inline] fn drain(&mut self, range: std::ops::Range) { - self.0.drain(range); + self.delete_range(range.start, range.end); } #[inline]