From f5ede5e683b341827d69637427bf3b1e3fefd4cb Mon Sep 17 00:00:00 2001 From: Davide Angelocola Date: Sun, 4 Oct 2026 08:24:30 +0200 Subject: [PATCH] feat(reader): read vortex.parquet.variant (#445) vortex-jni's default writer turns an Arrow arrow.parquet.variant column into vortex.variant over vortex.parquet.variant (core2026.08.3, Rust's default edition), which vortex-java could not decode. ParquetVariantEncodingDecoder follows Rust's ParquetVariant::deserialize: proto metadata {has_value, typed_value_dtype, value_nullable}, no buffers, children [validity?, metadata, value?, typed_value?]. It decodes to a StructArray {metadata, value?, typed_value?}, Arrow's own storage shape for the extension, masked by row validity; no new Array type. Rust dict-encodes the Binary value child, which exposed that DictEncodingDecoder only routed Utf8 through its VarBin path: a Binary dictionary hit a (DType.Primitive) cast and leaked ClassCastException. Binary now shares the Utf8 path, and any other non-primitive dtype throws VortexException. Writing vortex.parquet.variant stays out of scope (ADR 0014). Co-Authored-By: Claude Sonnet 5 --- CHANGELOG.md | 4 + .../dfa1/vortex/core/model/Editions.java | 10 +- .../dfa1/vortex/core/model/EncodingId.java | 5 + .../proto/ProtoParquetVariantMetadata.java | 73 ++++++ core/src/main/proto/encodings.proto | 6 + .../dfa1/vortex/core/model/EditionsTest.java | 2 +- docs/compatibility.md | 7 +- docs/reference.md | 6 +- .../ParquetVariantInteropIntegrationTest.java | 240 ++++++++++++++++++ .../dfa1/vortex/reader/ReadRegistry.java | 2 + .../reader/decode/DictEncodingDecoder.java | 15 +- .../decode/ParquetVariantEncodingDecoder.java | 101 ++++++++ .../decode/DictEncodingDecoderTest.java | 27 +- .../ParquetVariantEncodingDecoderTest.java | 225 ++++++++++++++++ 14 files changed, 705 insertions(+), 18 deletions(-) create mode 100644 core/src/main/java/io/github/dfa1/vortex/core/proto/ProtoParquetVariantMetadata.java create mode 100644 integration/src/test/java/io/github/dfa1/vortex/integration/ParquetVariantInteropIntegrationTest.java create mode 100644 reader/src/main/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoder.java create mode 100644 reader/src/test/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoderTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index e884dc29..a073fe3e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,7 +13,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `fastlanes.delta` and `vortex.patched` belong to no edition, as in Rust: emit them with the new `WriteOptions.withoutEditions()`, the counterpart of Rust's `disable_editions()` ([#441](https://github.com/dfa1/vortex-java/issues/441)). - Nullable low-cardinality string columns dict-encoded per chunk store null as a dictionary entry, as Rust does, instead of adding a row validity bitmap: about 20% smaller, now slightly below vortex-jni ([2091beb](https://github.com/dfa1/vortex-java/commit/2091beb3)). +### Added +- Read `vortex.parquet.variant`, the encoding vortex-jni's default writer uses for Arrow `arrow.parquet.variant` columns: each row's Apache Variant `metadata`/`value` binaries come back as a struct ([#445](https://github.com/dfa1/vortex-java/issues/445)). + ### Fixed +- Dictionary-encoded Binary columns (e.g. the `value` child of `vortex.parquet.variant`) failed to read with a `ClassCastException` ([#445](https://github.com/dfa1/vortex-java/issues/445)). - Null rows of Rust-written `vortex.list` columns read as empty lists: the decoder ignored the list's validity child. - CSV export (`CsvExporter`, `vortex export`) failed on list-view columns with `unsupported array type for CSV export: ListViewArray`. - `vortex inspect --html` grouped the Chunks panel by chunk index, so files whose columns chunk differently (e.g. Rust's `tpch_orders.compact`) showed overlapping row ranges and mixed sizes; the panel now lists one entry per distinct row range. diff --git a/core/src/main/java/io/github/dfa1/vortex/core/model/Editions.java b/core/src/main/java/io/github/dfa1/vortex/core/model/Editions.java index 2c9425f6..3b4ef9a6 100644 --- a/core/src/main/java/io/github/dfa1/vortex/core/model/Editions.java +++ b/core/src/main/java/io/github/dfa1/vortex/core/model/Editions.java @@ -23,10 +23,10 @@ /// zone-map aggregates) has no members here. Encodings in no edition at all — `fastlanes.delta`, /// `vortex.patched` — are writable only with the guard turned off, as in Rust. /// -/// vortex-java implements every `core`-family encoding except `vortex.parquet.variant`, and not -/// `zstd`'s `vortex.zstd_buffers`; both have no [EncodingId.WellKnown] constant yet and are named -/// as [EncodingId.Custom] instead: the -/// catalog mirrors upstream faithfully rather than being truncated to what is implemented today. +/// vortex-java implements every `core`-family encoding (`vortex.parquet.variant` read only), but not +/// `zstd`'s `vortex.zstd_buffers`, which has no [EncodingId.WellKnown] constant yet and is named as +/// an [EncodingId.Custom] instead: the catalog mirrors upstream faithfully rather than being +/// truncated to what is implemented today. public final class Editions { /// The baseline `core` edition: stable encodings writable by Vortex (Rust reference) 0.36.0. @@ -75,7 +75,7 @@ public final class Editions { /// `core` edition, and the one the default writer targets, as in Rust. public static final Edition CORE_2026_08_3 = new Edition( new EditionId(EditionFamily.CORE, YearMonth.of(2026, 8), 3), - Set.of(new EncodingId.Custom("vortex.parquet.variant"), EncodingId.VORTEX_VARIANT)); + Set.of(EncodingId.VORTEX_PARQUET_VARIANT, EncodingId.VORTEX_VARIANT)); /// The August 2026 draft edition of the `preview` family. Empty in Rust too: no component has /// entered preview yet. diff --git a/core/src/main/java/io/github/dfa1/vortex/core/model/EncodingId.java b/core/src/main/java/io/github/dfa1/vortex/core/model/EncodingId.java index 81a4313f..8baf60c2 100644 --- a/core/src/main/java/io/github/dfa1/vortex/core/model/EncodingId.java +++ b/core/src/main/java/io/github/dfa1/vortex/core/model/EncodingId.java @@ -112,6 +112,9 @@ enum WellKnown implements EncodingId { VORTEX_PATCHED("vortex.patched"), /// Variant logical encoding: canonical container over `core_storage` plus an optional shredded child. VORTEX_VARIANT("vortex.variant"), + /// Parquet Variant physical encoding (`vortex.parquet.variant`): per-row Apache Variant binary + /// `metadata` and `value` children, plus an optional shredded `typed_value` child. + VORTEX_PARQUET_VARIANT("vortex.parquet.variant"), ; // O(1) access to a WellKnown constant by its string representation @@ -253,4 +256,6 @@ public String toString() { WellKnown VORTEX_PATCHED = WellKnown.VORTEX_PATCHED; /// Well-known `vortex.variant` id. WellKnown VORTEX_VARIANT = WellKnown.VORTEX_VARIANT; + /// Well-known `vortex.parquet.variant` id. + WellKnown VORTEX_PARQUET_VARIANT = WellKnown.VORTEX_PARQUET_VARIANT; } diff --git a/core/src/main/java/io/github/dfa1/vortex/core/proto/ProtoParquetVariantMetadata.java b/core/src/main/java/io/github/dfa1/vortex/core/proto/ProtoParquetVariantMetadata.java new file mode 100644 index 00000000..0b579f31 --- /dev/null +++ b/core/src/main/java/io/github/dfa1/vortex/core/proto/ProtoParquetVariantMetadata.java @@ -0,0 +1,73 @@ +package io.github.dfa1.vortex.core.proto; + +import java.io.IOException; +import java.lang.foreign.MemorySegment; +import javax.annotation.processing.Generated; + +/// Generated from proto3 message {@code vortex.encodings.ParquetVariantMetadata}. +/// Do not edit by hand — regenerate via {@code ./mvnw generate-sources -pl core -P regenerate-sources}. +/// @param has_value field tag 1 +/// @param typed_value_dtype field tag 2 +/// @param value_nullable field tag 3 +@Generated("io.github.dfa1.vortex.protogen.CodeGen") +public record ProtoParquetVariantMetadata( + boolean has_value, + ProtoDType typed_value_dtype, + boolean value_nullable +) { + + /// Decodes a {@code vortex.encodings.ParquetVariantMetadata} from a slice of a memory segment. + /// @param __seg backing segment + /// @param __off start offset in bytes + /// @param __len payload length in bytes + /// @return decoded record + /// @throws IOException if the slice is malformed or truncated + public static ProtoParquetVariantMetadata decode(MemorySegment __seg, long __off, long __len) throws IOException { + ProtoReader r = new ProtoReader(__seg, __off, __len); + boolean has_value = false; + ProtoDType typed_value_dtype = null; + boolean value_nullable = false; + while (r.hasMore()) { + int tag = r.readVarint32(); + switch (tag >>> 3) { + case 1 -> { + has_value = r.readBool(); + } + case 2 -> { + MemorySegment __slice = r.readLenDelimSegment(); + typed_value_dtype = ProtoDType.decode(__slice, 0, __slice.byteSize()); + } + case 3 -> { + value_nullable = r.readBool(); + } + default -> r.skipField(tag & 7); + } + } + return new ProtoParquetVariantMetadata(has_value, typed_value_dtype, value_nullable); + } + + /// Encodes this record to a proto3-wire-format byte array. + /// @return encoded bytes + public byte[] encode() { + ProtoWriter w = new ProtoWriter(); + encodeTo(w); + return w.toByteArray(); + } + + void encodeTo(ProtoWriter w) { + if (has_value) { + w.writeTag(1, 0); + w.writeBool(has_value); + } + if (typed_value_dtype != null) { + w.writeTag(2, 2); + int __mark = w.beginLenDelim(); + typed_value_dtype.encodeTo(w); + w.endLenDelim(__mark); + } + if (value_nullable) { + w.writeTag(3, 0); + w.writeBool(value_nullable); + } + } +} diff --git a/core/src/main/proto/encodings.proto b/core/src/main/proto/encodings.proto index 88a46333..a9feb647 100644 --- a/core/src/main/proto/encodings.proto +++ b/core/src/main/proto/encodings.proto @@ -148,6 +148,12 @@ message VariantMetadata { optional vortex.dtype.DType shredded_dtype = 1; } +message ParquetVariantMetadata { + bool has_value = 1; + optional vortex.dtype.DType typed_value_dtype = 2; + bool value_nullable = 3; +} + message OnPairMetadata { vortex.dtype.PType uncompressed_lengths_ptype = 1; uint32 dict_size = 3; diff --git a/core/src/test/java/io/github/dfa1/vortex/core/model/EditionsTest.java b/core/src/test/java/io/github/dfa1/vortex/core/model/EditionsTest.java index 8a5c2aa1..a41713fb 100644 --- a/core/src/test/java/io/github/dfa1/vortex/core/model/EditionsTest.java +++ b/core/src/test/java/io/github/dfa1/vortex/core/model/EditionsTest.java @@ -96,7 +96,7 @@ void core2026_08_3_isTheFullCoreSet() { EncodingId.FASTLANES_RLE, EncodingId.VORTEX_FIXED_SIZE_LIST, EncodingId.VORTEX_LISTVIEW, EncodingId.VORTEX_MASKED, EncodingId.VORTEX_ONPAIR, EncodingId.VORTEX_MAP, - EncodingId.VORTEX_VARIANT, new EncodingId.Custom("vortex.parquet.variant"))); + EncodingId.VORTEX_VARIANT, EncodingId.VORTEX_PARQUET_VARIANT)); } @Test diff --git a/docs/compatibility.md b/docs/compatibility.md index 4fc937a3..c782c955 100644 --- a/docs/compatibility.md +++ b/docs/compatibility.md @@ -38,7 +38,7 @@ only the built-in decoders in `reader`; no encoder class is loaded. |------|------------|-------------| | `DType::Union` (`fbs.DType.Type.Union = 12`) | Rust 0.71.0 | ❌ Decode throws `VortexException("unsupported DType typeType=12")`. No `DType.Union` variant in Java's sealed type. | | `vortex.onpair` experimental string encoding | Rust 0.74.0 | ✅ Read and written. In `core2026.08.1`, so default cascading writes offer it, as Rust's default compressor does; the trained dictionary is valid for Rust but not byte-identical to Rust's (Rust's sampling RNG is not portable). | -| `vortex.variant` arbitrary nested objects | Rust (`vortex.parquet.variant`) | ⚠️ Java encodes/decodes variant columns of **typed scalar** values (constant / chunked-of-constants core, optional shredded child); Java↔Rust round-trip verified. Arbitrary nested JSON objects and real path-based shredding need the `vortex.parquet.variant` physical encoding — deferred ([ADR 0014](../adr/0014-variant-encoding-strategy.md)). ❌ Reading is a real gap: vortex-jni's default writer turns an Arrow `arrow.parquet.variant` column into `vortex.parquet.variant` (`core2026.08.3`), which vortex-java has no decoder for. | +| `vortex.variant` arbitrary nested objects | Rust (`vortex.parquet.variant`) | ⚠️ Java reads `vortex.parquet.variant`, the Apache Variant binary encoding vortex-jni's default writer uses for Arrow `arrow.parquet.variant` columns, as a struct of per-row `metadata`/`value` binaries; it does not interpret the Variant binary itself. Java writes variant columns of **typed scalar** values only (constant / chunked-of-constants core, optional shredded child), via `vortex.variant`; writing `vortex.parquet.variant` is not implemented ([ADR 0014](../adr/0014-variant-encoding-strategy.md)). | | Arrow extension array import affecting Variant shape | Rust 0.74.0 (#8125) | Untested against the currently pinned v0.85.0 fixtures; #8125 not yet re-verified. | | `vortex.dict` **layout** over a values pool that is neither VarBin- nor primitive-shaped (e.g. a dict-encoded `vortex.uuid`, whose storage is `FixedSizeList(U8, 16)`) | Not written by Rust: its dict layout admits only `Primitive \| Utf8 \| Binary` (`dict_layout_supported`), and its dict compressor schemes only integers, floats and strings | ⚠️ Unreachable from Rust- or Java-written files. Every type vortex-jni writes reads back exactly, dict-encoded or not (`DictAllTypesInteropIntegrationTest`, which also fails if a Rust bump starts dict-encoding another type). A foreign file with such a pool fails with `VortexException("unsupported dict values shape: …")`. | | Duplicate struct field names | Rust writer rejects ("StructLayout must have unique field names"); Rust reader tolerates foreign files (first-match access) | ⚠️ Deliberate divergence on read: Java rejects such files with `VortexException("duplicate field name in file schema")` instead of tolerating them — the name-keyed `Chunk` API cannot represent both columns, and silent column loss is worse than a loud failure on a file the reference writer refuses to produce. Java's writer mirrors the Rust writer's rejection. | @@ -122,6 +122,7 @@ decimals ([#430](https://github.com/dfa1/vortex-java/pull/430)) and nulls in nul | `fastlanes.rle` | `RleEncodingDecoder` | `RleEncodingEncoder` | ✅ | ✅ | Chunk-based RLE. Integers and floats (Rust's int and float RLE schemes); float runs compare raw bits, so -0.0 and NaN payloads round-trip. Cascades values/indices/offsets | | `vortex.patched` | `PatchedEncodingDecoder` | `PatchedEncodingEncoder` | ✅ | ✅ | Primitive PTypes; base + chunked patches (1024-elem blocks) | | `vortex.variant` | `VariantEncodingDecoder` | `VariantEncodingEncoder` | ✅ | ✅ | Canonical container; constant / chunked-of-constants core + optional shredded child. Typed-scalar values only — nested objects need `parquet.variant` (ADR 0014) | +| `vortex.parquet.variant` | `ParquetVariantEncodingDecoder` | — | ✅ | ❌ | Read only: decodes to a struct `{metadata, value?, typed_value?}` of Apache Variant binaries (Arrow's `arrow.parquet.variant` storage shape), masked by row validity. Writing is not implemented; Java writes typed-scalar variants through `vortex.variant` (ADR 0014) | | `vortex.onpair` | `OnPairEncodingDecoder` | `OnPairEncodingEncoder` | ✅ | ✅ | Utf8, Binary; `core2026.08.1`, a cascade candidate under the default edition (competes with FSST, as in Rust) | ### Decode shape @@ -215,7 +216,7 @@ guarantee once frozen (ADR 0023) — a write-time/read-time policy, not part of | `core2026.08.0` | `core` | no encoding (Rust adds the `vortex.zoned` layout and zone-map aggregates, not modelled) | | `core2026.08.1` | `core` | `vortex.onpair` | | `core2026.08.2` | `core` | `vortex.map` | -| `core2026.08.3` | `core` | `vortex.variant`, `vortex.parquet.variant` ❌ not implemented — **default write target** | +| `core2026.08.3` | `core` | `vortex.variant`, `vortex.parquet.variant` (read only) — **default write target** | | `preview2026.08.0` | `preview` | nothing yet | | `zstd2026.02.0` | `zstd` | `vortex.zstd_buffers` ❌ not implemented (buffer-level Zstd for GPU decode; declared by Rust's `vortex-zstd` plugin, opt-in) | @@ -345,5 +346,5 @@ Cross-language round-trips tested against Rust-written fixture files hosted at | `clickbench_hits_5k.regular.vortex` | ✅ | Full scan of every Utf8 column (`vortex.onpair`) compared against vortex-jni | | `masked.vortex` | ❓ | No fixture through v0.86.1, and vortex-jni's writer folds validity into each encoding rather than emitting `vortex.masked` for nullable primitive, string, struct or list input | | `patched.vortex` | ❓ | No fixture through v0.86.1; `vortex.patched` is in no edition, so vortex-jni's default writer never emits it | -| `variant.vortex` | ❌ | No fixture through v0.86.1, but vortex-jni writes an Arrow `arrow.parquet.variant` column as `vortex.variant` over `vortex.parquet.variant`, which vortex-java cannot decode | +| `variant.vortex` | ✅ | No fixture through v0.86.1; covered instead by `ParquetVariantInteropIntegrationTest`, which has vortex-jni write an Arrow `arrow.parquet.variant` column (`vortex.variant` over `vortex.parquet.variant`, plain and shredded) | | `map.vortex` | ✅ | New in v0.86.1; also covered directly (both directions, nullable) by the vortex-jni oracle (issue #351) | diff --git a/docs/reference.md b/docs/reference.md index 87e4061f..c7950cbe 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -347,9 +347,9 @@ targeted edition is ever persisted into a `.vortex` file. | `cumulativeMembers(Edition)` | The edition's own additions plus every earlier same-family edition's | | `owningEdition(EncodingId)` | The edition an id first joined, or empty if it belongs to none | -vortex-java implements every `core`-family encoding except `vortex.parquet.variant`, and not `zstd`'s -`vortex.zstd_buffers`; both resolve to `EncodingId.Custom` and are stored in the catalog anyway, -mirroring upstream faithfully. +vortex-java implements every `core`-family encoding (`vortex.parquet.variant` read only), but not +`zstd`'s `vortex.zstd_buffers`, which resolves to `EncodingId.Custom` and is stored in the catalog +anyway, mirroring upstream faithfully. `fastlanes.delta` and `vortex.patched` belong to no edition, as in Rust. ### Writer integration (`WriteOptions#editions()`) diff --git a/integration/src/test/java/io/github/dfa1/vortex/integration/ParquetVariantInteropIntegrationTest.java b/integration/src/test/java/io/github/dfa1/vortex/integration/ParquetVariantInteropIntegrationTest.java new file mode 100644 index 00000000..66555285 --- /dev/null +++ b/integration/src/test/java/io/github/dfa1/vortex/integration/ParquetVariantInteropIntegrationTest.java @@ -0,0 +1,240 @@ +package io.github.dfa1.vortex.integration; + +import dev.vortex.api.Session; +import dev.vortex.api.VortexWriter; +import dev.vortex.arrow.ArrowAllocation; +import dev.vortex.jni.NativeLoader; +import io.github.dfa1.vortex.inspect.InspectorTree; +import io.github.dfa1.vortex.reader.ReadRegistry; +import io.github.dfa1.vortex.reader.ScanOptions; +import io.github.dfa1.vortex.reader.VortexReader; +import io.github.dfa1.vortex.reader.array.Array; +import io.github.dfa1.vortex.reader.array.IntArray; +import io.github.dfa1.vortex.reader.array.MaskedArray; +import io.github.dfa1.vortex.reader.array.StructArray; +import io.github.dfa1.vortex.reader.array.VarBinArray; +import io.github.dfa1.vortex.reader.array.VariantArray; +import org.apache.arrow.c.ArrowArray; +import org.apache.arrow.c.ArrowSchema; +import org.apache.arrow.c.Data; +import org.apache.arrow.memory.BufferAllocator; +import org.apache.arrow.vector.IntVector; +import org.apache.arrow.vector.VarBinaryVector; +import org.apache.arrow.vector.VectorSchemaRoot; +import org.apache.arrow.vector.complex.StructVector; +import org.apache.arrow.vector.types.pojo.ArrowType; +import org.apache.arrow.vector.types.pojo.Field; +import org.apache.arrow.vector.types.pojo.FieldType; +import org.apache.arrow.vector.types.pojo.Schema; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +/// vortex-jni's default writer turns an Arrow `arrow.parquet.variant` column into +/// `vortex.variant` over `vortex.parquet.variant` (`core2026.08.3`, Rust's default edition); +/// vortex-java must read every row's Variant binaries back (#445). +/// +/// Rows mix an int8 and a short-string Variant so a reader that assumed one value type per +/// column would fail, and every 7th row is null so the row validity child is exercised. +class ParquetVariantInteropIntegrationTest { + + private static final Session SESSION = Session.create(); + private static final BufferAllocator ALLOCATOR = ArrowAllocation.rootAllocator(); + private static final int ROWS = 1_000; + /// Variant metadata with an empty dictionary: version 1, no keys. + private static final byte[] EMPTY_METADATA = {0x01, 0x00}; + + static { + NativeLoader.loadJni(); + } + + @Test + void jniWritesParquetVariant_javaReadsEveryRow(@TempDir Path tmp) throws IOException { + // Given + Path file = tmp.resolve("parquet_variant.vortex"); + writeJni(file); + + // When + List result = readValues(file); + + // Then — vortex-jni chose the encoding under test, and every row's binary survives + try (var vf = VortexReader.open(file, ReadRegistry.loadAll())) { + assertThat(InspectorTree.build(vf).usedEncodings()).contains("vortex.parquet.variant"); + } + assertThat(result).hasSize(ROWS); + for (int i = 0; i < ROWS; i++) { + assertThat(result.get(i)).as("row %d", i).isEqualTo(expectedValue(i)); + } + } + + @Test + void jniWritesShreddedParquetVariant_javaReadsValueAndTypedValue(@TempDir Path tmp) throws IOException { + // Given — shredded storage: even rows carry the int in `typed_value` with a null `value`, + // odd rows the reverse, so the nullable `value` child must decode + Path file = tmp.resolve("parquet_variant_shredded.vortex"); + writeShreddedJni(file); + + // When + List result = readShredded(file); + + // Then + assertThat(result).hasSize(ROWS); + for (int i = 0; i < ROWS; i++) { + assertThat(result.get(i)).as("row %d", i).isEqualTo(i % 2 == 0 ? "typed:" + i : "value:" + i); + } + } + + /// Apache Variant value for row `i`, or `null` for a null row: an int8 primitive (header + /// `0x0c` = type 3 << 2 | basic type 0) on even rows, a short string (header = length << 2 | + /// basic type 1) on odd rows. + private static byte[] expectedValue(int i) { + if (i % 7 == 0) { + return null; + } + return int8OrString(i); + } + + private static byte[] int8OrString(int i) { + if (i % 2 == 0) { + return new byte[]{0x0c, (byte) i}; + } + byte[] s = ("s" + i % 10).getBytes(StandardCharsets.UTF_8); + byte[] v = new byte[s.length + 1]; + v[0] = (byte) (s.length << 2 | 1); + System.arraycopy(s, 0, v, 1, s.length); + return v; + } + + private static void writeJni(Path file) throws IOException { + Field metadata = Field.notNullable("metadata", ArrowType.Binary.INSTANCE); + Field value = Field.notNullable("value", ArrowType.Binary.INSTANCE); + FieldType variantType = new FieldType(true, ArrowType.Struct.INSTANCE, null, + Map.of("ARROW:extension:name", "arrow.parquet.variant", "ARROW:extension:metadata", "")); + Schema schema = new Schema(List.of(new Field("v", variantType, List.of(metadata, value)))); + try (VortexWriter writer = VortexWriter.builder(SESSION, file.toUri().toString(), schema, ALLOCATOR).build(); + VectorSchemaRoot root = VectorSchemaRoot.create(schema, ALLOCATOR)) { + var struct = (StructVector) root.getVector("v"); + struct.allocateNew(); + var metadataVector = (VarBinaryVector) struct.getChild("metadata"); + var valueVector = (VarBinaryVector) struct.getChild("value"); + for (int i = 0; i < ROWS; i++) { + byte[] v = expectedValue(i); + // Non-null children under a null struct row still need a value; any valid Variant does. + metadataVector.setSafe(i, EMPTY_METADATA); + valueVector.setSafe(i, v == null ? new byte[]{0x00} : v); + if (v == null) { + struct.setNull(i); + } else { + struct.setIndexDefined(i); + } + } + metadataVector.setValueCount(ROWS); + valueVector.setValueCount(ROWS); + struct.setValueCount(ROWS); + root.setRowCount(ROWS); + try (ArrowArray arr = ArrowArray.allocateNew(ALLOCATOR); + ArrowSchema arrowSchema = ArrowSchema.allocateNew(ALLOCATOR)) { + Data.exportVectorSchemaRoot(ALLOCATOR, root, null, arr, arrowSchema); + writer.writeBatch(arr.memoryAddress(), arrowSchema.memoryAddress()); + } + } + } + + /// Reads each row's Variant `value` binary (asserting its `metadata` on the way), or `null` + /// for a null row. + private static List readValues(Path file) throws IOException { + List rows = new ArrayList<>(ROWS); + try (var reader = VortexReader.open(file, ReadRegistry.loadAll()); + var iter = reader.scan(ScanOptions.columns("v"))) { + while (iter.hasNext()) { + Array column = iter.next().column("v"); + Array storage = ((VariantArray) column).coreStorage(); + MaskedArray masked = storage instanceof MaskedArray m ? m : null; + StructArray struct = (StructArray) (masked != null ? masked.inner() : storage); + VarBinArray metadata = (VarBinArray) struct.field(0); + VarBinArray value = (VarBinArray) struct.field(1); + for (long i = 0; i < struct.length(); i++) { + if (masked != null && !masked.isValid(i)) { + rows.add(null); + continue; + } + assertThat(metadata.getBytes(i)).isEqualTo(EMPTY_METADATA); + rows.add(value.getBytes(i)); + } + } + } + return rows; + } + private static void writeShreddedJni(Path file) throws IOException { + Field metadata = Field.notNullable("metadata", ArrowType.Binary.INSTANCE); + Field value = Field.nullable("value", ArrowType.Binary.INSTANCE); + Field typedValue = Field.nullable("typed_value", new ArrowType.Int(32, true)); + FieldType variantType = new FieldType(false, ArrowType.Struct.INSTANCE, null, + Map.of("ARROW:extension:name", "arrow.parquet.variant", "ARROW:extension:metadata", "")); + Schema schema = new Schema(List.of(new Field("v", variantType, List.of(metadata, value, typedValue)))); + try (VortexWriter writer = VortexWriter.builder(SESSION, file.toUri().toString(), schema, ALLOCATOR).build(); + VectorSchemaRoot root = VectorSchemaRoot.create(schema, ALLOCATOR)) { + var struct = (StructVector) root.getVector("v"); + struct.allocateNew(); + var metadataVector = (VarBinaryVector) struct.getChild("metadata"); + var valueVector = (VarBinaryVector) struct.getChild("value"); + var typedVector = (IntVector) struct.getChild("typed_value"); + for (int i = 0; i < ROWS; i++) { + metadataVector.setSafe(i, EMPTY_METADATA); + if (i % 2 == 0) { + valueVector.setNull(i); + typedVector.setSafe(i, i); + } else { + valueVector.setSafe(i, int8OrString(i)); + typedVector.setNull(i); + } + struct.setIndexDefined(i); + } + metadataVector.setValueCount(ROWS); + valueVector.setValueCount(ROWS); + typedVector.setValueCount(ROWS); + struct.setValueCount(ROWS); + root.setRowCount(ROWS); + try (ArrowArray arr = ArrowArray.allocateNew(ALLOCATOR); + ArrowSchema arrowSchema = ArrowSchema.allocateNew(ALLOCATOR)) { + Data.exportVectorSchemaRoot(ALLOCATOR, root, null, arr, arrowSchema); + writer.writeBatch(arr.memoryAddress(), arrowSchema.memoryAddress()); + } + } + } + + /// Renders each row as `typed:` when the shredded child holds it, else `value:` after + /// checking the `value` binary is the one written for row `i`. Rust's writer moves the Arrow + /// `typed_value` out of `vortex.parquet.variant` into the canonical `vortex.variant` + /// container's shredded child, so that is where the ints come back from. + private static List readShredded(Path file) throws IOException { + List rows = new ArrayList<>(ROWS); + try (var reader = VortexReader.open(file, ReadRegistry.loadAll()); + var iter = reader.scan(ScanOptions.columns("v"))) { + while (iter.hasNext()) { + VariantArray column = iter.next().column("v"); + StructArray struct = (StructArray) column.coreStorage(); + VarBinArray value = (VarBinArray) struct.field(1); + MaskedArray typed = (MaskedArray) column.shredded(); + for (long i = 0; i < struct.length(); i++) { + int row = rows.size(); + if (typed.isValid(i)) { + rows.add("typed:" + ((IntArray) typed.inner()).getInt(i)); + } else { + assertThat(value.getBytes(i)).isEqualTo(int8OrString(row)); + rows.add("value:" + row); + } + } + } + } + return rows; + } +} diff --git a/reader/src/main/java/io/github/dfa1/vortex/reader/ReadRegistry.java b/reader/src/main/java/io/github/dfa1/vortex/reader/ReadRegistry.java index 7aedfe52..5fccf0bf 100644 --- a/reader/src/main/java/io/github/dfa1/vortex/reader/ReadRegistry.java +++ b/reader/src/main/java/io/github/dfa1/vortex/reader/ReadRegistry.java @@ -31,6 +31,7 @@ import io.github.dfa1.vortex.reader.decode.MaskedEncodingDecoder; import io.github.dfa1.vortex.reader.decode.NullEncodingDecoder; import io.github.dfa1.vortex.reader.decode.OnPairEncodingDecoder; +import io.github.dfa1.vortex.reader.decode.ParquetVariantEncodingDecoder; import io.github.dfa1.vortex.reader.decode.PatchedEncodingDecoder; import io.github.dfa1.vortex.reader.decode.PcoEncodingDecoder; import io.github.dfa1.vortex.reader.decode.PrimitiveEncodingDecoder; @@ -231,6 +232,7 @@ public Builder registerDefaults() { .register(new MaskedEncodingDecoder()) .register(new NullEncodingDecoder()) .register(new OnPairEncodingDecoder()) + .register(new ParquetVariantEncodingDecoder()) .register(new PatchedEncodingDecoder()) .register(new PcoEncodingDecoder()) .register(new PrimitiveEncodingDecoder()) diff --git a/reader/src/main/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoder.java b/reader/src/main/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoder.java index 84bfd18d..d237898d 100644 --- a/reader/src/main/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoder.java +++ b/reader/src/main/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoder.java @@ -53,7 +53,9 @@ public EncodingId encodingId() { public Array decode(DecodeContext ctx) { MemorySegment meta = ctx.metadata(); - if (ctx.dtype() instanceof DType.Utf8) { + // Binary dictionaries share the Utf8 shape (VarBin values + codes); Rust dict-encodes + // e.g. the `value` child of `vortex.parquet.variant` this way. + if (ctx.dtype() instanceof DType.Utf8 || ctx.dtype() instanceof DType.Binary) { if (ctx.node().children().length == 0) { if (meta == null || meta.byteSize() == 0) { throw new VortexException(EncodingId.VORTEX_DICT, "missing metadata for legacy utf8 dict"); @@ -78,7 +80,7 @@ public Array decode(DecodeContext ctx) { private static Array decodeLegacyJava(DecodeContext ctx, byte codeTypeByte) { PType codePType = PType.fromOrdinal(Byte.toUnsignedInt(codeTypeByte)); - PType valPType = ((DType.Primitive) ctx.dtype()).ptype(); + PType valPType = valuePType(ctx.dtype()); long rowCount = ctx.rowCount(); requireUnsignedCodePType(codePType); @@ -108,7 +110,7 @@ private static Array decodeRustProto(DecodeContext ctx, MemorySegment metaBuf) { PType codePType = PType.fromOrdinal(meta.codes_ptype().value()); long valuesLen = meta.values_len(); long rowCount = ctx.rowCount(); - PType valPType = ((DType.Primitive) ctx.dtype()).ptype(); + PType valPType = valuePType(ctx.dtype()); requireUnsignedCodePType(codePType); // Row validity mirrors the Rust reference: a DictArray row is null when its CODE @@ -149,6 +151,13 @@ private static Array decodeRustProto(DecodeContext ctx, MemorySegment metaBuf) { /// @param values dictionary pool, already mask-unwrapped /// @param codes per-row codes, already mask-unwrapped /// @return the lazy dict array + private static PType valuePType(DType dtype) { + if (!(dtype instanceof DType.Primitive p)) { + throw new VortexException(EncodingId.VORTEX_DICT, "unsupported dict dtype: " + dtype); + } + return p.ptype(); + } + private static Array buildLazyDict(DType dtype, PType valPType, long n, Array values, Array codes) { try { return switch (valPType) { diff --git a/reader/src/main/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoder.java b/reader/src/main/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoder.java new file mode 100644 index 00000000..8639f28a --- /dev/null +++ b/reader/src/main/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoder.java @@ -0,0 +1,101 @@ +package io.github.dfa1.vortex.reader.decode; + +import io.github.dfa1.vortex.core.error.VortexException; +import io.github.dfa1.vortex.core.model.ColumnName; +import io.github.dfa1.vortex.core.model.DType; +import io.github.dfa1.vortex.core.model.EncodingId; +import io.github.dfa1.vortex.core.proto.ProtoParquetVariantMetadata; +import io.github.dfa1.vortex.reader.array.Array; +import io.github.dfa1.vortex.reader.array.MaskedArray; +import io.github.dfa1.vortex.reader.array.StructArray; + +import java.io.IOException; +import java.lang.foreign.MemorySegment; +import java.util.ArrayList; +import java.util.List; + +/// Read-only decoder for `vortex.parquet.variant`: per-row [Apache Variant](https://github.com/apache/parquet-format/blob/master/VariantEncoding.md) +/// binaries, as Rust's `encodings/parquet-variant` `deserialize` lays them out. No buffers; +/// children `[validity?, metadata, value?, typed_value?]`, with `value` / `typed_value` presence +/// and their dtypes declared in [ProtoParquetVariantMetadata]. +/// +/// Decodes to a [StructArray] `{metadata, value?, typed_value?}`: the shape Arrow's own +/// `arrow.parquet.variant` storage uses, so callers read the Variant binaries per field. Row +/// validity, when present, wraps the struct in a [MaskedArray]. +public final class ParquetVariantEncodingDecoder implements EncodingDecoder { + + private static final DType BINARY = new DType.Binary(false); + + @Override + public EncodingId encodingId() { + return EncodingId.VORTEX_PARQUET_VARIANT; + } + + @Override + public Array decode(DecodeContext ctx) { + if (!(ctx.dtype() instanceof DType.Variant)) { + throw new VortexException(EncodingId.VORTEX_PARQUET_VARIANT, + "expected variant dtype, got " + ctx.dtype()); + } + ProtoParquetVariantMetadata meta = parseMetadata(ctx.metadata()); + DType typedValueDtype = meta.typed_value_dtype() == null + ? null : VariantEncodingDecoder.dtypeFromProto(meta.typed_value_dtype()); + if (!meta.has_value() && typedValueDtype == null) { + throw new VortexException(EncodingId.VORTEX_PARQUET_VARIANT, + "at least one of value or typed_value must be present"); + } + if (ctx.node().bufferIndices().length != 0) { + throw new VortexException(EncodingId.VORTEX_PARQUET_VARIANT, + "expected 0 buffers, got " + ctx.node().bufferIndices().length); + } + + int expected = 1 + (meta.has_value() ? 1 : 0) + (typedValueDtype != null ? 1 : 0); + int numChildren = ctx.node().children().length; + if (numChildren != expected && numChildren != expected + 1) { + throw new VortexException(EncodingId.VORTEX_PARQUET_VARIANT, + "expected " + expected + " or " + (expected + 1) + " children, got " + numChildren); + } + + long n = ctx.rowCount(); + int child = 0; + Array validity = numChildren == expected ? null : ctx.decodeChild(child++, DType.BOOL, n); + + List names = new ArrayList<>(3); + List types = new ArrayList<>(3); + List fields = new ArrayList<>(3); + names.add(ColumnName.of("metadata")); + types.add(BINARY); + fields.add(ctx.decodeChild(child++, BINARY, n)); + if (meta.has_value()) { + DType valueDtype = new DType.Binary(meta.value_nullable()); + names.add(ColumnName.of("value")); + types.add(valueDtype); + fields.add(ctx.decodeChild(child++, valueDtype, n)); + } + if (typedValueDtype != null) { + names.add(ColumnName.of("typed_value")); + types.add(typedValueDtype); + fields.add(ctx.decodeChild(child, typedValueDtype, n)); + } + + StructArray struct = new StructArray(new DType.Struct(names, types, false), n, fields); + if (validity == null) { + return struct; + } + return new MaskedArray(struct, + MaskedArray.requireBoolArray(validity, EncodingId.VORTEX_PARQUET_VARIANT, "validity child")); + } + + private static ProtoParquetVariantMetadata parseMetadata(MemorySegment rawMeta) { + if (rawMeta == null || rawMeta.byteSize() == 0) { + // proto3: an all-default message encodes to zero bytes (has_value=false, no typed_value), + // which the presence check above then rejects, as Rust does. + return new ProtoParquetVariantMetadata(false, null, false); + } + try { + return ProtoParquetVariantMetadata.decode(rawMeta, 0, rawMeta.byteSize()); + } catch (IOException e) { + throw new VortexException(EncodingId.VORTEX_PARQUET_VARIANT, "invalid metadata", e); + } + } +} diff --git a/reader/src/test/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoderTest.java b/reader/src/test/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoderTest.java index c2f03336..f768fbca 100644 --- a/reader/src/test/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoderTest.java +++ b/reader/src/test/java/io/github/dfa1/vortex/reader/decode/DictEncodingDecoderTest.java @@ -160,6 +160,19 @@ void unexpectedCodeType_throws() { .hasMessageContaining("unexpected code type"); } + @Test + void nonPrimitiveNonVarBinDtype_throwsVortexException() { + // Given — a dtype with no dict carrier must fail as a VortexException, not a raw + // ClassCastException from a (DType.Primitive) cast + MemorySegment codes = u8Codes(0, 1); + MemorySegment values = TestSegments.leInts(1, 2); + + // When / Then + assertThatThrownBy(() -> decodeProtoSegments(DType.BOOL, PType.U8, codes, values, 2, 2)) + .isInstanceOf(VortexException.class) + .hasMessageContaining("unsupported dict dtype"); + } + @Test void unsupportedValuePType_throws() { // Given — F16 expands fine (2 bytes) but typedArray has no F16 mapping @@ -315,8 +328,12 @@ void legacyLayout_decodesStringsByCode() { assertThat(result.getString(2)).isEqualTo("cde"); } - @Test - void protoLayout_decodesStringsByCode() { + /// Binary shares the Utf8 dictionary shape; it used to fall into the primitive path and + /// leak a ClassCastException (Rust dict-encodes the `value` child of + /// `vortex.parquet.variant` as a Binary dictionary, #445). + @ParameterizedTest + @MethodSource("io.github.dfa1.vortex.reader.decode.DictEncodingDecoderTest#varBinDtypes") + void protoLayout_decodesStringsByCode(DType dtype) { // Given — children present: child[0]=codes, child[1]=varbin dictionary values byte[] dictBytes = "fizzbuzz".getBytes(StandardCharsets.UTF_8); // "fizz","buzz" MemorySegment bytes = MemorySegment.ofArray(dictBytes); @@ -335,7 +352,7 @@ void protoLayout_decodesStringsByCode() { ArrayNode dictNode = new ArrayNode(EncodingId.VORTEX_DICT, dictMeta, new ArrayNode[]{codesNode, valuesNode}, new int[]{}); - DecodeContext ctx = new DecodeContext(dictNode, DType.UTF8, 3, + DecodeContext ctx = new DecodeContext(dictNode, dtype, 3, segs, REGISTRY, Arena.ofAuto()); // When @@ -984,4 +1001,8 @@ private static void assertLongValues(Array array, PType valPType, long[] expecte assertThat(actual).as("index %d", i).isEqualTo(expected[i]); } } + + static Stream varBinDtypes() { + return Stream.of(DType.UTF8, DType.BINARY); + } } diff --git a/reader/src/test/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoderTest.java b/reader/src/test/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoderTest.java new file mode 100644 index 00000000..9d2bc28e --- /dev/null +++ b/reader/src/test/java/io/github/dfa1/vortex/reader/decode/ParquetVariantEncodingDecoderTest.java @@ -0,0 +1,225 @@ +package io.github.dfa1.vortex.reader.decode; + +import io.github.dfa1.vortex.core.error.VortexException; +import io.github.dfa1.vortex.core.model.DType; +import io.github.dfa1.vortex.core.model.EncodingId; +import io.github.dfa1.vortex.core.model.PType; +import io.github.dfa1.vortex.core.proto.ProtoDType; +import io.github.dfa1.vortex.core.proto.ProtoPType; +import io.github.dfa1.vortex.core.proto.ProtoParquetVariantMetadata; +import io.github.dfa1.vortex.core.proto.ProtoPrimitive; +import io.github.dfa1.vortex.core.proto.ProtoVarBinMetadata; +import io.github.dfa1.vortex.core.testing.TestSegments; +import io.github.dfa1.vortex.reader.ReadRegistry; +import io.github.dfa1.vortex.reader.array.Array; +import io.github.dfa1.vortex.reader.array.IntArray; +import io.github.dfa1.vortex.reader.array.MaskedArray; +import io.github.dfa1.vortex.reader.array.StructArray; +import io.github.dfa1.vortex.reader.array.VarBinArray; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; + +import java.lang.foreign.Arena; +import java.lang.foreign.MemorySegment; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/// Child shapes mirror Rust's `ParquetVariant::deserialize`: `[validity?, metadata, value?, +/// typed_value?]`. The interop test covers what vortex-jni writes; these cover the optional +/// children Rust's writer does not emit on its own (`typed_value` stays in the canonical +/// container there) and every malformed shape Rust rejects. +class ParquetVariantEncodingDecoderTest { + + private static final ParquetVariantEncodingDecoder SUT = new ParquetVariantEncodingDecoder(); + private static final ReadRegistry REGISTRY = TestRegistry.ofDecoders( + SUT, new PrimitiveEncodingDecoder(), new VarBinEncodingDecoder(), new BoolEncodingDecoder()); + private static final int N = 2; + + // Segments: 0 = metadata bytes, 1 = metadata offsets, 2 = value bytes, 3 = value offsets, + // 4 = typed_value ints, 5 = validity bits. + private static final MemorySegment[] SEGMENTS = { + MemorySegment.ofArray(new byte[]{0x01, 0x00, 0x01, 0x00}), TestSegments.leLongs(0, 2, 4), + MemorySegment.ofArray(new byte[]{0x0c, 0x2a, 0x0c, 0x07}), TestSegments.leLongs(0, 2, 4), + TestSegments.leInts(10, 20), + MemorySegment.ofArray(new byte[]{0b01}), + }; + + @Test + void encodingId_isParquetVariant() { + // When + EncodingId result = SUT.encodingId(); + + // Then + assertThat(result).isEqualTo(EncodingId.VORTEX_PARQUET_VARIANT); + } + + @Test + void metadataAndValue_decodeToStructOfBinaries() { + // Given + ArrayNode node = parquetVariant(meta(true, null, false), metadataNode(), valueNode()); + + // When + Array result = SUT.decode(ctx(node)); + + // Then + StructArray struct = (StructArray) result; + assertThat(((DType.Struct) struct.dtype()).fieldNames()).extracting(Object::toString).containsExactly("metadata", "value"); + assertThat(((VarBinArray) struct.field(0)).getBytes(1)).containsExactly(0x01, 0x00); + assertThat(((VarBinArray) struct.field(1)).getBytes(0)).containsExactly(0x0c, 0x2a); + assertThat(((VarBinArray) struct.field(1)).getBytes(1)).containsExactly(0x0c, 0x07); + } + + @Test + void typedValueOnly_decodesTheShreddedChildAtItsDeclaredDtype() { + // Given — no `value`: every row is fully shredded into an I32 `typed_value` + ProtoDType i32 = ProtoDType.ofPrimitive(new ProtoPrimitive(ProtoPType.I32, false)); + ArrayNode node = parquetVariant(meta(false, i32, false), metadataNode(), primitiveNode(4)); + + // When + Array result = SUT.decode(ctx(node)); + + // Then + StructArray struct = (StructArray) result; + assertThat(((DType.Struct) struct.dtype()).fieldNames()).extracting(Object::toString).containsExactly("metadata", "typed_value"); + assertThat(((DType.Struct) struct.dtype()).fieldTypes().get(1)).isEqualTo(new DType.Primitive(PType.I32, false)); + assertThat(((IntArray) struct.field(1)).getInt(1)).isEqualTo(20); + } + + @Test + void leadingValidityChild_masksTheStruct() { + // Given — one extra child beyond the expected count is the row validity, first + ArrayNode validity = new ArrayNode(EncodingId.VORTEX_BOOL, null, new ArrayNode[0], new int[]{5}); + ArrayNode node = parquetVariant(meta(true, null, false), validity, metadataNode(), valueNode()); + + // When + Array result = SUT.decode(ctx(node)); + + // Then + MaskedArray masked = (MaskedArray) result; + assertThat(masked.isValid(0)).isTrue(); + assertThat(masked.isValid(1)).isFalse(); + assertThat(masked.inner()).isInstanceOf(StructArray.class); + } + + @Nested + class Malformed { + + @Test + void nonVariantDtype_throws() { + // Given + ArrayNode node = parquetVariant(meta(true, null, false), metadataNode(), valueNode()); + DecodeContext ctx = new DecodeContext(node, DType.BINARY, N, SEGMENTS, REGISTRY, Arena.ofAuto()); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx)) + .isInstanceOf(VortexException.class) + .hasMessageContaining("expected variant dtype"); + } + + @Test + void neitherValueNorTypedValue_throws() { + // Given — empty metadata decodes as all-defaults: no value, no typed_value + ArrayNode node = new ArrayNode(EncodingId.VORTEX_PARQUET_VARIANT, null, + new ArrayNode[]{metadataNode()}, new int[0]); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx(node))) + .isInstanceOf(VortexException.class) + .hasMessageContaining("at least one of value or typed_value"); + } + + @Test + void buffers_throw() { + // Given + ArrayNode node = new ArrayNode(EncodingId.VORTEX_PARQUET_VARIANT, meta(true, null, false), + new ArrayNode[]{metadataNode(), valueNode()}, new int[]{0}); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx(node))) + .isInstanceOf(VortexException.class) + .hasMessageContaining("expected 0 buffers"); + } + + @Test + void tooFewChildren_throws() { + // Given — metadata declares `value`, but only the metadata child is present + ArrayNode node = parquetVariant(meta(true, null, false), metadataNode()); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx(node))) + .isInstanceOf(VortexException.class) + .hasMessageContaining("expected 2 or 3 children, got 1"); + } + + @Test + void tooManyChildren_throws() { + // Given + ArrayNode node = parquetVariant(meta(true, null, false), + metadataNode(), valueNode(), valueNode(), valueNode()); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx(node))) + .isInstanceOf(VortexException.class) + .hasMessageContaining("expected 2 or 3 children, got 4"); + } + + @Test + void typedValueDtypeWithUnknownPType_throwsVortexException() { + // Given — has_value=true, typed_value_dtype = Primitive{type: 99}: hand-encoded because + // the generated enum cannot name an out-of-range PType. Parsed off the wire, it must not + // leak an index or argument exception from the PType lookup. + MemorySegment meta = MemorySegment.ofArray(new byte[]{ + 0x08, 0x01, // has_value = true + 0x12, 0x04, // typed_value_dtype, 4 bytes + 0x1a, 0x02, // DType.primitive, 2 bytes + 0x08, 0x63}); // Primitive.type = 99 + ArrayNode node = parquetVariant(meta, metadataNode(), valueNode(), primitiveNode(4)); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx(node))).isInstanceOf(VortexException.class); + } + + @Test + void truncatedMetadata_throws() { + // Given — a length-delimited field (tag 2) claiming more bytes than follow + MemorySegment meta = MemorySegment.ofArray(new byte[]{0x12, 0x7f}); + ArrayNode node = parquetVariant(meta, metadataNode(), valueNode()); + + // When / Then + assertThatThrownBy(() -> SUT.decode(ctx(node))) + .isInstanceOf(VortexException.class) + .hasMessageContaining("invalid metadata"); + } + } + + private static DecodeContext ctx(ArrayNode node) { + return new DecodeContext(node, DType.VARIANT, N, SEGMENTS, REGISTRY, Arena.ofAuto()); + } + + private static MemorySegment meta(boolean hasValue, ProtoDType typedValueDtype, boolean valueNullable) { + return MemorySegment.ofArray(new ProtoParquetVariantMetadata(hasValue, typedValueDtype, valueNullable).encode()); + } + + private static ArrayNode parquetVariant(MemorySegment meta, ArrayNode... children) { + return new ArrayNode(EncodingId.VORTEX_PARQUET_VARIANT, meta, children, new int[0]); + } + + private static ArrayNode metadataNode() { + return varBinNode(0, 1); + } + + private static ArrayNode valueNode() { + return varBinNode(2, 3); + } + + private static ArrayNode varBinNode(int bytesSegment, int offsetsSegment) { + MemorySegment varBinMeta = MemorySegment.ofArray(new ProtoVarBinMetadata(ProtoPType.I64).encode()); + return new ArrayNode(EncodingId.VORTEX_VARBIN, varBinMeta, + new ArrayNode[]{primitiveNode(offsetsSegment)}, new int[]{bytesSegment}); + } + + private static ArrayNode primitiveNode(int segment) { + return new ArrayNode(EncodingId.VORTEX_PRIMITIVE, null, new ArrayNode[0], new int[]{segment}); + } +}