#![cfg(feature = "pg-walstream")]
extern crate alloc;
use alloc::sync::Arc;
use alloc::vec::Vec;
use sqlite_diff_rs::pg_walstream::{ColumnValue, ConversionError, EventType, PgWalstream, RowData};
use sqlite_diff_rs::{ChangeSet, ChangesetOp, DecodeError, PatchSet, TypeMap, Value};
mod common;
use common::{TestUsersTable, test_schema};
fn default_adapter() -> TypeMap<PgWalstream, String, Vec<u8>> {
TypeMap::defaults()
}
fn row_data(id: i64, name: &str, active: bool) -> RowData {
let mut data = RowData::new();
data.push(Arc::from("id"), ColumnValue::text(&id.to_string()));
data.push(Arc::from("name"), ColumnValue::text(name));
data.push(
Arc::from("active"),
ColumnValue::text(if active { "t" } else { "f" }),
);
data
}
#[test]
fn pg_changeset_insert() {
let schema = test_schema();
let adapter = default_adapter();
let data = row_data(1, "Alice", true);
let event = EventType::Insert {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
data,
};
let cs: ChangeSet<TestUsersTable, String, Vec<u8>> =
ChangeSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = cs.build();
assert!(!bytes.is_empty(), "changeset must contain data");
assert_eq!(bytes[0], b'T', "changeset marker must be 'T'");
}
#[test]
fn pg_changeset_update() {
let schema = test_schema();
let adapter = default_adapter();
let old = row_data(1, "Alice", true);
let new = row_data(1, "Alicia", true);
let event = EventType::Update {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: Some(old),
new_data: new,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let cs: ChangeSet<TestUsersTable, String, Vec<u8>> =
ChangeSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = cs.build();
assert!(!bytes.is_empty(), "changeset must contain data");
assert_eq!(bytes[0], b'T', "changeset marker must be 'T'");
}
#[test]
fn pg_changeset_delete() {
let schema = test_schema();
let adapter = default_adapter();
let data = row_data(1, "Alice", true);
let event = EventType::Delete {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: data,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let cs: ChangeSet<TestUsersTable, String, Vec<u8>> =
ChangeSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = cs.build();
assert!(!bytes.is_empty(), "changeset must contain data");
assert_eq!(bytes[0], b'T', "changeset marker must be 'T'");
}
#[test]
fn pg_patchset_insert() {
let schema = test_schema();
let adapter = default_adapter();
let data = row_data(1, "Alice", true);
let event = EventType::Insert {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
data,
};
let ps: PatchSet<TestUsersTable, String, Vec<u8>> =
PatchSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = ps.build();
assert!(!bytes.is_empty(), "patchset must contain data");
assert_eq!(bytes[0], b'P', "patchset marker must be 'P'");
}
#[test]
fn pg_patchset_update() {
let schema = test_schema();
let adapter = default_adapter();
let old = row_data(1, "Alice", true);
let new = row_data(1, "Alicia", true);
let event = EventType::Update {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: Some(old),
new_data: new,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let ps: PatchSet<TestUsersTable, String, Vec<u8>> =
PatchSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = ps.build();
assert!(!bytes.is_empty(), "patchset must contain data");
assert_eq!(bytes[0], b'P', "patchset marker must be 'P'");
}
#[test]
fn pg_patchset_delete() {
let schema = test_schema();
let adapter = default_adapter();
let data = row_data(1, "Alice", true);
let event = EventType::Delete {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: data,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let ps: PatchSet<TestUsersTable, String, Vec<u8>> =
PatchSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = ps.build();
assert!(!bytes.is_empty(), "patchset must contain data");
assert_eq!(bytes[0], b'P', "patchset marker must be 'P'");
}
#[test]
fn pg_table_not_found_is_error() {
let schema = test_schema();
let adapter = default_adapter();
let data = row_data(1, "Alice", true);
let event = EventType::Insert {
schema: Arc::from("public"),
table: Arc::from("nonexistent"),
relation_oid: 1,
data,
};
let result: Result<ChangeSet<TestUsersTable, String, Vec<u8>>, ConversionError> =
ChangeSet::new().digest(&event, &schema, &adapter);
match result {
Err(ConversionError::TableNotFound(n)) => assert_eq!(n, "nonexistent"),
Err(other) => panic!("expected TableNotFound, got {other:?}"),
Ok(_) => panic!("expected error"),
}
}
#[test]
fn pg_column_not_found_is_error() {
let schema = test_schema();
let adapter = default_adapter();
let mut data = RowData::new();
data.push(Arc::from("id"), ColumnValue::text("1"));
data.push(Arc::from("missing_col"), ColumnValue::text("val"));
let event = EventType::Insert {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
data,
};
let result: Result<ChangeSet<TestUsersTable, String, Vec<u8>>, ConversionError> =
ChangeSet::new().digest(&event, &schema, &adapter);
match result {
Err(ConversionError::ColumnNotFound(n)) => assert!(n.contains("missing_col")),
Err(other) => panic!("expected ColumnNotFound, got {other:?}"),
Ok(_) => panic!("expected error"),
}
}
#[test]
fn pg_decode_error_is_propagated() {
let adapter: TypeMap<PgWalstream, String, Vec<u8>> = TypeMap::new();
let schema = test_schema();
let data = row_data(1, "Alice", true);
let event = EventType::Insert {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
data,
};
let result: Result<ChangeSet<TestUsersTable, String, Vec<u8>>, ConversionError> =
ChangeSet::new().digest(&event, &schema, &adapter);
match result {
Err(ConversionError::Decode(DecodeError::NoDecoderForType { column })) => {
assert!(!column.is_empty());
}
Err(other) => panic!("expected Decode(NoDecoderForType), got {other:?}"),
Ok(_) => panic!("expected error"),
}
}
#[test]
fn pg_missing_old_data_still_works_for_changeset_update() {
let schema = test_schema();
let adapter = default_adapter();
let new = row_data(1, "Alice", true);
let event = EventType::Update {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: None,
new_data: new,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let cs: ChangeSet<TestUsersTable, String, Vec<u8>> =
ChangeSet::new().digest(&event, &schema, &adapter).unwrap();
let bytes: Vec<u8> = cs.build();
assert!(
!bytes.is_empty(),
"changeset should produce output with None old_data"
);
}
#[test]
fn pg_update_with_no_old_data_patchset_still_works() {
let schema = test_schema();
let new = row_data(1, "Alice", true);
let event = EventType::Update {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: None,
new_data: new,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let ps: PatchSet<TestUsersTable, String, Vec<u8>> = PatchSet::new()
.digest(&event, &schema, &default_adapter())
.unwrap();
let bytes: Vec<u8> = ps.build();
assert!(
!bytes.is_empty(),
"patchset update with no old data must produce output"
);
}
#[test]
fn pg_changeset_update_captures_old_pk_when_old_data_absent() {
let schema = test_schema();
let adapter = default_adapter();
let new = row_data(1, "Alicia", true);
let event = EventType::Update {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: None,
new_data: new,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let cs: ChangeSet<TestUsersTable, String, Vec<u8>> =
ChangeSet::new().digest(&event, &schema, &adapter).unwrap();
let ops: Vec<_> = cs.iter().collect();
assert_eq!(ops.len(), 1);
match &ops[0] {
ChangesetOp::Update { values, .. } => {
assert_eq!(
values[0].0,
Some(Value::Integer(1)),
"old primary key must be captured from the new tuple when old_data is absent"
);
assert_eq!(
values[1].0, None,
"non-key old absent under default identity"
);
assert_eq!(
values[1].1,
Some(Value::Text("Alicia".to_string())),
"new name"
);
}
other => panic!("expected update, got {other:?}"),
}
}
#[test]
fn pg_changeset_update_captures_changed_pk() {
let schema = test_schema();
let adapter = default_adapter();
let old = row_data(1, "Alice", true);
let new = row_data(2, "Alice", true);
let event = EventType::Update {
schema: Arc::from("public"),
table: Arc::from("users"),
relation_oid: 1,
old_data: Some(old),
new_data: new,
replica_identity: pg_walstream::ReplicaIdentity::Default,
key_columns: alloc::vec![Arc::from("id")],
};
let cs: ChangeSet<TestUsersTable, String, Vec<u8>> =
ChangeSet::new().digest(&event, &schema, &adapter).unwrap();
let ops: Vec<_> = cs.iter().collect();
match &ops[0] {
ChangesetOp::Update { values, .. } => {
assert_eq!(values[0].0, Some(Value::Integer(1)), "old key");
assert_eq!(values[0].1, Some(Value::Integer(2)), "new key");
}
other => panic!("expected update, got {other:?}"),
}
}