From fdefeaf308f188b233b45616b3b677eb7052e569 Mon Sep 17 00:00:00 2001 From: Hasyimi Bahrudin Date: Mon, 5 Oct 2026 22:01:35 +0800 Subject: [PATCH] fix: schema changes applied ahead of rows staged before them Relation messages evolved the Iceberg schema the moment they arrived, ahead of rows already staged under the old schema and keyed by column name. A column dropped and re-added took those rows' values for the dropped one. Schema changes now go through the change log, in order with the rows: - The pipeline stages a Relation message that changes a table's columns as a schema event (staged op `R`) where the stream put it, starting a new log entry. - The materializer, reaching an entry that starts with one, commits the rows before it, then applies the change against the catalog's current schema. A crash in between re-reads that entry, and applying the change again is a no-op. - The lifecycle no longer applies Relation messages itself (`Materializer::apply_relation` is gone). Co-Authored-By: Claude Opus 5.5 --- crates/pg2iceberg-logical/src/lib.rs | 1 + crates/pg2iceberg-logical/src/materializer.rs | 148 ++++++++--------- crates/pg2iceberg-logical/src/pipeline.rs | 71 +++++++-- .../pg2iceberg-logical/src/relation_event.rs | 43 +++++ crates/pg2iceberg-logical/src/sink.rs | 107 +++++++++++-- crates/pg2iceberg-logical/tests/end_to_end.rs | 24 +-- .../tests/materializer_e2e.rs | 5 +- crates/pg2iceberg-pg/src/lib.rs | 11 +- crates/pg2iceberg-pg/src/prod/stream.rs | 2 +- crates/pg2iceberg-sim/src/postgres.rs | 16 +- crates/pg2iceberg-stream/src/codec.rs | 75 +++++---- crates/pg2iceberg-tests/tests/dst.rs | 150 ++++++++++++++---- crates/pg2iceberg-tests/tests/dst_decimal.rs | 7 - .../tests/dst_dml_correctness.rs | 11 +- .../pg2iceberg-tests/tests/dst_meta_tables.rs | 4 - .../tests/dst_schema_evolution.rs | 99 +++++++++--- crates/pg2iceberg-validate/src/runtime.rs | 18 +-- 17 files changed, 549 insertions(+), 243 deletions(-) create mode 100644 crates/pg2iceberg-logical/src/relation_event.rs diff --git a/crates/pg2iceberg-logical/src/lib.rs b/crates/pg2iceberg-logical/src/lib.rs index 3c07698..fcbe02d 100644 --- a/crates/pg2iceberg-logical/src/lib.rs +++ b/crates/pg2iceberg-logical/src/lib.rs @@ -14,6 +14,7 @@ pub mod materializer; pub mod pipeline; +pub mod relation_event; pub mod runner; pub mod sink; diff --git a/crates/pg2iceberg-logical/src/materializer.rs b/crates/pg2iceberg-logical/src/materializer.rs index 6ce1cac..867ef1f 100644 --- a/crates/pg2iceberg-logical/src/materializer.rs +++ b/crates/pg2iceberg-logical/src/materializer.rs @@ -23,6 +23,7 @@ //! 7. `Coordinator::set_cursor` to the highest end_offset processed. //! 8. Update FileIndex with the new data file + removed PKs. +use crate::relation_event; use async_trait::async_trait; use bytes::Bytes; use pg2iceberg_coord::{Coordinator, LogEntry}; @@ -461,15 +462,6 @@ pub struct Materializer { /// controls how long a missed heartbeat survives before the /// worker drops out of the active list. distributed: Option, - /// PG → Iceberg ident translation, mirroring `Pipeline`'s map. - /// `apply_relation` is called with the PG-side ident (the - /// pgoutput stream's view), but `tables` is keyed by the - /// Iceberg-side ident the lifecycle registered with. Without - /// this translation, schema-evolution messages - /// (`ALTER TABLE ... ADD COLUMN`, type promotions, drops) would - /// silently no-op, and post-ALTER inserts would write to the - /// pre-ALTER iceberg schema. - table_translation: BTreeMap, } /// Distributed-mode parameters. Built by @@ -589,19 +581,6 @@ impl Materializer { meta_marker: None, meta_recorder: None, distributed: None, - table_translation: BTreeMap::new(), - } - } - - /// Register a translation from the PG-side ident (carried in - /// pgoutput Relation messages) to the Iceberg-side ident this - /// materializer's `tables` map is keyed by. Lets schema- - /// evolution events from the WAL stream find their target - /// table even when `sink.namespace` differs from the PG schema. - /// No-op when the two idents are equal. - pub fn register_table_translation(&mut self, pg_ident: TableIdent, iceberg_ident: TableIdent) { - if pg_ident != iceberg_ident { - self.table_translation.insert(pg_ident, iceberg_ident); } } @@ -843,7 +822,7 @@ impl Materializer { // land under the wrong Iceberg column. Nor does discovery's view // of the source's *current* columns apply yet: the stream may // still replay transactions from before a change. Its Relation - // messages evolve the schema in order ([`Self::apply_relation`]). + // messages, staged in the log, evolve the schema in order. let schema = match self.catalog.load_table(&ident).await? { None => { self.catalog.create_table(&schema).await?; @@ -931,71 +910,81 @@ impl Materializer { } } - /// Apply a pgoutput Relation message: bring the table's Iceberg - /// schema in line with `incoming_columns` ([`reconcile_columns`]), - /// call `Catalog::evolve_schema`, and rebuild the in-memory - /// `TableEntry::schema` and `TableWriter` from the result, so - /// subsequent materialize cycles encode rows with the new shape. + /// Apply a staged relation event: bring the table's Iceberg schema + /// in line with the source columns it carries ([`reconcile_columns`]), + /// and rebuild the in-memory `TableEntry::schema` and `TableWriter` + /// from the result, so the rows after it encode with the new shape. /// /// - **AddColumn** for a new name, with a fresh field id. - /// - **DropColumn** (soft) for a non-PK column `incoming` lacks: it + /// - **DropColumn** (soft) for a non-PK column the source lacks: it /// stays, nullable, so older data files keep resolving. /// - **PromoteColumnType** for a legal Iceberg promotion (int→long, /// float→double, decimal precision increase). Illegal type /// changes (e.g. long→int, text→int) are rejected with - /// `MaterializerError::Catalog` so the lifecycle fails loudly - /// rather than silently truncate or coerce downstream readers. + /// `MaterializerError::Catalog` so the table fails loudly rather + /// than silently truncate or coerce downstream readers. /// - **RenameColumn + AddColumn** for a column the source dropped /// and re-added: Postgres gives it no values, so it's a new /// column, and the dropped one keeps its values under another /// name. /// - /// The incoming `is_primary_key`/`nullable` are ignored: - /// pgoutput's key flag means "part of REPLICA IDENTITY" (every - /// column under `REPLICA IDENTITY FULL`), not "is primary key". - /// - /// **No-op when the table isn't registered** (the lifecycle - /// only registers tables in YAML — incoming Relations for - /// untracked tables, e.g. `_pg2iceberg.markers`, are silently - /// skipped). - /// - /// **No-op when columns match.** pgoutput re-emits Relation - /// messages liberally (e.g. on every cache invalidation, even - /// for unchanged schemas); the diff just produces an empty - /// change list and we return without touching the catalog. - pub async fn apply_relation( + /// Reconciles against the catalog's schema, not the in-memory one: + /// the event may already have been applied — before a crash, whose + /// retry re-reads it, or by another worker — and applying it again + /// is then a no-op. + async fn apply_columns( &mut self, ident: &TableIdent, - incoming_columns: &[pg2iceberg_pg::RelationColumn], + columns: &relation_event::Columns, ) -> Result<()> { - // Pgoutput emits Relation messages keyed by PG schema + - // table name; our `tables` map is keyed by Iceberg-side - // ident. Translate when registered, otherwise fall through. - let lookup_ident = self.table_translation.get(ident).unwrap_or(ident).clone(); - let entry = match self.tables.get_mut(&lookup_ident) { - Some(e) => e, - None => return Ok(()), - }; - let incoming: Vec<(String, pg2iceberg_core::IcebergType)> = incoming_columns - .iter() - .map(|c| (c.name.clone(), c.ty)) - .collect(); - let changes = reconcile_columns(&entry.schema, &incoming, &entry.dropped) + let current = self + .catalog + .load_table(ident) + .await? + .ok_or_else(|| MaterializerError::UnknownTable(ident.clone()))? + .schema; + let entry = self + .tables + .get_mut(ident) + .ok_or_else(|| MaterializerError::UnknownTable(ident.clone()))?; + let changes = reconcile_columns(¤t, columns, &entry.dropped) .map_err(MaterializerError::Catalog)?; - if !changes.is_empty() { - // Catalog-side first so a failure leaves the materializer's - // in-memory state unchanged (the next Relation re-triggers - // the diff). Use the Iceberg-side ident (`lookup_ident`) — - // `ident` is the PG-side key from the pgoutput Relation - // message, which the catalog wouldn't recognise. - self.catalog - .evolve_schema(&lookup_ident, changes.clone()) - .await?; - pg2iceberg_iceberg::apply_schema_changes(&mut entry.schema, &changes) - .map_err(MaterializerError::Catalog)?; - entry.writer = TableWriter::new(entry.schema.clone()); + let current = if changes.is_empty() { + current + } else { + self.catalog.evolve_schema(ident, changes).await?.schema + }; + entry.schema.columns = current.columns; + entry.writer = TableWriter::new(entry.schema.clone()); + entry.dropped = absent_from(&entry.schema, columns); + Ok(()) + } + + /// Buffer a log entry's events into `unit`, applying its relation + /// events in place: rows before one are staged under the old schema + /// first. + async fn buffer_entry( + &mut self, + ident: &TableIdent, + unit: &mut Unit, + events: Vec, + ) -> Result<()> { + for evt in events { + if evt.op != Op::Relation { + unit.buf.push(evt); + continue; + } + if !unit.buf.is_empty() { + self.prepare_step(ident, unit).await?; + } + let columns = relation_event::decode(&evt.row).ok_or_else(|| { + MaterializerError::Blob(StreamError::Decode(format!( + "malformed relation event for {ident}: {:?}", + evt.row + ))) + })?; + self.apply_columns(ident, &columns).await?; } - entry.dropped = absent_from(&entry.schema, &incoming); Ok(()) } @@ -1404,6 +1393,19 @@ impl Materializer { for e in entries { let bytes = self.blob_store.get(&e.s3_path).await?; let events = decode_chunk(&bytes)?; + // A schema change starts its entry (the pipeline stages it + // so). Commit the rows before it — written under the old + // schema — before applying it: a crash in between then + // re-reads this entry, and applying the change again is a + // no-op. Within a transaction (DDL mid-transaction) the + // rows so far can only be staged as a step. + if events.first().is_some_and(|e| e.op == Op::Relation) { + if is_tx_boundary(unit.buf.last(), events.first()) { + folded += self.commit_unit(ident, &mut unit).await?; + } else { + self.prepare_step(ident, &mut unit).await?; + } + } if !unit.buf.is_empty() && unit.buf.len() + events.len() > self.batch_rows { if is_tx_boundary(unit.buf.last(), events.first()) { folded += self.commit_unit(ident, &mut unit).await?; @@ -1418,7 +1420,7 @@ impl Materializer { } } events_read += events.len(); - unit.buf.extend(events); + self.buffer_entry(ident, &mut unit, events).await?; unit.end_offset = Some(e.end_offset); after = e.end_offset; } diff --git a/crates/pg2iceberg-logical/src/pipeline.rs b/crates/pg2iceberg-logical/src/pipeline.rs index 38c0e31..6ae8d97 100644 --- a/crates/pg2iceberg-logical/src/pipeline.rs +++ b/crates/pg2iceberg-logical/src/pipeline.rs @@ -4,13 +4,16 @@ //! `flush()` drains, uploads, claims, and advances `flushedLSN` — the last //! step gated by [`CoordCommitReceipt`]. +use crate::relation_event; use crate::sink::{FlushOutput, Sink, SinkError, TableChunk}; use async_trait::async_trait; use pg2iceberg_coord::{ CommitBatch, CoordCommitReceipt, CoordError, Coordinator, MarkerInfo, OffsetClaim, }; use pg2iceberg_core::metrics::{names, Labels}; -use pg2iceberg_core::{ColumnName, Lsn, Metrics, NoopMetrics, Op, PgValue, TableIdent}; +use pg2iceberg_core::{ + ChangeEvent, ColumnName, Lsn, Metrics, NoopMetrics, Op, PgValue, TableIdent, Timestamp, +}; use pg2iceberg_pg::DecodedMessage; use pg2iceberg_stream::codec::EncodedChunk; use pg2iceberg_stream::{BlobStore, StreamError}; @@ -149,6 +152,10 @@ pub struct Pipeline { /// This pipeline consumes the replication stream (see /// [`Self::track_replication`]). replication: bool, + /// The open transaction: xid and `Begin`'s LSN. + open_tx: Option<(u32, Lsn)>, + /// Each table's columns as last staged (see [`Self::stage_relation`]). + relations: BTreeMap, } impl Pipeline { @@ -192,6 +199,8 @@ impl Pipeline { primary_keys: BTreeMap::new(), table_translation: BTreeMap::new(), replication: false, + open_tx: None, + relations: BTreeMap::new(), } } @@ -262,8 +271,12 @@ impl Pipeline { return Ok(()); } match msg { - DecodedMessage::Begin { xid, .. } => self.sink.begin_tx(xid), + DecodedMessage::Begin { xid, final_lsn } => { + self.open_tx = Some((xid, final_lsn)); + self.sink.begin_tx(xid); + } DecodedMessage::Commit { xid, commit_lsn } => { + self.open_tx = None; // Drain any markers observed in this tx before // committing. Flushed atomically with the rest of // the tx via the next claim_offsets call. @@ -386,10 +399,8 @@ impl Pipeline { self.sink.record_change(evt)?; self.spill_if_full(xid).await?; } - DecodedMessage::Relation { .. } => { - // Schema evolution is applied via - // `Materializer::apply_relation`, which the lifecycle - // calls *before* this dispatch. Nothing to do here. + DecodedMessage::Relation { ident, columns } => { + self.stage_relation(ident, &columns)?; } DecodedMessage::Keepalive { wal_end, .. } => { // Only trustworthy between transactions: pgoutput sends a @@ -404,6 +415,47 @@ impl Pipeline { Ok(()) } + /// Stage the table's columns, if they changed, as a relation event + /// where the stream put the Relation message — before the + /// transaction's changes to the table — so the materializer applies + /// the schema change between the rows staged before and after it. + /// Applied when the message arrived instead, ahead of rows already + /// staged under the old schema, a dropped and re-added column would + /// take those rows' values for the dropped one. + fn stage_relation( + &mut self, + ident: TableIdent, + columns: &[pg2iceberg_pg::RelationColumn], + ) -> Result<()> { + let table = self.table_translation.get(&ident).cloned().unwrap_or(ident); + if self.markers_table.as_ref() == Some(&table) { + return Ok(()); + } + let columns: relation_event::Columns = + columns.iter().map(|c| (c.name.clone(), c.ty)).collect(); + // pgoutput resends a table's Relation, unchanged, at every new + // session and cache invalidation; those say nothing new. + if self.relations.get(&table) == Some(&columns) { + return Ok(()); + } + let (xid, lsn) = match self.open_tx { + Some((xid, lsn)) => (Some(xid), lsn), + None => (None, Lsn::ZERO), + }; + self.sink.record_change(ChangeEvent { + table: table.clone(), + op: Op::Relation, + lsn, + commit_ts: Timestamp(0), + xid, + before: None, + after: Some(relation_event::encode(&columns)), + unchanged_cols: Vec::new(), + })?; + self.relations.insert(table, columns); + Ok(()) + } + /// Drain all committed-but-unflushed transactions: encode → upload → /// `claim_offsets` → advance `flushedLSN` via the receipt. No-op if /// nothing is ready (no staged events, no markers awaiting emission, @@ -520,11 +572,8 @@ impl Pipeline { self.blob_store.put(&path, chunk.bytes).await?; let mut labels = Labels::new(); labels.insert("table".into(), table.name.clone()); - self.metrics.counter( - names::PIPELINE_ROWS_STAGED_TOTAL, - &labels, - chunk.record_count, - ); + self.metrics + .counter(names::PIPELINE_ROWS_STAGED_TOTAL, &labels, chunk.row_count); Ok(OffsetClaim { table, record_count: chunk.record_count, diff --git a/crates/pg2iceberg-logical/src/relation_event.rs b/crates/pg2iceberg-logical/src/relation_event.rs new file mode 100644 index 0000000..3912713 --- /dev/null +++ b/crates/pg2iceberg-logical/src/relation_event.rs @@ -0,0 +1,43 @@ +//! A source table's columns, as staged in the change log. +//! +//! The pipeline stages each pgoutput Relation message that changes a +//! table's columns as an `Op::Relation` event, in stream order, so the +//! materializer applies schema changes between the rows staged before +//! and after them. Its row carries the column list under one key, in the +//! source's column order (which re-add detection reads — see +//! [`pg2iceberg_iceberg::reconcile_columns`]). + +use pg2iceberg_core::{ColumnName, IcebergType, PgValue, Row}; + +const COLUMNS: &str = "columns"; + +/// A table's columns: name and type, in the source's order. +pub type Columns = Vec<(String, IcebergType)>; + +pub fn encode(columns: &Columns) -> Row { + let json = serde_json::to_string(columns).expect("columns serialize"); + Row::from([(ColumnName(COLUMNS.into()), PgValue::Json(json))]) +} + +/// `None` if `row` isn't an encoded column list. +pub fn decode(row: &Row) -> Option { + match row.get(&ColumnName(COLUMNS.into()))? { + PgValue::Json(json) => serde_json::from_str(json).ok(), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn columns_round_trip_in_order() { + let cols: Columns = vec![ + ("id".into(), IcebergType::Int), + ("qty".into(), IcebergType::Long), + ("note".into(), IcebergType::String), + ]; + assert_eq!(decode(&encode(&cols)), Some(cols)); + } +} diff --git a/crates/pg2iceberg-logical/src/sink.rs b/crates/pg2iceberg-logical/src/sink.rs index e6ae939..6f706a2 100644 --- a/crates/pg2iceberg-logical/src/sink.rs +++ b/crates/pg2iceberg-logical/src/sink.rs @@ -56,6 +56,9 @@ pub struct TableChunk { pub struct Sink { /// Per-table rolling writers, lazily created on first event. table_writers: BTreeMap, + /// Chunks closed early, each because a relation event starts a new + /// one ([`Self::append_event`]); the next flush emits them first. + closed: Vec, /// In-flight transactions keyed by xid. open_txns: BTreeMap, /// Committed-but-unflushed transactions, in commit order. @@ -69,6 +72,7 @@ impl Sink { assert!(flush_threshold > 0); Self { table_writers: BTreeMap::new(), + closed: Vec::new(), open_txns: BTreeMap::new(), committed: VecDeque::new(), flush_threshold, @@ -101,12 +105,15 @@ impl Sink { match evt.xid { Some(xid) if self.open_txns.contains_key(&xid) => { let tx = self.open_txns.get_mut(&xid).expect("checked above"); - tx.commit_ts = Some(evt.commit_ts); + // A relation event carries no commit timestamp. + if evt.op != Op::Relation { + tx.commit_ts = Some(evt.commit_ts); + } tx.events.push(evt); } _ => { // No tx: e.g. a snapshot-phase event. Stage immediately. - self.append_event(evt); + self.append_event(evt)?; } } Ok(()) @@ -149,15 +156,25 @@ impl Sink { for evt in tx.events.drain(..) { by_table.entry(evt.table.clone()).or_default().push(evt); } - by_table - .into_iter() - .map(|(table, events)| { - Ok(TableChunk { - table, - chunk: encode_chunk(&events)?, - }) - }) - .collect() + let mut out = Vec::new(); + for (table, events) in by_table { + // As in `append_event`: a relation event starts a chunk. + let mut part: Vec = Vec::new(); + for evt in events { + if evt.op == Op::Relation && !part.is_empty() { + out.push(TableChunk { + table: table.clone(), + chunk: encode_chunk(&std::mem::take(&mut part))?, + }); + } + part.push(evt); + } + out.push(TableChunk { + table, + chunk: encode_chunk(&part)?, + }); + } + Ok(out) } /// Change events currently held in memory: open transactions, @@ -189,12 +206,13 @@ impl Sink { } } for evt in tx.events { - self.append_event(evt); + self.append_event(evt)?; } } - // Flush every writer that has buffered rows. - let mut chunks = Vec::new(); + // Flush every writer that has buffered rows, after the chunks + // closed early (which hold earlier events). + let mut chunks = std::mem::take(&mut self.closed); for (table, writer) in self.table_writers.iter_mut() { if let Some(chunk) = writer.flush()? { chunks.push(TableChunk { @@ -210,13 +228,22 @@ impl Sink { })) } - fn append_event(&mut self, evt: ChangeEvent) { + /// A relation event (a schema change) starts a new chunk: the + /// materializer applies it at a log entry's start, once the entries + /// before it — rows written under the old schema — are committed. + fn append_event(&mut self, evt: ChangeEvent) -> Result<()> { let table = evt.table.clone(); let writer = self .table_writers - .entry(table) + .entry(table.clone()) .or_insert_with(|| RollingWriter::new(self.flush_threshold)); + if evt.op == Op::Relation { + if let Some(chunk) = writer.flush()? { + self.closed.push(TableChunk { table, chunk }); + } + } writer.append(evt); + Ok(()) } } @@ -239,6 +266,54 @@ mod tests { r } + fn relation(xid: u32, lsn: u64) -> ChangeEvent { + ChangeEvent { + op: Op::Relation, + after: Some(crate::relation_event::encode(&vec![])), + ..insert(xid, lsn, 0) + } + } + + /// Each chunk's ops, decoded. + fn ops(chunks: &[TableChunk]) -> Vec> { + chunks + .iter() + .map(|c| { + pg2iceberg_stream::codec::decode_chunk(&c.chunk.bytes) + .unwrap() + .into_iter() + .map(|e| e.op) + .collect() + }) + .collect() + } + + #[test] + fn a_relation_event_starts_a_chunk() { + let mut sink = Sink::new(100); + for (xid, lsn) in [(1, 10), (2, 20)] { + sink.begin_tx(xid); + sink.record_change(relation(xid, lsn)).unwrap(); + sink.record_change(insert(xid, lsn, xid as i32)).unwrap(); + sink.commit_tx(xid, Lsn(lsn)); + } + let out = sink.flush().unwrap().unwrap(); + assert_eq!( + ops(&out.chunks), + [[Op::Relation, Op::Insert], [Op::Relation, Op::Insert]] + ); + + // Likewise when an open transaction is spilled. + sink.begin_tx(3); + sink.record_change(insert(3, 30, 3)).unwrap(); + sink.record_change(relation(3, 30)).unwrap(); + sink.record_change(insert(3, 30, 4)).unwrap(); + assert_eq!( + ops(&sink.spill_open_tx(3).unwrap()), + vec![vec![Op::Insert], vec![Op::Relation, Op::Insert]] + ); + } + fn insert(xid: u32, lsn: u64, id: i32) -> ChangeEvent { ChangeEvent { table: ident(), diff --git a/crates/pg2iceberg-logical/tests/end_to_end.rs b/crates/pg2iceberg-logical/tests/end_to_end.rs index a899298..5239d1c 100644 --- a/crates/pg2iceberg-logical/tests/end_to_end.rs +++ b/crates/pg2iceberg-logical/tests/end_to_end.rs @@ -141,12 +141,13 @@ fn flushed_lsn_covers_the_commit_after_single_tx() { assert_eq!(advanced, Some(wal_end)); assert_eq!(h.pipeline.flushed_lsn(), wal_end); - // Coord state: one log_index row for "orders" at offset [0, 1). + // Coord state: one log_index row for "orders" at offset [0, 2): the + // table's columns (the session's first Relation message), then the row. let entries = block_on(h.coord.read_log(&ident("orders"), 0, 100)).unwrap(); assert_eq!(entries.len(), 1); assert_eq!(entries[0].start_offset, 0); - assert_eq!(entries[0].end_offset, 1); - assert_eq!(entries[0].record_count, 1); + assert_eq!(entries[0].end_offset, 2); + assert_eq!(entries[0].record_count, 2); // Blob store has the staged file; it round-trips via the codec. let paths = h.blob_store.paths(); @@ -154,9 +155,10 @@ fn flushed_lsn_covers_the_commit_after_single_tx() { assert_eq!(entries[0].s3_path, paths[0]); let bytes = block_on(h.blob_store.get(&paths[0])).unwrap(); let mat = decode_chunk(&bytes).unwrap(); - assert_eq!(mat.len(), 1); - assert_eq!(mat[0].op, Op::Insert); - assert_eq!(mat[0].lsn, change_lsn_before_commit(commit_lsn)); + assert_eq!(mat.len(), 2); + assert_eq!(mat[0].op, Op::Relation); + assert_eq!(mat[1].op, Op::Insert); + assert_eq!(mat[1].lsn, change_lsn_before_commit(commit_lsn)); } /// In a single-row tx the WAL is `Begin / Insert / Commit` with consecutive @@ -187,10 +189,10 @@ fn flushed_lsn_advances_to_max_across_multiple_txns() { assert_eq!(advanced, Some(wal_end)); assert_eq!(h.pipeline.flushed_lsn(), wal_end); - // One chunk covers all three rows: [0, 3). + // One chunk covers the table's columns and all three rows: [0, 4). let entries = block_on(h.coord.read_log(&ident("orders"), 0, 100)).unwrap(); assert_eq!(entries.len(), 1); - assert_eq!((entries[0].start_offset, entries[0].end_offset), (0, 3)); + assert_eq!((entries[0].start_offset, entries[0].end_offset), (0, 4)); } #[test] @@ -263,9 +265,11 @@ fn second_flush_after_more_txns_appends_offsets_contiguously() { block_on(h.pipeline.flush()).unwrap(); let entries = block_on(h.coord.read_log(&ident("orders"), 0, 100)).unwrap(); + // The first holds the table's columns too; the second needs none + // (unchanged). assert_eq!(entries.len(), 2); - assert_eq!((entries[0].start_offset, entries[0].end_offset), (0, 1)); - assert_eq!((entries[1].start_offset, entries[1].end_offset), (1, 3)); + assert_eq!((entries[0].start_offset, entries[0].end_offset), (0, 2)); + assert_eq!((entries[1].start_offset, entries[1].end_offset), (2, 4)); assert!(h.pipeline.flushed_lsn() > last); } diff --git a/crates/pg2iceberg-logical/tests/materializer_e2e.rs b/crates/pg2iceberg-logical/tests/materializer_e2e.rs index 29dd96c..22b9c82 100644 --- a/crates/pg2iceberg-logical/tests/materializer_e2e.rs +++ b/crates/pg2iceberg-logical/tests/materializer_e2e.rs @@ -290,9 +290,10 @@ fn flushed_lsn_and_cursor_advance_independently_but_correctly() { run_materializer(&mut h); - // Cursor now points past the only log entry. + // Cursor now points past the only log entry: the table's columns + // and the row. let cur_after = block_on(h.coord.get_cursor("default", &ident())).unwrap(); - assert_eq!(cur_after, Some(1)); + assert_eq!(cur_after, Some(2)); } #[test] diff --git a/crates/pg2iceberg-pg/src/lib.rs b/crates/pg2iceberg-pg/src/lib.rs index b42a834..84ed3b9 100644 --- a/crates/pg2iceberg-pg/src/lib.rs +++ b/crates/pg2iceberg-pg/src/lib.rs @@ -125,12 +125,11 @@ pub enum DecodedMessage { commit_lsn: Lsn, xid: u32, }, - /// Schema for a relation. Sent before any change events for that - /// relation, and re-sent after `ALTER TABLE` invalidates PG's - /// cache. The lifecycle compares incoming columns against the - /// materializer's registered schema; differences become - /// `SchemaChange::AddColumn` / `DropColumn` and trigger - /// `Catalog::evolve_schema`. + /// Schema for a relation. Sent before the first change to it in a + /// session or after its cache entry is invalidated (`ALTER TABLE`, + /// `TRUNCATE`, …). The pipeline stages a changed one in the log, in + /// order with the rows; the materializer applies it there + /// (`Catalog::evolve_schema`). Relation { ident: TableIdent, columns: Vec, diff --git a/crates/pg2iceberg-pg/src/prod/stream.rs b/crates/pg2iceberg-pg/src/prod/stream.rs index 6bdda59..5968f0f 100644 --- a/crates/pg2iceberg-pg/src/prod/stream.rs +++ b/crates/pg2iceberg-pg/src/prod/stream.rs @@ -286,7 +286,7 @@ impl ReaderState { // pgoutput doesn't carry nullability; we stamp // nullable=true on every non-PK column. PG only // permits adding non-nullable columns with a - // DEFAULT, and for that case `apply_relation` + // DEFAULT, and for that case the materializer // would still see a nullable add (which Iceberg // tolerates). let ty = pg2iceberg_core::map_pg_to_iceberg(pg_type) diff --git a/crates/pg2iceberg-sim/src/postgres.rs b/crates/pg2iceberg-sim/src/postgres.rs index 53a099b..f05d481 100644 --- a/crates/pg2iceberg-sim/src/postgres.rs +++ b/crates/pg2iceberg-sim/src/postgres.rs @@ -515,8 +515,8 @@ impl SimPostgres { /// table's schema and emits a fresh Relation WAL event so any /// active replication stream picks up the change. The column /// auto-allocates the next field id (matching Iceberg's - /// monotonic-only field-id rule). Used by DST to drive - /// `Materializer::apply_relation` end-to-end. + /// monotonic-only field-id rule). Used by DST to drive schema + /// evolution end-to-end. pub fn alter_add_column( &self, ident: &TableIdent, @@ -613,11 +613,11 @@ impl SimPostgres { /// Test hook: `ALTER TABLE … ALTER COLUMN … TYPE …`. Mutates the /// named column's `IcebergType` in place (preserving field id and /// nullability) and emits a fresh Relation event so the - /// materializer's `apply_relation` diff sees the type change. Used + /// materializer's schema diff sees the type change. Used /// by DST to drive the legal-promotion + illegal-narrowing paths. /// The sim doesn't validate the change itself (real PG would do its /// own type-cast checks); validation happens downstream in - /// `apply_relation` / `apply_schema_changes`. + /// `reconcile_columns` / `apply_schema_changes`. pub fn alter_column_type( &self, ident: &TableIdent, @@ -878,6 +878,14 @@ impl SimPostgres { } /// Snapshot of a table's rows in PK order, for tests / verify. + /// `ident`'s column names as of WAL position `at`. + pub fn columns_at(&self, ident: &TableIdent, at: Lsn) -> Vec { + let s = self.state.lock().unwrap(); + s.columns_at(ident, at) + .map(|cols| cols.iter().map(|c| c.name.clone()).collect()) + .unwrap_or_default() + } + /// Whether `ident`'s relation changed (DDL, or an invalidation) after /// its last row change. pgoutput tells a consumer about a schema /// change only with the table's next change, so it hasn't told one. diff --git a/crates/pg2iceberg-stream/src/codec.rs b/crates/pg2iceberg-stream/src/codec.rs index 9a64be2..b212ed8 100644 --- a/crates/pg2iceberg-stream/src/codec.rs +++ b/crates/pg2iceberg-stream/src/codec.rs @@ -84,9 +84,10 @@ fn op_to_str(op: Op) -> Option<&'static str> { // current FileIndex. Without staging, a TRUNCATE in PG would // silently leave Iceberg with the pre-truncate rows. Op::Truncate => Some("T"), - // Begin/Commit/Relation are control events consumed by the - // pipeline before reaching the writer. - Op::Relation => None, + // A source schema change, staged in order with the rows so the + // materializer applies it between the rows written before and + // after it. `_data` carries the table's columns. + Op::Relation => Some("R"), } } @@ -96,18 +97,19 @@ fn op_from_str(s: &str) -> Option { "U" => Some(Op::Update), "D" => Some(Op::Delete), "T" => Some(Op::Truncate), + "R" => Some(Op::Relation), _ => None, } } -/// Pick the row that gets staged for a given event: `after` for insert/update, -/// `before` for delete. Returns `None` for non-DML ops, which the caller -/// should filter before calling. +/// Pick the row that gets staged for a given event: `after` for +/// insert/update (and a relation's column list), `before` for delete. +/// `None` for a truncate, which carries no row. fn staged_row(evt: &ChangeEvent) -> Option<&Row> { match evt.op { - Op::Insert | Op::Update => evt.after.as_ref(), + Op::Insert | Op::Update | Op::Relation => evt.after.as_ref(), Op::Delete => evt.before.as_ref(), - Op::Relation | Op::Truncate => None, + Op::Truncate => None, } } @@ -206,7 +208,10 @@ pub fn encode_batch(events: &[ChangeEvent]) -> Result { #[derive(Clone, Debug, PartialEq, Eq)] pub struct EncodedChunk { pub bytes: Bytes, + /// Records staged — every event, so log offsets count them all. pub record_count: u64, + /// Row changes among them (records less schema events). + pub row_count: u64, pub max_lsn: Lsn, } @@ -219,6 +224,10 @@ pub fn encode_chunk(events: &[ChangeEvent]) -> Result { .max() .unwrap_or(Lsn::ZERO); let record_count = batch.num_rows() as u64; + let row_count = events + .iter() + .filter(|e| op_to_str(e.op).is_some() && e.op != Op::Relation) + .count() as u64; let mut buf = Vec::::new(); let props = WriterProperties::builder().build(); @@ -236,6 +245,7 @@ pub fn encode_chunk(events: &[ChangeEvent]) -> Result { Ok(EncodedChunk { bytes: Bytes::from(buf), record_count, + row_count, max_lsn, }) } @@ -462,32 +472,29 @@ mod tests { } #[test] - fn non_dml_events_are_skipped() { - let evts = vec![ - ChangeEvent { - table: ident(), - op: Op::Relation, - lsn: Lsn(1), - commit_ts: Timestamp(0), - xid: None, - before: None, - after: None, - unchanged_cols: vec![], - }, - ChangeEvent { - table: ident(), - op: Op::Insert, - lsn: Lsn(2), - commit_ts: Timestamp(0), - xid: None, - before: None, - after: Some(row(&[("id", PgValue::Int4(1))])), - unchanged_cols: vec![], - }, - ]; - let chunk = encode_chunk(&evts).unwrap(); - assert_eq!(chunk.record_count, 1); - assert_eq!(chunk.max_lsn, Lsn(2)); + fn relation_events_round_trip_in_order_with_rows() { + let relation = ChangeEvent { + table: ident(), + op: Op::Relation, + lsn: Lsn(1), + commit_ts: Timestamp(0), + xid: Some(7), + before: None, + after: Some(row(&[("columns", PgValue::Json("[]".into()))])), + unchanged_cols: vec![], + }; + let insert = ChangeEvent { + op: Op::Insert, + lsn: Lsn(2), + after: Some(row(&[("id", PgValue::Int4(1))])), + ..relation.clone() + }; + let chunk = encode_chunk(&[relation.clone(), insert]).unwrap(); + assert_eq!(chunk.record_count, 2); + let decoded = decode_chunk(&chunk.bytes).unwrap(); + assert_eq!(decoded[0].op, Op::Relation); + assert_eq!(Some(&decoded[0].row), relation.after.as_ref()); + assert_eq!(decoded[1].op, Op::Insert); } fn update(id: i32, before: Option) -> ChangeEvent { diff --git a/crates/pg2iceberg-tests/tests/dst.rs b/crates/pg2iceberg-tests/tests/dst.rs index 21a3bdd..3640b2f 100644 --- a/crates/pg2iceberg-tests/tests/dst.rs +++ b/crates/pg2iceberg-tests/tests/dst.rs @@ -207,17 +207,26 @@ fn discovered_schema(db: &SimPostgres) -> TableSchema { /// a row lacks reads as NULL): Iceberg keeps dropped columns, and a row /// written before a column existed doesn't carry it. fn on_source_columns(db: &SimPostgres, rows: Vec) -> Vec { - let cols: Vec = db + let cols: Vec = db .table_schema(&ident()) .expect("source table") .columns .into_iter() - .map(|c| ColumnName(c.name)) + .map(|c| c.name) .collect(); + on_columns(&cols, rows) +} + +/// `rows` on just `cols`, NULL where a row lacks one. +fn on_columns(cols: &[String], rows: Vec) -> Vec { rows.into_iter() .map(|r| { cols.iter() - .map(|c| (c.clone(), r.get(c).cloned().unwrap_or(PgValue::Null))) + .map(|c| { + let c = ColumnName(c.clone()); + let v = r.get(&c).cloned().unwrap_or(PgValue::Null); + (c, v) + }) .collect() }) .collect() @@ -283,7 +292,6 @@ fn published() -> Vec { fn register_other(m: &mut Materializer) { if SECOND_TABLE.get() > 0 { block_on(m.register_table(other_schema())).unwrap(); - m.register_table_translation(other_pg_ident(), other_ident()); } } @@ -918,6 +926,7 @@ impl DstHarness { db: db.clone(), violations: Mutex::new(Vec::new()), fail_next_commit: Default::default(), + fail_commit_after_schema_change: Default::default(), audit_paused: Default::default(), lose_next_response: Default::default(), }); @@ -1027,6 +1036,7 @@ impl DstHarness { db: db.clone(), violations: Mutex::new(Vec::new()), fail_next_commit: Default::default(), + fail_commit_after_schema_change: Default::default(), audit_paused: Default::default(), lose_next_response: Default::default(), }); @@ -1114,12 +1124,9 @@ impl DstHarness { self.stream.recv() } - /// What the lifecycle does with a message: a schema change reaches - /// the materializer before the pipeline. + /// What the lifecycle does with a message: hand it to the pipeline, + /// which stages schema changes in order with the rows. fn process(&mut self, msg: DecodedMessage) { - if let DecodedMessage::Relation { ident, columns } = &msg { - block_on(self.materializer.apply_relation(ident, columns)).unwrap(); - } block_on(self.pipeline.process(msg)).unwrap(); } @@ -1483,6 +1490,12 @@ impl DstHarness { } Step::AddNote => { if !NOTE_PRESENT.get() { + // A column re-added in last place can't be told from + // one never dropped unless pg2iceberg saw the drop, + // which pgoutput reports only with the table's next + // change: make it first. (Without it — no change + // between drop and re-add — the re-add is invisible.) + write_after_schema_change(self); let col = schema().columns.into_iter().find(|c| c.name == "note"); self.db.alter_add_column(&ident(), col.unwrap()).unwrap(); NOTE_PRESENT.set(true); @@ -1697,10 +1710,24 @@ async fn atomic_visibility(storage: &Storage, db: &SimPostgres) -> Result<(), St _ => i32::MAX, }; let mut state: BTreeMap = BTreeMap::new(); - let mut boundaries: Vec> = vec![Vec::new()]; + // Each boundary's rows, with the source's columns then. + let mut boundaries: Vec<(Vec, Vec)> = vec![(Vec::new(), Vec::new())]; let mut i = 0; while i < events.len() { let xid = events[i].xid; + // A column dropped since the last transaction takes its values + // with it: one re-added later starts out NULL. Iceberg can show + // that state too, once it has applied the drop. + let columns = db.columns_at(&ident(), events[i].lsn); + let dropped = state + .values() + .any(|r| r.keys().any(|c| !columns.contains(&c.0))); + if dropped { + for row in state.values_mut() { + row.retain(|c, _| columns.contains(&c.0)); + } + boundaries.push((state.values().cloned().collect(), columns.clone())); + } while i < events.len() && events[i].xid == xid { let e = &events[i]; match e.op { @@ -1728,16 +1755,16 @@ async fn atomic_visibility(storage: &Storage, db: &SimPostgres) -> Result<(), St } i += 1; } - boundaries.push(state.values().cloned().collect()); + boundaries.push((state.values().cloned().collect(), columns)); } - // Compare on the source's current columns: Iceberg keeps dropped - // ones, and older boundaries predate added ones. - let iceberg = on_source_columns(db, iceberg); - let boundaries: Vec> = boundaries + // Compare on the source's columns at each boundary: Iceberg keeps + // dropped columns, and learns of a schema change only with the + // table's next change (a column dropped then re-added still reads + // as the dropped one until then). + let matches = boundaries .into_iter() - .map(|b| on_source_columns(db, b)) - .collect(); - if !boundaries.contains(&iceberg) { + .any(|(rows, cols)| on_columns(&cols, iceberg.clone()) == on_columns(&cols, rows)); + if !matches { return Err(format!( "invariant 10 (atomic visibility): Iceberg state matches no transaction boundary: {iceberg:?}" )); @@ -1806,6 +1833,9 @@ struct AuditedCatalog { violations: Mutex>, /// When set, the next multi-step commit fails without committing. fail_next_commit: std::sync::atomic::AtomicBool, + /// When set, the commit after the next schema change fails without + /// committing — a crash between the two. + fail_commit_after_schema_change: std::sync::atomic::AtomicBool, /// When set, commits aren't audited — for tests that count blob reads /// (an audit reads the whole table). audit_paused: std::sync::atomic::AtomicBool, @@ -1902,7 +1932,15 @@ impl Catalog for AuditedCatalog { ident: &TableIdent, changes: Vec, ) -> pg2iceberg_iceberg::Result { - self.inner.evolve_schema(ident, changes).await + let meta = self.inner.evolve_schema(ident, changes).await?; + if self + .fail_commit_after_schema_change + .swap(false, std::sync::atomic::Ordering::SeqCst) + { + self.fail_next_commit + .store(true, std::sync::atomic::Ordering::SeqCst); + } + Ok(meta) } async fn expire_snapshots( &self, @@ -2053,7 +2091,8 @@ fn check_invariants(h: &mut DstHarness) -> Result<(), String> { last = Some((xid, commit)); } let mut staged: BTreeMap> = BTreeMap::new(); - for e in staged_events { + // Staged schema changes aren't WAL changes. + for e in staged_events.into_iter().filter(|e| e.op != Op::Relation) { staged.entry(e.xid.unwrap_or(0)).or_default().push(e); } // What staging should hold for each WAL event: its row as sent, except @@ -2598,7 +2637,13 @@ fn pipeline_crash_mid_transaction_keeps_it_atomic() { #[test] fn file_index_stays_true_after_a_lost_compaction_response() { let mut h = DstHarness::boot(); - h.run_step(&Step::BigTx { inserts: 2, qty: 0 }); + // One data file holding a live row (2) and a deleted one (1): the + // pass rewrites it, moving row 2 to a new file. + h.run_step(&Step::Insert { id: 1, qty: 10 }); + h.run_step(&Step::Insert { id: 2, qty: 20 }); + h.run_step(&Step::DriveFlush); + h.run_step(&Step::MaterializerCycle); + h.run_step(&Step::Update { id: 1, qty: 11 }); h.run_step(&Step::DriveFlush); h.run_step(&Step::MaterializerCycle); h.run_step(&Step::LoseCommitResponse); @@ -2658,6 +2703,45 @@ fn rows_staged_before_a_re_add_keep_their_values_out_of_the_new_column() { check_invariants(&mut h).unwrap(); } +/// The materializer fails between applying a schema change and committing +/// the rows after it. The retry re-reads the log from the last commit, so +/// that commit must cover every row staged before the change: re-read +/// under the changed schema, a dropped and re-added column would take +/// their values. +#[test] +fn a_failed_commit_after_a_schema_change_retries_cleanly() { + let mut h = DstHarness::boot(); + h.run_step(&Step::Insert { id: 1, qty: 10 }); + // Drop and re-add `note` with no write between (unlike `AddNote`): + // seen by order, as it moves last, and staged in the same flush as + // row 1. + h.db.alter_drop_column(&ident(), "note").unwrap(); + let note = schema().columns.into_iter().find(|c| c.name == "note"); + h.db.alter_add_column(&ident(), note.unwrap()).unwrap(); + h.run_step(&Step::Insert { id: 2, qty: 20 }); + h.run_step(&Step::DriveFlush); + h.audited + .fail_commit_after_schema_change + .store(true, std::sync::atomic::Ordering::SeqCst); + let err = block_on(h.materializer.cycle()).unwrap_err(); + assert!(err.to_string().contains("injected"), "{err}"); + h.run_step(&Step::MaterializerCycle); + check_invariants(&mut h).unwrap(); + // Nor may the retry re-apply the older schema change on top of the + // newer one: that renames columns that were never re-added. + let schema = block_on(h.audited.load_table(&ident())) + .unwrap() + .unwrap() + .schema; + let columns: Vec<(String, i32)> = schema + .columns + .into_iter() + .map(|c| (c.name, c.field_id)) + .collect(); + let want = [("id", 1), ("note__dropped_2", 2), ("qty", 3), ("note", 4)]; + assert_eq!(columns, want.map(|(n, i)| (n.to_string(), i))); +} + /// Compaction and TOAST resolution read older data files, which hold the /// dropped column's values under the re-added column's name: they must /// read columns by field id, or those values come back. @@ -3011,6 +3095,16 @@ fn crash_after_some_inserts_then_more_inserts() { check_invariants(&mut h).unwrap(); } +/// Row changes staged for the main table (schema events aside). +fn staged_changes(h: &DstHarness) -> usize { + block_on(h.coord.read_log(&ident(), 0, 1_000_000)) + .unwrap() + .iter() + .flat_map(|e| decode_chunk(&block_on(h.blob_store.get(&e.s3_path)).unwrap()).unwrap()) + .filter(|e| e.op != Op::Relation) + .count() +} + #[test] fn rollback_does_not_appear_in_staged_or_coord() { let mut h = DstHarness::boot(); @@ -3018,10 +3112,11 @@ fn rollback_does_not_appear_in_staged_or_coord() { h.run_step(&Step::Insert { id: 1, qty: 10 }); h.run_step(&Step::DriveFlush); check_invariants(&mut h).unwrap(); - - let entries = block_on(h.coord.read_log(&ident(), 0, 100)).unwrap(); - let total: u64 = entries.iter().map(|e| e.record_count).sum(); - assert_eq!(total, 1, "only the committed insert should be staged"); + assert_eq!( + staged_changes(&h), + 1, + "only the committed insert should be staged" + ); } #[test] @@ -3032,10 +3127,7 @@ fn update_then_delete_round_trips_to_staged() { h.run_step(&Step::Delete { id: 1 }); h.run_step(&Step::DriveFlush); check_invariants(&mut h).unwrap(); - - let entries = block_on(h.coord.read_log(&ident(), 0, 100)).unwrap(); - let total: u64 = entries.iter().map(|e| e.record_count).sum(); - assert_eq!(total, 3, "I + U + D"); + assert_eq!(staged_changes(&h), 3, "I + U + D"); } #[test] diff --git a/crates/pg2iceberg-tests/tests/dst_decimal.rs b/crates/pg2iceberg-tests/tests/dst_decimal.rs index 0b87e84..6252bce 100644 --- a/crates/pg2iceberg-tests/tests/dst_decimal.rs +++ b/crates/pg2iceberg-tests/tests/dst_decimal.rs @@ -25,7 +25,6 @@ use pg2iceberg_core::{ use pg2iceberg_iceberg::read_materialized_state; use pg2iceberg_logical::materializer::{CounterMaterializerNamer, Materializer}; use pg2iceberg_logical::pipeline::{CounterBlobNamer, Pipeline}; -use pg2iceberg_pg::DecodedMessage; use pg2iceberg_sim::blob::MemoryBlobStore; use pg2iceberg_sim::catalog::MemoryCatalog; use pg2iceberg_sim::clock::TestClock; @@ -155,9 +154,6 @@ impl Harness { fn drive_then_materialize(&mut self) { while let Some(msg) = self.stream.recv() { - if let DecodedMessage::Relation { ident, columns } = &msg { - block_on(self.materializer.apply_relation(ident, columns)).unwrap(); - } block_on(self.pipeline.process(msg)).unwrap(); } block_on(self.pipeline.flush()).unwrap(); @@ -239,9 +235,6 @@ fn decimal_lossy_downscale_refuses_with_error() { tx.commit(Timestamp(0)).unwrap(); while let Some(msg) = h.stream.recv() { - if let DecodedMessage::Relation { ident, columns } = &msg { - block_on(h.materializer.apply_relation(ident, columns)).unwrap(); - } block_on(h.pipeline.process(msg)).unwrap(); } block_on(h.pipeline.flush()).unwrap(); diff --git a/crates/pg2iceberg-tests/tests/dst_dml_correctness.rs b/crates/pg2iceberg-tests/tests/dst_dml_correctness.rs index d6c7467..3ebd5c6 100644 --- a/crates/pg2iceberg-tests/tests/dst_dml_correctness.rs +++ b/crates/pg2iceberg-tests/tests/dst_dml_correctness.rs @@ -22,7 +22,6 @@ use pg2iceberg_core::{ use pg2iceberg_iceberg::read_materialized_state; use pg2iceberg_logical::materializer::{CounterMaterializerNamer, Materializer}; use pg2iceberg_logical::pipeline::{CounterBlobNamer, Pipeline}; -use pg2iceberg_pg::DecodedMessage; use pg2iceberg_sim::blob::MemoryBlobStore; use pg2iceberg_sim::catalog::MemoryCatalog; use pg2iceberg_sim::clock::TestClock; @@ -134,13 +133,9 @@ impl Harness { // Drain replication stream → pipeline → flush → materialize. // SimReplicationStream::recv is sync and returns None when // there are no more events to deliver. Mirrors the - // lifecycle's main loop: Relation messages drive schema - // evolution via `materializer.apply_relation` BEFORE the - // pipeline forwards them onward. + // lifecycle's main loop: every message goes to the pipeline, + // which stages schema changes in order with the rows. while let Some(msg) = self.stream.recv() { - if let DecodedMessage::Relation { ident, columns } = &msg { - block_on(self.materializer.apply_relation(ident, columns)).unwrap(); - } block_on(self.pipeline.process(msg)).unwrap(); } block_on(self.pipeline.flush()).unwrap(); @@ -275,7 +270,7 @@ fn update_with_pk_change_deletes_old_pk_in_iceberg() { fn alter_table_add_column_mid_stream_propagates_to_iceberg() { // Operator runs `ALTER TABLE orders ADD COLUMN note TEXT`, // followed by an INSERT carrying the new column. - // Goes through `materializer.apply_relation` → + // Goes through the staged relation event → // `Catalog::evolve_schema` → updated `entry.schema` + writer // rebuild. The post-evolve INSERT must persist the new // column's value in Iceberg. diff --git a/crates/pg2iceberg-tests/tests/dst_meta_tables.rs b/crates/pg2iceberg-tests/tests/dst_meta_tables.rs index fdf3268..9bf6a2a 100644 --- a/crates/pg2iceberg-tests/tests/dst_meta_tables.rs +++ b/crates/pg2iceberg-tests/tests/dst_meta_tables.rs @@ -34,7 +34,6 @@ use pg2iceberg_iceberg::meta::{ use pg2iceberg_iceberg::{read_data_file, Catalog}; use pg2iceberg_logical::materializer::{CounterMaterializerNamer, Materializer}; use pg2iceberg_logical::pipeline::{CounterBlobNamer, Pipeline}; -use pg2iceberg_pg::DecodedMessage; use pg2iceberg_sim::blob::MemoryBlobStore; use pg2iceberg_sim::catalog::MemoryCatalog; use pg2iceberg_sim::clock::TestClock; @@ -143,9 +142,6 @@ impl Harness { fn drive_then_materialize(&mut self) { while let Some(msg) = self.stream.recv() { - if let DecodedMessage::Relation { ident, columns } = &msg { - block_on(self.materializer.apply_relation(ident, columns)).unwrap(); - } block_on(self.pipeline.process(msg)).unwrap(); } block_on(self.pipeline.flush()).unwrap(); diff --git a/crates/pg2iceberg-tests/tests/dst_schema_evolution.rs b/crates/pg2iceberg-tests/tests/dst_schema_evolution.rs index 5ae7a1d..9dea478 100644 --- a/crates/pg2iceberg-tests/tests/dst_schema_evolution.rs +++ b/crates/pg2iceberg-tests/tests/dst_schema_evolution.rs @@ -1,7 +1,8 @@ //! DST coverage for the full surface of schema-evolution scenarios that -//! pgoutput can produce. Each test seeds a sim PG, drives ALTERs through -//! `materializer.apply_relation`, and asserts the Iceberg-side schema / -//! data state. +//! pgoutput can produce. Each test seeds a sim PG, runs ALTERs, streams +//! the resulting Relation messages through the pipeline (which stages +//! them in the log) and the materializer (which applies them in order), +//! and asserts the Iceberg-side schema / data state. //! //! Companion to `dst_dml_correctness.rs` (which covers DML correctness, //! including the basic ADD/DROP COLUMN happy paths). This file focuses @@ -21,7 +22,6 @@ use pg2iceberg_iceberg::{read_materialized_state, Catalog}; use pg2iceberg_logical::materializer::{CounterMaterializerNamer, Materializer}; use pg2iceberg_logical::pipeline::{CounterBlobNamer, Pipeline}; use pg2iceberg_logical::MaterializerError; -use pg2iceberg_pg::DecodedMessage; use pg2iceberg_sim::blob::MemoryBlobStore; use pg2iceberg_sim::catalog::MemoryCatalog; use pg2iceberg_sim::clock::TestClock; @@ -151,29 +151,21 @@ impl Harness { /// asserting every step succeeded. Use this for the happy path. fn drive_then_materialize(&mut self) { while let Some(msg) = self.stream.recv() { - if let DecodedMessage::Relation { ident, columns } = &msg { - block_on(self.materializer.apply_relation(ident, columns)).unwrap(); - } block_on(self.pipeline.process(msg)).unwrap(); } block_on(self.pipeline.flush()).unwrap(); block_on(self.materializer.cycle()).unwrap(); } - /// Drain the replication stream, but capture the first - /// `apply_relation` error and return it instead of panicking. - /// Used by the illegal-type-change tests to assert that the - /// lifecycle fails loudly rather than silently coercing. + /// Drain the replication stream and run a materialize cycle, + /// returning the cycle's error: an illegal schema change fails the + /// cycle that reaches it, rather than silently coercing. fn drive_capturing_relation_error(&mut self) -> Option { while let Some(msg) = self.stream.recv() { - if let DecodedMessage::Relation { ident, columns } = &msg { - if let Err(e) = block_on(self.materializer.apply_relation(ident, columns)) { - return Some(e); - } - } block_on(self.pipeline.process(msg)).unwrap(); } - None + block_on(self.pipeline.flush()).unwrap(); + block_on(self.materializer.cycle()).err() } /// Insert a row into `ident`. pgoutput sends a table's Relation only @@ -360,7 +352,7 @@ fn pk_type_change_rejected_even_for_legal_promotion() { fn single_relation_with_multiple_new_columns_adds_all() { // PG can run `ALTER TABLE … ADD COLUMN a TEXT, ADD COLUMN b INT` // in one statement, producing a single Relation message that - // adds two columns at once. apply_relation must emit two + // adds two columns at once. The materializer must emit two // AddColumn changes and apply them atomically (both succeed or // neither does). let s = schema_with("orders", "qty", IcebergType::Int); @@ -568,6 +560,68 @@ fn add_then_drop_with_no_write_between_never_reaches_iceberg() { assert!(evolved.columns.iter().all(|c| c.name != "tmp")); } +#[test] +fn rows_before_a_re_add_keep_their_values_out_of_the_new_column() { + // A row written, `note` dropped (seen: a write follows), re-added + // (seen: a write follows), all materialized in one cycle. The first + // row's value belongs to the dropped column, now renamed; the new + // `note` is empty for it. + let s = schema_with("orders", "note", IcebergType::String); + let mut h = Harness::boot(std::slice::from_ref(&s)); + h.write( + &s.ident, + &[ + ("id", PgValue::Int4(1)), + ("note", PgValue::Text("a".into())), + ], + ); + h.db.alter_drop_column(&s.ident, "note").unwrap(); + h.write(&s.ident, &[("id", PgValue::Int4(2))]); + h.db.alter_add_column( + &s.ident, + ColumnSchema { + name: "note".into(), + field_id: 0, + ty: IcebergType::String, + nullable: true, + is_primary_key: false, + }, + ) + .unwrap(); + h.write( + &s.ident, + &[ + ("id", PgValue::Int4(3)), + ("note", PgValue::Text("c".into())), + ], + ); + h.drive_then_materialize(); + + let schema = h.iceberg_schema(&s.ident); + let mut rows = block_on(read_materialized_state( + h.catalog.as_ref(), + h.blob.as_ref(), + &s.ident, + &schema, + &[col("id")], + )) + .unwrap(); + rows.sort_by_key(|r| format!("{:?}", r[&col("id")])); + let notes: Vec<(PgValue, PgValue)> = rows + .iter() + .map(|r| (r[&col("note__dropped_2")].clone(), r[&col("note")].clone())) + .collect(); + let text = |v: &str| PgValue::Text(v.into()); + assert_eq!( + notes, + vec![ + (text("a"), PgValue::Null), + (PgValue::Null, PgValue::Null), + (PgValue::Null, text("c")), + ] + ); +} + // ── Multi-table evolution ──────────────────────────────────────────── #[test] @@ -656,11 +710,10 @@ fn repeated_relation_with_same_schema_is_no_op() { #[test] fn add_column_then_insert_uses_new_column() { - // The lifecycle's main loop routes Relation → apply_relation - // before pipeline.process. So even if a Relation arrives - // immediately before an Insert (e.g. mid-transaction in real - // PG), the Iceberg schema is updated in time for the staged - // Insert to encode the new column. + // The pipeline stages a Relation in order with the rows, so a + // Relation arriving immediately before an Insert (e.g. + // mid-transaction in real PG) is applied in time for that Insert + // to encode the new column. let s = schema_with("orders", "qty", IcebergType::Int); let mut h = Harness::boot(std::slice::from_ref(&s)); h.drive_then_materialize(); diff --git a/crates/pg2iceberg-validate/src/runtime.rs b/crates/pg2iceberg-validate/src/runtime.rs index 61b4514..11740b1 100644 --- a/crates/pg2iceberg-validate/src/runtime.rs +++ b/crates/pg2iceberg-validate/src/runtime.rs @@ -494,7 +494,6 @@ where // applied here so schema-evolution Relation messages from // pgoutput find their target table when `sink.namespace` // differs from the PG schema. - materializer.register_table_translation(schema.pg_ident(), schema.ident.clone()); } if let Some(meta_ns) = &lc.meta_namespace { // Enable meta-marker emission. The materializer creates the @@ -906,20 +905,9 @@ where } res = loop_state.stream.recv() => { let msg = res.map_err(|e| MainLoopError::Recv(e.to_string()))?; - // Relation messages drive schema evolution: diff - // incoming columns against the materializer's - // registered schema, call `Catalog::evolve_schema` - // for any AddColumn / DropColumn, and update the - // materializer's in-memory schema + writer. Forward - // a payload-less Relation to the pipeline anyway so - // future hooks (metrics, etc.) keep firing. - if let pg2iceberg_pg::DecodedMessage::Relation { ident, columns } = &msg { - loop_state - .materializer - .apply_relation(ident, columns) - .await - .map_err(|e| MainLoopError::Catalog(e.to_string()))?; - } + // Relation messages too: the pipeline stages schema + // changes in order with the rows, and the materializer + // applies them there. loop_state .pipeline .process(msg)