mod common;
use common::pgwire_harness::TestServer;
async fn scan_ints(server: &TestServer, sql: &str) -> Vec<i64> {
let mut v: Vec<i64> = server
.query_text(sql)
.await
.unwrap()
.into_iter()
.map(|s| s.parse().unwrap())
.collect();
v.sort_unstable();
v
}
async fn setup(server: &TestServer, coll: &str, engine: &str) {
server
.exec(&format!(
"CREATE COLLECTION {coll} \
(id STRING NOT NULL PRIMARY KEY, n INT) WITH (engine='{engine}')"
))
.await
.unwrap();
for (id, n) in [("a", 1), ("b", 2), ("unrelated", 100)] {
server
.exec(&format!("INSERT INTO {coll} (id, n) VALUES ('{id}', {n})"))
.await
.unwrap();
}
}
async fn scan_sees_own_insert(engine: &str, coll: &str) {
let server = TestServer::start().await;
setup(&server, coll, engine).await;
server.exec("BEGIN").await.unwrap();
server
.exec(&format!("INSERT INTO {coll} (id, n) VALUES ('c', 3)"))
.await
.unwrap();
let seen = scan_ints(&server, &format!("SELECT n FROM {coll}")).await;
assert_eq!(
seen,
vec![1, 2, 3, 100],
"{engine}: in-tx scan must include the staged insert"
);
server.client.simple_query("COMMIT").await.unwrap();
let after = scan_ints(&server, &format!("SELECT n FROM {coll}")).await;
assert_eq!(
after,
vec![1, 2, 3, 100],
"{engine}: staged insert persists"
);
}
async fn scan_excludes_own_delete(engine: &str, coll: &str) {
let server = TestServer::start().await;
setup(&server, coll, engine).await;
server.exec("BEGIN").await.unwrap();
server
.exec(&format!("DELETE FROM {coll} WHERE id = 'b'"))
.await
.unwrap();
let seen = scan_ints(&server, &format!("SELECT n FROM {coll}")).await;
assert_eq!(
seen,
vec![1, 100],
"{engine}: in-tx scan must hide the staged delete"
);
server.client.simple_query("ROLLBACK").await.unwrap();
let after = scan_ints(&server, &format!("SELECT n FROM {coll}")).await;
assert_eq!(
after,
vec![1, 2, 100],
"{engine}: ROLLBACK restores the base row"
);
}
async fn scan_reflects_predicate_move(engine: &str, coll: &str) {
let server = TestServer::start().await;
setup(&server, coll, engine).await;
let base = scan_ints(&server, &format!("SELECT n FROM {coll} WHERE n = 1")).await;
assert_eq!(base, vec![1], "{engine}: base predicate match");
server.exec("BEGIN").await.unwrap();
server
.exec(&format!("UPDATE {coll} SET n = 7 WHERE id = 'a'"))
.await
.unwrap();
server
.exec(&format!("UPDATE {coll} SET n = 1 WHERE id = 'b'"))
.await
.unwrap();
let matched = scan_ints(&server, &format!("SELECT n FROM {coll} WHERE n = 1")).await;
assert_eq!(
matched,
vec![1],
"{engine}: predicate must re-eval staged bodies — 'b' now matches, 'a' does not"
);
let moved_out = scan_ints(&server, &format!("SELECT n FROM {coll} WHERE n = 7")).await;
assert_eq!(
moved_out,
vec![7],
"{engine}: staged-updated row visible under its new predicate value"
);
server.client.simple_query("ROLLBACK").await.unwrap();
}
async fn scan_order_by_limit(engine: &str, coll: &str) {
let server = TestServer::start().await;
setup(&server, coll, engine).await;
server.exec("BEGIN").await.unwrap();
server
.exec(&format!("INSERT INTO {coll} (id, n) VALUES ('c', 3)"))
.await
.unwrap();
server
.exec(&format!("DELETE FROM {coll} WHERE id = 'a'"))
.await
.unwrap();
let ordered: Vec<i64> = server
.query_text(&format!("SELECT n FROM {coll} ORDER BY n ASC LIMIT 2"))
.await
.unwrap()
.into_iter()
.map(|s| s.parse().unwrap())
.collect();
assert_eq!(
ordered,
vec![2, 3],
"{engine}: merged result must sort + limit correctly"
);
server.client.simple_query("ROLLBACK").await.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn schemaless_scan_sees_own_insert() {
scan_sees_own_insert("document_schemaless", "sc_ins").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn schemaless_scan_excludes_own_delete() {
scan_excludes_own_delete("document_schemaless", "sc_del").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn schemaless_scan_reflects_predicate_move() {
scan_reflects_predicate_move("document_schemaless", "sc_mov").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn schemaless_scan_order_by_limit() {
scan_order_by_limit("document_schemaless", "sc_lim").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn strict_scan_sees_own_insert() {
scan_sees_own_insert("document_strict", "st_ins").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn strict_scan_excludes_own_delete() {
scan_excludes_own_delete("document_strict", "st_del").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn strict_scan_reflects_predicate_move() {
scan_reflects_predicate_move("document_strict", "st_mov").await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn strict_scan_order_by_limit() {
scan_order_by_limit("document_strict", "st_lim").await;
}
async fn setup_bitemporal(server: &TestServer, coll: &str) {
server
.exec(&format!(
"CREATE COLLECTION {coll} (id STRING PRIMARY KEY, n INT) \
WITH (engine='document_schemaless', bitemporal=true)"
))
.await
.unwrap();
for (id, n) in [("a", 1), ("b", 2), ("unrelated", 100)] {
server
.exec(&format!("INSERT INTO {coll} (id, n) VALUES ('{id}', {n})"))
.await
.unwrap();
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn bitemporal_scan_sees_own_writes() {
let server = TestServer::start().await;
let coll = "bt_scan";
setup_bitemporal(&server, coll).await;
server.exec("BEGIN").await.unwrap();
server
.exec(&format!("INSERT INTO {coll} (id, n) VALUES ('c', 3)"))
.await
.unwrap();
server
.exec(&format!("DELETE FROM {coll} WHERE id = 'b'"))
.await
.unwrap();
server
.exec(&format!("UPDATE {coll} SET n = 9 WHERE id = 'a'"))
.await
.unwrap();
let seen = scan_ints(&server, &format!("SELECT n FROM {coll}")).await;
assert_eq!(
seen,
vec![3, 9, 100],
"bitemporal current-version scan must merge staged insert+delete+update"
);
server.client.simple_query("ROLLBACK").await.unwrap();
let after = scan_ints(&server, &format!("SELECT n FROM {coll}")).await;
assert_eq!(
after,
vec![1, 2, 100],
"ROLLBACK: bitemporal scan sees base only"
);
}