Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
102 changes: 100 additions & 2 deletions crates/adapters/src/integrated/postgres/prepared_statements.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> = 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();
Expand Down Expand Up @@ -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}"#,
);
}
}
Expand Down Expand Up @@ -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);
}
}
119 changes: 119 additions & 0 deletions crates/adapters/src/integrated/postgres/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -275,9 +275,27 @@ mod pg {

impl TempPgTable {
fn new(name: &str, uri: String, pk: bool, tls: &Option<PostgresTlsConfig>) -> 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<PostgresTlsConfig>,
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(
Expand Down Expand Up @@ -322,6 +340,7 @@ CREATE TABLE {name} (
map_ JSONB,
__feldera_op CHAR,
__feldera_ts BIGINT
{extra_columns}
)"#
),
&[],
Expand Down Expand Up @@ -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<PostgresTlsConfig>,
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<String>) {
Expand Down Expand Up @@ -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<PostgresTestStruct> = (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::<Vec<_>>();
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<chrono::NaiveDateTime>>("served_at")
.is_some()),
"served_at should be populated by its DEFAULT"
);
}

#[test]
#[serial]
fn test_pg_insert() {
Expand Down
Loading