moonice

    MoonBit-native Iceberg v2 reading with partition evolution, delete semantics and explainable scan planning.

    iceberg
    snapshot
    data-lake
    scan-planning
    Download zip
    Author
    Version
    0.2.1
    License
    Apache-2.0
    Last updated
    9 hours ago
    Downloads
    14

    #MoonIce API examples

    For setup, supported formats and the interactive inspector, see README.md.

    #Evaluate a partition transform

    ///|
    test "negative timestamps floor to the previous day and strings count code points" {
    assert_eq(
    @moonice.transform_partition(
    @moonice.Integer(-1L),
    Json::string("timestamp"),
    "day",
    ),
    Some(@moonice.Integer(-1L)),
    )
    assert_eq(
    @moonice.transform_partition(
    @moonice.Text("月🌙冰"),
    Json::string("string"),
    "truncate[2]",
    ),
    Some(@moonice.Text("月🌙")),
    )
    }

    #Process bounded batches

    Run moon run --target js examples/batch_read from the repository root. The example reads an independently generated table whose partition specification changes from day to month, filters epoch microseconds, and sums IDs without collecting a complete result array. Expected: two batches, two rows, sum 21. scan_batches applies deletes and residual filters before invoking the callback; arrays are not reused. A later error can follow already delivered batches, so atomic consumers must stage their output. Individual Parquet files still decode in memory. See compatibility limits.

    #Bind predicates to stable field IDs

    ///|
    test "bind a filter using the current column name" {
    let schema : @moonice.Schema = {
    id: 1,
    fields: [
    {
    id: 2,
    name: "region",
    required: false,
    field_type: Json::string("string"),
    },
    ],
    }
    let p = @moonice.parse_predicate(
    "{\"field\":\"region\",\"op\":\"=\",\"value\":\"北京\"}", schema,
    )
    assert_true(p.matches(Map([(2, @moonice.Text("北京"))])))
    assert_false(p.matches(Map([(2, @moonice.Missing)])))
    }

    #Preserve row identity across schema changes

    ///|
    test "renaming preserves field values and new columns null-fill" {
    let row : @moonice.DataRow = {
    file_path: "memory://table/data.parquet",
    position: 0L,
    values: Map([(2, @moonice.Text("北京"))]),
    }
    let schema : @moonice.Schema = {
    id: 1,
    fields: [
    {
    id: 2,
    name: "region",
    required: false,
    field_type: Json::string("string"),
    },
    {
    id: 3,
    name: "note",
    required: false,
    field_type: Json::string("string"),
    },
    ],
    }
    let result = row.project(schema)
    assert_eq(result.get("region"), Some(@moonice.Text("北京")))
    assert_eq(result.get("note"), Some(@moonice.Missing))
    }

    #Integration flow

    ///|
    let metadata = @moonice.parse_metadata(raw_metadata_json)

    ///|
    let state = @moonice.load_snapshot(metadata, path => storage_read(path))

    ///|
    let predicate = @moonice.parse_predicate(filter_json, state.schema)

    ///|
    let plan = @moonice.plan_scan(metadata, state, predicate)

    ///|
    let result = @moonice.scan_rows(metadata, state, predicate, path => {
    storage_read(path)
    })

    storage_read is an application-provided (String) -> Bytes raise @moonice.IceError callback. It receives original object URIs and does not need directory listing. Core APIs never open sockets or local files themselves. For a complete executable integration using public library APIs and standard file bytes, run moon run --target js examples/library_read. See the source and integration notes.

    IceError

    pub(all) suberror IceError {
    Invalid(String, String, String)
    } derive(
    Debug
    )

    A rejected input, with a stable machine-readable code and source path.
    impl Show for IceError

    IceError::output

    fn IceError::output(self : IceError, logger : &Logger) -> Unit

    IceError::to_repr

    IceError::to_string

    fn IceError::to_string(self : IceError) -> String

    Bundle

    pub(all) struct Bundle {
    metadata : TableMetadata
    files : Map[String, Bytes]
    } derive(
    Debug
    )

    An offline transport for standard Iceberg bytes, not a replacement table format.

    Bundle::diagnose

    fn Bundle::diagnose(self : Bundle) -> Diagnosis

    Inspect every retained snapshot. Failures are localized; other snapshots are still inspected. Unreferenced bundled objects are reported, never deleted.

    Bundle::load_snapshot

    fn Bundle::load_snapshot(self : Bundle, snapshot_id? : Int64) -> SnapshotState raise IceError

    Bundle::read_file

    fn Bundle::read_file(self : Bundle, path : String) -> Bytes raise IceError

    Bundle::scan

    fn Bundle::scan(self : Bundle, predicate : Predicate, snapshot_id? : Int64, limit? : Int) -> ScanResult raise IceError

    Bundle::scan_batches

    fn Bundle::scan_batches(self : Bundle, predicate : Predicate, emit : (Array[DataRow]) -> Unit raise IceError, snapshot_id? : Int64, batch_size? : Int) -> ScanSummary raise IceError

    Bundle::to_repr

    CompareOp

    pub(all) enum CompareOp {
    Eq
    Lt
    Le
    Gt
    Ge
    } derive(Eq, ToJson,
    Debug
    )

    CompareOp::equal

    fn CompareOp::equal(CompareOp, CompareOp) -> Bool

    CompareOp::not_equal

    fn CompareOp::not_equal(x : CompareOp, y : CompareOp) -> Bool

    CompareOp::to_json

    fn CompareOp::to_json(CompareOp) -> Json

    DataFile

    pub(all) struct DataFile {
    content : Int
    path : String
    format : String
    partition : Map[String, Scalar]
    record_count : Int64
    size_bytes : Int64
    lower_bounds : Map[Int, Bytes]
    upper_bounds : Map[Int, Bytes]
    null_counts : Map[Int, Int64]
    nan_counts : Map[Int, Int64]
    equality_ids : Array[Int]
    } derive(Eq, ToJson,
    Debug
    )

    File content: 0=data, 1=position deletes, 2=equality deletes.

    DataFile::equal

    fn DataFile::equal(DataFile, DataFile) -> Bool

    DataFile::not_equal

    fn DataFile::not_equal(x : DataFile, y : DataFile) -> Bool

    DataFile::to_json

    fn DataFile::to_json(DataFile) -> Json

    DataFile::to_repr

    DataRow

    pub(all) struct DataRow {
    file_path : String
    position : Int64
    values : Map[Int, Scalar]
    } derive(Eq, ToJson,
    Debug
    )

    DataRow::equal

    fn DataRow::equal(DataRow, DataRow) -> Bool

    DataRow::not_equal

    fn DataRow::not_equal(x : DataRow, y : DataRow) -> Bool

    DataRow::project

    fn DataRow::project(self : DataRow, schema : Schema) -> Map[String, Scalar] raise IceError

    Project by field ID, so renamed columns survive and newly added columns null-fill.

    DataRow::to_json

    fn DataRow::to_json(DataRow) -> Json

    DataRow::to_repr

    DeleteDecision

    pub(all) struct DeleteDecision {
    data_path : String
    delete_path : String
    applies : Bool
    reason_code : String
    explanation : String
    } derive(Eq, ToJson,
    Debug
    )

    DeleteDecision::equal

    DeleteDecision::not_equal

    fn DeleteDecision::not_equal(x : DeleteDecision, y : DeleteDecision) -> Bool

    DeleteDecision::to_json

    Diagnosis

    pub(all) struct Diagnosis {
    checked_snapshots : Int
    checked_objects : Int
    findings : Array[Finding]
    } derive(ToJson,
    Debug
    )

    Diagnosis::to_json

    fn Diagnosis::to_json(Diagnosis) -> Json

    Field

    pub(all) struct Field {
    id : Int
    name : String
    required : Bool
    field_type : Json
    } derive(Eq, ToJson,
    Debug
    )

    Field::equal

    fn Field::equal(Field, Field) -> Bool

    Field::not_equal

    fn Field::not_equal(x : Field, y : Field) -> Bool

    Field::to_json

    fn Field::to_json(Field) -> Json

    Field::to_repr

    FileDecision

    pub(all) struct FileDecision {
    entry : ManifestEntry
    kept : Bool
    reason_code : String
    explanation : String
    } derive(ToJson,
    Debug
    )

    FileDecision::to_json

    Finding

    pub(all) struct Finding {
    severity : String
    code : String
    path : String
    explanation : String
    snapshot_id : Int64?
    } derive(Eq, ToJson,
    Debug
    )

    Finding::equal

    fn Finding::equal(Finding, Finding) -> Bool

    Finding::not_equal

    fn Finding::not_equal(x : Finding, y : Finding) -> Bool

    Finding::to_json

    fn Finding::to_json(Finding) -> Json

    Finding::to_repr

    Manifest

    pub(all) struct Manifest {
    path : String
    length : Int64
    spec_id : Int
    content : Int
    sequence_number : Int64
    min_sequence_number : Int64
    added_snapshot_id : Int64
    } derive(Eq, ToJson,
    Debug
    )

    Manifest::equal

    fn Manifest::equal(Manifest, Manifest) -> Bool

    Manifest::not_equal

    fn Manifest::not_equal(x : Manifest, y : Manifest) -> Bool

    Manifest::to_json

    fn Manifest::to_json(Manifest) -> Json

    Manifest::to_repr

    ManifestEntry

    pub(all) struct ManifestEntry {
    manifest_path : String
    spec_id : Int
    status : Int
    snapshot_id : Int64
    sequence_number : Int64
    file_sequence_number : Int64
    file : DataFile
    } derive(Eq, ToJson,
    Debug
    )

    Status: 0=existing, 1=added, 2=deleted. Sequence inheritance is resolved.

    ManifestEntry::equal

    ManifestEntry::not_equal

    fn ManifestEntry::not_equal(x : ManifestEntry, y : ManifestEntry) -> Bool

    ManifestEntry::to_json

    PartitionField

    pub(all) struct PartitionField {
    source_id : Int
    field_id : Int
    name : String
    transform : String
    } derive(Eq, ToJson,
    Debug
    )

    PartitionField::equal

    PartitionField::not_equal

    fn PartitionField::not_equal(x : PartitionField, y : PartitionField) -> Bool

    PartitionField::to_json

    PartitionSpec

    pub(all) struct PartitionSpec {
    id : Int
    fields : Array[PartitionField]
    } derive(Eq, ToJson,
    Debug
    )

    PartitionSpec::equal

    PartitionSpec::not_equal

    fn PartitionSpec::not_equal(x : PartitionSpec, y : PartitionSpec) -> Bool

    PartitionSpec::to_json

    Predicate

    pub(all) enum Predicate {
    All
    Compare(Int, CompareOp, Scalar)
    IsNull(Int)
    NotNull(Int)
    And(Predicate, Predicate)
    Or(Predicate, Predicate)
    } derive(Eq, ToJson,
    Debug
    )

    Predicates bind to stable field IDs. Missing values never satisfy comparisons.

    Predicate::equal

    fn Predicate::equal(Predicate, Predicate) -> Bool

    Predicate::matches

    fn Predicate::matches(self : Predicate, row : Map[Int, Scalar]) -> Bool

    Predicate::not_equal

    fn Predicate::not_equal(x : Predicate, y : Predicate) -> Bool

    Predicate::to_json

    fn Predicate::to_json(Predicate) -> Json

    Predicate::validate

    fn Predicate::validate(self : Predicate, schema : Schema) -> Unit raise IceError

    Scalar

    pub(all) enum Scalar {
    Missing
    Boolean(Bool)
    Integer(Int64)
    Real(Double)
    Text(String)
    Binary(Bytes)
    } derive(Eq, ToJson,
    Debug
    )

    Scalar::display_json

    fn Scalar::display_json(self : Scalar) -> Json

    Scalar::equal

    fn Scalar::equal(Scalar, Scalar) -> Bool

    Scalar::not_equal

    fn Scalar::not_equal(x : Scalar, y : Scalar) -> Bool

    Scalar::to_json

    fn Scalar::to_json(Scalar) -> Json

    Scalar::to_repr

    ScanPlan

    pub(all) struct ScanPlan {
    snapshot_id : Int64
    schema : Schema
    predicate : Predicate
    decisions : Array[FileDecision]
    retained_bytes : Int64
    pruned_bytes : Int64
    } derive(ToJson,
    Debug
    )

    ScanPlan::to_json

    fn ScanPlan::to_json(ScanPlan) -> Json

    ScanPlan::to_repr

    ScanResult

    pub(all) struct ScanResult {
    snapshot_id : Int64
    schema : Schema
    rows : Array[DataRow]
    read_rows : Int
    deleted_rows : Int
    matched_rows : Int
    truncated : Bool
    } derive(ToJson,
    Debug
    )

    ScanResult::to_json

    fn ScanResult::to_json(ScanResult) -> Json

    ScanSummary

    pub(all) struct ScanSummary {
    snapshot_id : Int64
    schema : Schema
    read_rows : Int
    deleted_rows : Int
    matched_rows : Int
    scanned_files : Int
    delete_files : Int
    } derive(ToJson,
    Debug
    )

    Counters for a complete scan. Unlike ScanResult this contains no row array.

    ScanSummary::to_json

    fn ScanSummary::to_json(ScanSummary) -> Json

    Schema

    pub(all) struct Schema {
    id : Int
    fields : Array[Field]
    } derive(Eq, ToJson,
    Debug
    )

    Schema::equal

    fn Schema::equal(Schema, Schema) -> Bool

    Schema::field

    fn Schema::field(self : Schema, id : Int) -> Field?

    Schema::field_named

    fn Schema::field_named(self : Schema, name : String) -> Field raise IceError

    Schema::not_equal

    fn Schema::not_equal(x : Schema, y : Schema) -> Bool

    Schema::to_json

    fn Schema::to_json(Schema) -> Json

    Schema::to_repr

    SchemaChange

    pub(all) struct SchemaChange {
    field_id : Int
    kind : String
    before : Field?
    after : Field?
    } derive(ToJson,
    Debug
    )

    SchemaChange::to_json

    Snapshot

    pub(all) struct Snapshot {
    id : Int64
    parent_id : Int64?
    sequence_number : Int64
    timestamp_ms : Int64
    manifest_list : String
    schema_id : Int?
    operation : String
    } derive(Eq, ToJson,
    Debug
    )

    Snapshot::equal

    fn Snapshot::equal(Snapshot, Snapshot) -> Bool

    Snapshot::not_equal

    fn Snapshot::not_equal(x : Snapshot, y : Snapshot) -> Bool

    Snapshot::to_json

    fn Snapshot::to_json(Snapshot) -> Json

    Snapshot::to_repr

    SnapshotDiff

    pub(all) struct SnapshotDiff {
    from_id : Int64
    to_id : Int64
    added : Array[ManifestEntry]
    removed : Array[ManifestEntry]
    retained_files : Int
    added_records : Int64
    removed_records : Int64
    schema_changes : Array[SchemaChange]
    } derive(ToJson,
    Debug
    )

    SnapshotDiff::to_json

    SnapshotRef

    pub(all) struct SnapshotRef {
    name : String
    snapshot_id : Int64
    kind : String
    } derive(Eq, ToJson,
    Debug
    )

    SnapshotRef::equal

    fn SnapshotRef::equal(SnapshotRef, SnapshotRef) -> Bool

    SnapshotRef::not_equal

    fn SnapshotRef::not_equal(x : SnapshotRef, y : SnapshotRef) -> Bool

    SnapshotRef::to_json

    fn SnapshotRef::to_json(SnapshotRef) -> Json

    SnapshotState

    pub(all) struct SnapshotState {
    snapshot : Snapshot
    schema : Schema
    manifests : Array[Manifest]
    entries : Array[ManifestEntry]
    } derive(ToJson,
    Debug
    )

    SnapshotState::to_json

    TableMetadata

    pub(all) struct TableMetadata {
    uuid : String
    location : String
    current_schema_id : Int
    default_spec_id : Int
    last_sequence_number : Int64
    current_snapshot_id : Int64?
    schemas : Array[Schema]
    partition_specs : Array[PartitionSpec]
    snapshots : Array[Snapshot]
    refs : Array[SnapshotRef]
    } derive(ToJson,
    Debug
    )

    Parsed Iceberg v2 metadata. IDs remain signed 64-bit integers on every target.

    TableMetadata::schema

    fn TableMetadata::schema(self : TableMetadata, id : Int) -> Schema raise IceError

    TableMetadata::snapshot

    fn TableMetadata::snapshot(self : TableMetadata, id? : Int64, ref_name? : String) -> Snapshot raise IceError

    TableMetadata::to_json

    decode_manifest

    fn decode_manifest(data : Bytes, manifest : Manifest) -> Array[ManifestEntry] raise IceError

    Decode a v2 manifest, resolving null sequence numbers using its list entry.

    decode_manifest_list

    fn decode_manifest_list(data : Bytes) -> Array[Manifest] raise IceError

    Read a standard Avro manifest list. Only null/deflate OCF codecs are supported.

    diff_snapshots

    fn diff_snapshots(before : SnapshotState, after : SnapshotState) -> SnapshotDiff

    Physical file/record changes, not a logical row-level changelog. Rewrites may remove and re-add the same logical rows; delete files are reported separately.

    explain_delete

    fn explain_delete(metadata : TableMetadata, data : ManifestEntry, deletion : ManifestEntry) -> DeleteDecision raise IceError

    Decide delete-file applicability by Iceberg v2 data sequence and partition. Position file paths are matched against each delete row during execution.

    load_snapshot

    fn load_snapshot(metadata : TableMetadata, read_file : (String) -> Bytes raise IceError, snapshot_id? : Int64) -> SnapshotState raise IceError

    Load one snapshot using caller-provided storage. No directory listing is required. Only live entries are returned; deleted manifest entries are history.

    open_bundle

    fn open_bundle(source : String) -> Bundle raise IceError

    Open an offline bundle. Raw metadata is a string to preserve 64-bit JSON IDs when a browser or another JSON consumer transports the outer document.

    parquet_field_ids

    fn parquet_field_ids(bytes : Bytes) -> Array[Int] raise IceError

    Return flat leaf field IDs in Parquet column order; never guess from names.

    parse_metadata

    fn parse_metadata(source : String) -> TableMetadata raise IceError

    Parse v2 table metadata without I/O. Malformed input and other format versions are rejected; missing historical parents are allowed because snapshots expire.

    parse_predicate

    fn parse_predicate(source : String, schema : Schema) -> Predicate raise IceError

    JSON grammar: {field: name, op: =|<|<=|>|>=|is_null|not_null, value: literal}, {and: [predicate,predicate]} or {or: [...]}; null means all rows.

    plan_deletes

    fn plan_deletes(metadata : TableMetadata, state : SnapshotState, plan : ScanPlan) -> Array[DeleteDecision] raise IceError

    plan_scan

    fn plan_scan(metadata : TableMetadata, state : SnapshotState, predicate : Predicate) -> ScanPlan raise IceError

    Inclusive planning: every retained file still requires residual row filtering.

    read_data_rows

    fn read_data_rows(bytes : Bytes, file : DataFile) -> Array[DataRow] raise IceError

    Materialize a bounded, flat Parquet file with original positions and IDs.

    run_request

    fn run_request(bundle_json : String, options_json : String) -> String

    Transport-neutral JSON interface used by both the CLI and browser. Int64 output values are decimal strings. Errors always include code/path/message.

    scan_batches

    fn scan_batches(metadata : TableMetadata, state : SnapshotState, predicate : Predicate, read_file : (String) -> Bytes raise IceError, emit : (Array[DataRow]) -> Unit raise IceError, batch_size? : Int) -> ScanSummary raise IceError

    Emit matching rows in batches, after deletes and schema validation. Each emitted array is owned by the caller and is never reused or cleared. A callback error stops further scanning. Earlier batches remain delivered if a later read fails; callers needing atomic output must stage it themselves. Files are still decoded in memory (at most 100,000 rows/64 MiB per file). This is incremental delivery between files, not a Parquet range-read API.

    scan_rows

    fn scan_rows(metadata : TableMetadata, state : SnapshotState, predicate : Predicate, read_file : (String) -> Bytes raise IceError, limit? : Int) -> ScanResult raise IceError

    Apply deletes before residual filtering and projection. The limit limits returned rows; counters still describe the complete bounded scan.

    transform_partition

    fn transform_partition(value : Scalar, field_type : Json, transform : String) -> Scalar?

    Evaluate a supported Iceberg v2 transform. None means unsupported type, unknown/malformed transform, or arithmetic outside its representable domain. Some(Missing) is a known null partition, distinct from an unknown result. Timestamp inputs are signed microseconds since the Unix epoch, date inputs are signed days. String truncation counts Unicode code points.