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)