Warning

This is not authoritative documentation. It describes a plan of work that is not yet implemented and will change as it lands. Once the work is complete this document should be removed; the finished system is described by the Iceberg storage engine spec, the plugin’s user documentation, and the repos themselves.

Task 4: Read path — planner, executor, cursor, pinning

Repo: https://opendev.org/drizzle/drizzle. Depends on task 3. This is the core read deliverable: after it, use case 1 (SQL over pipeline output) works, full-scan. Three commits.

Design constraints from the design doc, restated as binding:

  • Snapshot resolution is an explicit planning step producing a value.

  • Planner output must be constructible from a received task list as well as from local planning (the future distribution seam) — enforce by giving the Planner a constructor/factory taking vector<FileScanTask> directly, used by tests today.

  • Pinning is per transaction (autocommit → per statement). Never per cursor.

Commit 1: planner and executor (no cursor yet; unit-tested standalone)

  • planner.{h,cc}: (table handle, snapshot-id, projected iceberg field-ids) → immutable ordered vector<FileScanTask> via iceberg-cpp scan planning (V2 deletes handled by the library). Plus the from-existing-tasks factory.

  • executor.{h,cc}: one FileScanTask → row stream. Internally: Arrow record-batch reader with column projection by field-id; exposes next_row() filling a caller-provided callback/row-view, plus seek(row_ordinal) used by rnd_pos (implement as re-open + skip if the reader lacks random access — correctness first; note the cost in a comment).

  • Value conversion Arrow → Drizzle Field::store lives here, the inverse of task 3’s schema map and colocated with it conceptually: cover every mapped type, with explicit UTC handling for timestamptz→EPOCH and exactness tests for decimal.

Commit 2: the cursor, wired

cursor.{h,cc} (IcebergCursor : Cursor), modeled on plugin/function_engine/cursor.{h,cc} for shape:

  • doOpen: resolve the per-table share (schema map, field-id table).

  • doStartTableScan(bool): obtain (snapshot-id, plan) through the session pin (below); position at task 0.

  • rnd_next(buf): drain executors task by task; pack via the read map — honor Table::read_set (table.h:497) so unread columns are neither fetched from Parquet nor stored. Add HTON_PARTIAL_COLUMN_READ to the engine flags in this commit.

  • position(record): encode (task ordinal: u32, row ordinal: u64) into ref; set ref_length accordingly at open.

  • rnd_pos(buf, pos): decode; executor seek within the pinned plan.

  • info(flag): stats.records from the pinned snapshot’s summary total-records (estimate; the engine deliberately does not set HTON_STATS_RECORDS_IS_EXACT), stats.mean_rec_length from the proto.

  • close() / doEndTableScan: release executors; the plan is owned by the session pin entry, not the cursor (two cursors on one table share one plan).

  • No index methods implemented; doStartIndexScan never reachable (no keys are ever declared in the synthesized proto).

Session pinning: session_state.{h,cc}

Per-session engine slot via Session::getEngineData(const plugin::MonitoredInTransaction*) (session.h:345; the innobase trx_t* pattern, ha_innodb.cc:1030). Holds map<table-uuid, PinEntry{snapshot-id, shared plan cache}>. First touch of a table in a transaction resolves current-snapshot and pins; subsequent touches (including second cursors in the same statement — the self-join case) reuse the pin. Clearing: engine doEndStatement (storage_engine.h:150-155) when no transaction is open (autocommit), and transaction end otherwise — in this task the engine is not yet transactional, so wire statement-end clearing now and leave a // task 5 moves clearing to doCommit/doRollback when a transaction is open marker that task 5 must resolve (no other tombstones).

Commit 3: test suites

  • Full-scan SELECT over each fixture table; checksums match pyiceberg-read values recorded in the fixtures (interop, not self-consistency).

  • The position-deletes fixture: deleted rows absent.

  • ORDER BY on a non-trivial table (exercises filesort → position/ rnd_pos).

  • Self-join on a fixture table while a background pyiceberg append commits mid-statement (the suite’s seeding hook can do this): both sides join the same snapshot. Also: two statements in one autocommit session straddling an external commit see old-then-new (statement-level pinning observable).

  • Projection: SELECT one_col FROM wide_fixture — assert via the callgrind CI job or a byte-counter that unprojected columns are not read (if neither is practical, a unit test on the executor’s projection set is the floor).

Verification

  • Build green in both switch states (--with-iceberg on/off) per commit; suites green; job voting.

  • Valgrind server start/scan/stop: no leaks attributable to the plugin; massif comparison sane on the wide-table scan (row-at-a-time packing must not accumulate whole files in memory — batch-at-a-time residency only).