Skip to content

feat(iceberg): replay-safe immutable history ingestion - #1077

Merged
iamgp merged 7 commits into
mainfrom
feat/1066-immutable-history
Oct 8, 2026
Merged

iamgp merged 7 commits into
mainfrom
feat/1066-immutable-history

Conversation

@iamgp

@iamgp iamgp commented Oct 7, 2026

Copy link
Copy Markdown
Collaborator

Scope

Closes #1066. Stacked on #1075 (feat/1065-iceberg-schema-policy, exact base e47b0b5), which follows #1064. This diff contains only immutable-history support, evidence, tests, and documentation; no merge or schema-policy reimplementation.

Behaviour

  • Add merge_strategy="history" through supported DLT ingestion and an optional provider capability; unsupported providers reject before table creation or writes.
  • Require composite entity/version identity and exactly one explicit payload-column set or supplied hash. Ignore arrival metadata, skip identical committed and batch observations, retain late versions, and reject whole-batch payload conflicts without publishing schema, policy, or data.
  • Bind canonical identity/payload policy and fingerprint atomically with first insertion. Reject unmarked nonempty adoption and incompatible policy changes; migration requires quiesced writers and validation.
  • Use one loaded Iceberg snapshot for lookup and append. Definite conflicts retry the complete operation at most three times; ambiguous outcomes reconcile from fresh metadata and fail explicitly without blind append or invented counts.
  • Report incoming inserted/skipped/conflicting counts and unknown-outcome reconciliation in normal operation evidence. Combine all staged files into one batch.

The uniqueness guarantee covers cooperating history writers on the same table/ref. Ordinary writes and property-only migration must not race history ingestion. Batches must fit memory. Domain identity generation, latest-version ordering, and live migration are outside scope.

Validation

  • make setup: passed.
  • make check: passed; 5,932 passed, 4 skipped, 229 deselected. An initial shallow-history header-test failure disappeared after fetching full repository history; the complete baseline was rerun successfully.
  • uv run --locked pytest packages/phlo-iceberg/tests packages/phlo-dlt/tests -m "not integration" --tb=short -q: 424 passed, 7 deselected, including 58 new history cases.
  • uv run --locked python scripts/run_integration.py: 111 passed, no skips; includes real Nessie/MinIO history replay, conflict rejection, and injected concurrent commit retry. An initial test-only catalog-instance injection mistake was corrected and the full lane rerun.
  • make test-core-regression: 360 passed, 2,381 deselected.
  • make docs-build: passed, 849 static pages; inspected rendered guide and API reference, including scrolling the code example to its final lines.
  • Changed-doc codespell, header checks, Ruff, typing, complexity, and git diff --check: passed.

Real-catalog cases cover committed and in-batch conflicts, composite identities, late versions, null/NaN payloads, supplied hashes, multifile input, bounded retries, concurrent empty/nonempty tables, additive races, competing policies, unknown landed/unlanded outcomes, and evidence. The public decorated ingestion test stages actual DLT data.

No dependency changes. PyIceberg 0.11.0 (declared minimum) and locked 0.11.1 both emit the snapshot-head assertion, including an absent initial head.

iamgp added 3 commits October 7, 2026 17:23
Validate and align incoming data before staging writes. Use one Iceberg transaction for every delete batch and the replacement append, and record post-commit evidence only after publication.

Add local catalog acceptance tests for multiple delete batches, rollback, replay, real commit conflicts, resource evidence, and DLT ingestion with ref isolation.

Closes #1064
Expose strict, additive, and drop_extra policies on ingestion and storage writes. Validate Arrow batches before mutation and stage nullable additions with data in the existing transaction.

Reconcile compatible concurrent additions, preserve unrelated provider defaults, and document the intentional compatibility change.
Compare explicit entity/version identity and payload policy across a complete staged batch. Bind policy and schema additions with snapshot-guarded insertion, retry complete operations on definite conflicts, and reconcile ambiguous outcomes without claiming counts.

Expose optional provider support, ingestion evidence, real-catalog and service acceptance tests, and history configuration and migration documentation.
@coldtea-pr-lens

coldtea-pr-lens Bot commented Oct 7, 2026 •

Copy link
Copy Markdown

Note

This drawing shows 3234d6e, and the branch has new commits since. Tick Redraw to draw the latest one

  • Redraw

Nothing flagged · reviewed 3234d6e


Architecture

Architecture diagram for phlohouse/phlo at 3234d6e

Play the walkthrough


Inside the changed components — 2 views

Component view — dlt Ingestion Pipeline

Internal decorator, executor, and helper components within phlo-dlt.

Architecture view of Component view — dlt Ingestion Pipeline in phlohouse/phlo

Component view — Iceberg History Engine

Internal resource adapter and history writer within phlo-iceberg.

Architecture view of Component view — Iceberg History Engine in phlohouse/phlo

Data flow

Data flow diagram for phlohouse/phlo at 3234d6e

Follow each request


The other flows — 1 sequence

Reconciling an ambiguous commit outcome

Sequence diagram of Reconciling an ambiguous commit outcome in phlohouse/phlo

View

  • Architecture lens
  • Data flow lens
  • Expand every detail

Tip

Open a diagram on the canvas, then press W or click play to walk through the change one step at a time

