diff --git a/crates/adapters/src/integrated/postgres/prepared_statements.rs b/crates/adapters/src/integrated/postgres/prepared_statements.rs index a6435f122cd..6dc15ce5feb 100644 --- a/crates/adapters/src/integrated/postgres/prepared_statements.rs +++ b/crates/adapters/src/integrated/postgres/prepared_statements.rs @@ -22,13 +22,29 @@ impl RawQueries { .map(|f| f.name.sql_name()) .collect(); + // List supplied columns explicitly so Postgres applies DEFAULTs to any + // omitted column, instead of SELECT * forcing NULL into them (#6694). + let mut insert_columns: Vec = value_schema + .fields + .iter() + .map(|f| f.name.sql_name()) + .chain(config.extra_columns.iter().map(|k| format!(r#""{k}""#))) + .collect(); + + // CDC rows also carry the op/ts metadata columns. + if matches!(config.mode, PostgresWriteMode::Cdc) { + insert_columns.push(format!(r#""{}""#, config.cdc_op_column)); + insert_columns.push(format!(r#""{}""#, config.cdc_ts_column)); + } + let insert_columns = insert_columns.join(", "); + let mut raw_queries = RawQueries::default(); match config.mode { PostgresWriteMode::Cdc => { // In CDC mode, everything is an INSERT into the event log raw_queries.insert = format!( - r#"INSERT INTO "{table}" SELECT * FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb)"#, + r#"INSERT INTO "{table}" ({insert_columns}) SELECT {insert_columns} FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb)"#, ); // For CDC mode, upsert and delete operations also use INSERT raw_queries.upsert = raw_queries.insert.clone(); @@ -59,7 +75,7 @@ impl RawQueries { }; raw_queries.insert = format!( - r#"INSERT INTO "{table}" SELECT * FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb) ON CONFLICT {on_conflict}"#, + r#"INSERT INTO "{table}" ({insert_columns}) SELECT {insert_columns} FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb) ON CONFLICT {on_conflict}"#, ); } } @@ -169,3 +185,85 @@ impl PreparedStatements { }) } } + +#[cfg(test)] +mod tests { + use super::*; + use feldera_types::program_schema::{ColumnType, Field}; + use std::collections::BTreeMap; + + fn relation(name: &str, columns: &[&str]) -> Relation { + Relation::new( + name.into(), + columns + .iter() + .map(|c| Field::new((*c).into(), ColumnType::varchar(true))) + .collect(), + false, + BTreeMap::new(), + ) + } + + fn writer_config(mode: PostgresWriteMode, extra_columns: &[&str]) -> PostgresWriterConfig { + serde_json::from_value(serde_json::json!({ + "uri": "postgres://localhost", + "table": "t", + "mode": mode.to_string(), + "extra_columns": extra_columns, + })) + .unwrap() + } + + // INSERT lists supplied columns, never SELECT * (#6694). + #[test] + fn materialized_insert_lists_supplied_columns() { + let key = relation("k", &["id"]); + let value = relation("v", &["id", "name"]); + let queries = RawQueries::new( + &key, + &value, + &writer_config(PostgresWriteMode::Materialized, &["audit"]), + ); + + assert!(!queries.insert.contains("SELECT *"), "{}", queries.insert); + assert!(queries.insert.contains(r#"(id, name, "audit")"#)); + assert!(queries.insert.contains(r#"SELECT id, name, "audit" FROM"#)); + } + + // Case-sensitive columns stay quoted. + #[test] + fn materialized_insert_quotes_case_sensitive_columns() { + let key = relation("k", &[r#""Id""#]); + let value = relation("v", &[r#""Id""#, r#""Name""#]); + let queries = RawQueries::new( + &key, + &value, + &writer_config(PostgresWriteMode::Materialized, &[]), + ); + + assert!(queries.insert.contains(r#"("Id", "Name")"#)); + assert!(queries.insert.contains(r#"SELECT "Id", "Name" FROM"#)); + } + + // CDC INSERT also lists the op/ts metadata columns. + #[test] + fn cdc_insert_lists_supplied_and_metadata_columns() { + let key = relation("k", &["id"]); + let value = relation("v", &["id", "name"]); + let queries = RawQueries::new(&key, &value, &writer_config(PostgresWriteMode::Cdc, &[])); + + assert!(!queries.insert.contains("SELECT *")); + assert!( + queries + .insert + .contains(r#"(id, name, "__feldera_op", "__feldera_ts")"#) + ); + assert!( + queries + .insert + .contains(r#"SELECT id, name, "__feldera_op", "__feldera_ts" FROM"#) + ); + assert_eq!(queries.insert, queries.upsert); + assert_eq!(queries.insert, queries.delete); + } +} diff --git a/crates/adapters/src/integrated/postgres/test.rs b/crates/adapters/src/integrated/postgres/test.rs index 533cf65759c..fbd1eabab5d 100644 --- a/crates/adapters/src/integrated/postgres/test.rs +++ b/crates/adapters/src/integrated/postgres/test.rs @@ -275,9 +275,27 @@ mod pg { impl TempPgTable { fn new(name: &str, uri: String, pk: bool, tls: &Option) -> Self { + Self::new_with_extra_columns(name, uri, pk, tls, "") + } + + /// Same as `new`, but appends `extra_columns` to the table: column + /// definitions absent from the Feldera schema, e.g. + /// `"served_at TIMESTAMP DEFAULT now() NOT NULL"`. + fn new_with_extra_columns( + name: &str, + uri: String, + pk: bool, + tls: &Option, + extra_columns: &str, + ) -> Self { let mut client = pg_connect(&uri, tls); let pk = if pk { "PRIMARY KEY" } else { "" }; + let extra_columns = if extra_columns.is_empty() { + String::new() + } else { + format!(",\n {extra_columns}") + }; client .execute( @@ -322,6 +340,7 @@ CREATE TABLE {name} ( map_ JSONB, __feldera_op CHAR, __feldera_ts BIGINT + {extra_columns} )"# ), &[], @@ -490,6 +509,18 @@ CREATE TABLE {name} ( TempPgTable::new(name, uri, pk, tls) } + /// Create the table with extra column definitions absent from the + /// Feldera schema. + pub fn create_table_with_extra_columns( + name: &str, + uri: String, + pk: bool, + tls: &Option, + extra_columns: &str, + ) -> TempPgTable { + TempPgTable::new_with_extra_columns(name, uri, pk, tls, extra_columns) + } + pub fn test_circuit( config: PipelineConfig, ) -> (Controller, crossbeam::channel::Receiver) { @@ -950,6 +981,94 @@ fn pg_insert0(mode: PostgresWriteMode) { .expect("timeout: failed to insert data into postgres"); } +// A NOT NULL DEFAULT column (served_at) omitted from the Feldera schema must be +// populated from its default, not NULL (#6694). +#[test] +#[serial] +fn test_pg_insert_omitted_default_column() { + let table_name = unique_pg_name("test_pg_default"); + let url = postgres_url(); + let verify_url = url.clone(); + let timeout_ms = 60_000; + + let mut data: Vec = (0..100).map(|_| rand::random()).collect(); + + let mut temp_file = NamedTempFile::new().unwrap(); + for datum in data.iter() { + let mut serializer = serde_json::Serializer::new(Vec::new()); + datum + .serialize_with_context(&mut serializer, &SqlSerdeConfig::default()) + .unwrap(); + temp_file + .as_file_mut() + .write_all(&serializer.into_inner()) + .unwrap(); + temp_file.write_all(b"\n").unwrap(); + } + + let config = serde_json::from_value(json!({ + "name": "test", + "workers": 4, + "inputs": { + "ins": { + "stream": "test_input1", + "transport": { "name": "file_input", "config": { "path": temp_file.path() } }, + "format": { "name": "json", "config": { "update_format": "raw", "array": false } } + } + }, + "outputs": { + "test_output1": { + "stream": "test_output1", + "transport": { "name": "postgres_output", "config": { "uri": url, "table": &table_name } }, + "index": "idx" + } + } + })) + .unwrap(); + + // The Feldera view (PostgresTestStruct) has no served_at column. + let mut table = PostgresTestStruct::create_table_with_extra_columns( + &table_name, + url, + true, + &None, + "served_at TIMESTAMP DEFAULT now() NOT NULL", + ); + + let (controller, err_receiver) = PostgresTestStruct::test_circuit(config); + controller.start(); + data.sort(); + + wait( + || { + let mut got = table + .query() + .into_iter() + .map(PostgresTestStruct::from) + .collect::>(); + got.sort(); + data == got || !err_receiver.is_empty() + }, + timeout_ms, + ) + .expect("timeout: failed to insert data into postgres"); + + assert!(err_receiver.is_empty(), "connector reported an error"); + + // Every row's served_at was filled by the DEFAULT rather than left NULL. + let mut client = pg::pg_connect(&verify_url, &None); + let rows = client + .query(&format!("SELECT served_at FROM {table_name}"), &[]) + .expect("failed to query table"); + assert_eq!(rows.len(), data.len()); + assert!( + rows.iter().all(|r| r + .get::<_, Option>("served_at") + .is_some()), + "served_at should be populated by its DEFAULT" + ); +} + #[test] #[serial] fn test_pg_insert() {