mod common;
use common::pgwire_harness::TestServer;
async fn array_cell_rows(server: &TestServer, array: &str, row: i64, col: i64) -> Vec<Vec<String>> {
server
.query_rows(&format!(
"SELECT * FROM ARRAY_SLICE('{array}', '{{\"row\":[{row},{row}],\"col\":[{col},{col}]}}', '*', 10)"
))
.await
.unwrap()
}
async fn setup_array(server: &TestServer, array: &str) {
server
.exec(&format!(
"CREATE ARRAY {array} \
DIMS (row INT64, col INT64) \
ATTRS (value FLOAT64) \
TILE_EXTENTS (10, 10)"
))
.await
.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn array_insert_rollback_discards_cell() {
let server = TestServer::start().await;
setup_array(&server, "arr_atomic_rb").await;
server.exec("BEGIN").await.unwrap();
server
.exec("INSERT INTO ARRAY arr_atomic_rb COORDS (1, 1) VALUES (7.0)")
.await
.unwrap();
server.client.simple_query("ROLLBACK").await.unwrap();
let rows = array_cell_rows(&server, "arr_atomic_rb", 1, 1).await;
assert!(
rows.is_empty(),
"ROLLBACK must discard a buffered ArrayOp::Put; found cell: {rows:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn array_insert_commit_persists_cell() {
let server = TestServer::start().await;
setup_array(&server, "arr_atomic_commit").await;
server.exec("BEGIN").await.unwrap();
server
.exec("INSERT INTO ARRAY arr_atomic_commit COORDS (2, 2) VALUES (9.0)")
.await
.unwrap();
server.client.simple_query("COMMIT").await.unwrap();
let rows = array_cell_rows(&server, "arr_atomic_commit", 2, 2).await;
assert_eq!(
rows.len(),
1,
"COMMIT must durably persist a buffered ArrayOp::Put; got {rows:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn array_delete_rollback_restores_cell() {
let server = TestServer::start().await;
setup_array(&server, "arr_atomic_del_rb").await;
server
.exec("INSERT INTO ARRAY arr_atomic_del_rb COORDS (3, 3) VALUES (5.0)")
.await
.unwrap();
server.exec("BEGIN").await.unwrap();
server
.exec("DELETE FROM ARRAY arr_atomic_del_rb WHERE COORDS IN ((3, 3))")
.await
.unwrap();
server.client.simple_query("ROLLBACK").await.unwrap();
let rows = array_cell_rows(&server, "arr_atomic_del_rb", 3, 3).await;
assert_eq!(
rows.len(),
1,
"ROLLBACK must discard a buffered ArrayOp::Delete; cell must survive, got {rows:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn array_insert_in_txn_not_visible_to_same_txn_read() {
let server = TestServer::start().await;
setup_array(&server, "arr_atomic_ryow").await;
server.exec("BEGIN").await.unwrap();
server
.exec("INSERT INTO ARRAY arr_atomic_ryow COORDS (4, 4) VALUES (1.0)")
.await
.unwrap();
let rows = array_cell_rows(&server, "arr_atomic_ryow", 4, 4).await;
assert!(
rows.is_empty(),
"a buffered ArrayOp::Put must NOT be visible to a read in the same \
transaction before COMMIT (RYOW loss is the documented trade-off \
for closing the atomicity gap); got {rows:?}"
);
server.client.simple_query("ROLLBACK").await.unwrap();
}
async fn create_merge_schema(server: &TestServer) {
server
.exec(
"CREATE COLLECTION merge_atomic_target (\
id TEXT PRIMARY KEY, \
name TEXT, \
score INT) WITH (engine='document_strict')",
)
.await
.unwrap();
server
.exec(
"CREATE COLLECTION merge_atomic_source (\
id TEXT PRIMARY KEY, \
name TEXT, \
score INT) WITH (engine='document_strict')",
)
.await
.unwrap();
server
.exec("INSERT INTO merge_atomic_target (id, name, score) VALUES ('a', 'alpha', 10)")
.await
.unwrap();
server
.exec("INSERT INTO merge_atomic_source (id, name, score) VALUES ('a', 'ALPHA_UPD', 99)")
.await
.unwrap();
}
const MERGE_SQL: &str = "MERGE INTO merge_atomic_target t \
USING merge_atomic_source s ON t.id = s.id \
WHEN MATCHED THEN UPDATE SET name = s.name, score = s.score";
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn document_merge_rollback_discards_changes() {
let server = TestServer::start().await;
create_merge_schema(&server).await;
server.exec("BEGIN").await.unwrap();
server.exec(MERGE_SQL).await.unwrap();
server.client.simple_query("ROLLBACK").await.unwrap();
let rows = server
.query_rows("SELECT name, score FROM merge_atomic_target WHERE id = 'a'")
.await
.unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(
rows[0],
vec!["alpha".to_string(), "10".to_string()],
"ROLLBACK must discard a buffered DocumentOp::Merge; got {rows:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn document_merge_commit_persists_changes() {
let server = TestServer::start().await;
create_merge_schema(&server).await;
server.exec("BEGIN").await.unwrap();
server.exec(MERGE_SQL).await.unwrap();
server.client.simple_query("COMMIT").await.unwrap();
let rows = server
.query_rows("SELECT name, score FROM merge_atomic_target WHERE id = 'a'")
.await
.unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(
rows[0],
vec!["ALPHA_UPD".to_string(), "99".to_string()],
"COMMIT must durably persist a buffered DocumentOp::Merge; got {rows:?}"
);
}