🪧 More tips
  • Run npx skills add coldteadotai/pr-lens, then tell your coding agent: "Diagram the change you just made with PR Lens and attach it to the pull request."
  • Run npx @coldtea/pr-lens-cli analyze --base origin/main on a branch, then npx @coldtea/pr-lens-cli render .pr-lens/graph.json. Same lenses, your own model key, before the pull request exists
  • Untick Architecture lens or Data flow lens under View to hide a diagram, or tick Expand every detail to open every section. The comment redraws in a few seconds
  • Click the link under each diagram to open it on a canvas you can zoom, pan and step through
  • The diagrams are links. Click one to open it on the canvas, then press W or click play to walk through the change
  • The CLI's render reads .github/pr-lens.yml and applies your renames, exclusions and lane pins at draw time
  • Set github.comment.collapsed: true in .github/pr-lens.yml to fold the comment behind one View architecture and data flow row. Drawing still runs as before
  • Set github.draw: on-demand in .github/pr-lens.yml and PR Lens stops drawing pull requests on its own. Comment @pr-lens draw on a pull request when you want that one drawn
  • Add .github/workflows/pr-lens.yml with coldteadotai/pr-lens/packages/action@v0 and your model provider's key as its api-key to run PR Lens from your own CI. Any /chat/completions endpoint works
  • Push a commit and the drawing stays, with a note that it is out of date. Tick Redraw in the note to draw the new head
  • Switch GitHub to dark mode and the diagrams follow. The moving dots are this pull request's data in motion

Thanks for using PR Lens! It's built by Coldtea, free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

@phlo-agent

phlo-agent Bot commented Oct 7, 2026

Copy link
Copy Markdown
Contributor

Reviewed 3234d6e against #1066: the full diff, src/phlo/capabilities/history.py, src/phlo/capabilities/interfaces.py, packages/phlo-iceberg/src/phlo_iceberg/history.py and resource.py, the DLT decorator/helper/executor wiring, both new test suites plus the integration case, and the three documentation updates. The atomicity, conflict-rejection, policy-binding and retry design look correct; I found no correctness, security or compatibility defects. Two low-severity findings, both in the new lookup and reconciliation evidence.

1. Lookup is filtered on the version key only, so read cost scales with the table rather than the batch

packages/phlo-iceberg/src/phlo_iceberg/history.py:137-148 builds versions from the batch's distinct version_key values and pushes down only In(term=Reference(policy.version_key), ...); the entity key is matched in Python afterwards, and to_pylist() materialises every row that predicate returns.

Under the identity scheme in #1066, version_id values such as V-1 are shared across entities, so a batch touching V-1…V-50 reads and materialises every stored row carrying any of those values across all entities, not just the incoming pairs. The reference docs state only that "batch size must fit available memory" (docs/reference/python-api.md), so a small batch against a large table can still exhaust memory in a way users are not warned about. Scoping the predicate to the incoming (entity_key, version_key) pairs — or documenting the real lookup cost — would close the gap.

2. reconciliation.rows_present counts incoming rows, not stored rows

packages/phlo-iceberg/src/phlo_iceberg/history.py:205-207 computes rows_present as sum(len(members) for key, members in groups.items() if key in existing) — incoming rows whose version key is already committed. test_unknown_commit_reconciles_and_fails_without_retry_or_invented_counts in packages/phlo-iceberg/tests/test_history.py shows the consequence directly: a two-row identical replay that lands one row reports rows_present: 2 while the table holds one row. versions_missing counts keys, so the two fields also use different units. Since the docs instruct operators to read this evidence before explicitly replaying an ambiguous commit, a stored-row count (or a name such as incoming_rows_already_present) would avoid overstating what landed. Related: in the "observed" branch rows_conflicting is always 0, because _compare raises on any conflict before this dict is built.

Checked and deliberately not flagged. The central concurrency claim holds. In PyIceberg 0.11.0/0.11.1 an ordinary append staged from a branch always emits AssertRefSnapshotId, using snapshot_id=None for a snapshot-less table (snapshot.py), so a second writer staged from the same state fails instead of double-appending. CommitStateUnknownException subclasses RESTError, not CommitFailedException (exceptions.py), so the ambiguous path is not swallowed by the retry branch, and Transaction.set_properties + append + commit_transaction() is a supported explicit-commit pattern. Multi-file batches are also safe: concat_tables(promote_options="default") null-fills missing columns, which is exactly why the per-file _incoming check must run before concatenation — and it does.

Validation not reproduced here. I ran no checks in this review; the make check, focused pytest, scripts/run_integration.py and make docs-build results in the description are taken as reported rather than verified.

iamgp added 3 commits October 8, 2026 07:34
Push exact composite identity predicates into snapshot-pinned scans in groups of at most 1000 pairs. Report distinct matching, missing, and conflicting version counts after ambiguous commits rather than counting duplicate observations as stored rows.

Exercise real scan results, crossed identities, batching, landed and unlanded outcomes, reconciliation conflicts, and public ingestion evidence. Document distinct-version reconciliation semantics.
Forward explicit migration policies across every chunk and Kafka policies only to opted-in stores. Formalise the optional provider contract and validate advertised signatures before writes while preserving legacy defaults.

Project unified carrier business columns before strict writes, exercise real materialisation and local-catalog wrapper behavior, and document configuration and compatibility boundaries.
Preserve the optional schema-policy signature negotiation and validate the selected history writer before table creation. Retain both published histories and add history-specific contract coverage.
@iamgp
iamgp changed the base branch from feat/1065-iceberg-schema-policy to main October 8, 2026 09:11
@iamgp
iamgp added this pull request to the merge queue Oct 8, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to a conflict with the base branch Oct 8, 2026
Preserve the approved history-only delta on top of the schema-policy squash and main export/service changes. Verify the resolved tree exactly equals main plus the previously approved history patch.
@iamgp
iamgp added this pull request to the merge queue Oct 8, 2026
Merged via the queue into main with commit dd07c57 Oct 8, 2026
28 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

feat(phlo-iceberg): replay-safe append-only version history ingestion

1 participant