mod common;
use common::pgwire_harness::TestServer;
use tokio_postgres::SimpleQueryMessage;
fn command_count(msgs: &[SimpleQueryMessage]) -> Option<u64> {
msgs.iter().find_map(|m| match m {
SimpleQueryMessage::CommandComplete(n) => Some(*n),
_ => None,
})
}
fn sorted_col(msgs: &[SimpleQueryMessage], col: &str) -> Vec<String> {
let mut v: Vec<String> = msgs
.iter()
.filter_map(|m| match m {
SimpleQueryMessage::Row(r) => r.get(col).map(str::to_string),
_ => None,
})
.collect();
v.sort();
v
}
async fn setup(server: &TestServer) {
server
.exec(
"CREATE COLLECTION sensors \
COLUMNS (id TEXT, ts BIGINT TIME_KEY, v INT) \
WITH (engine='timeseries')",
)
.await
.unwrap();
for (i, ts, v) in [(1u32, 1000u64, 10u32), (2, 2000, 20), (3, 3000, 30)] {
server
.exec(&format!(
"INSERT INTO sensors (id, ts, v) VALUES ('b{i}', {ts}, {v})"
))
.await
.unwrap();
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn timeseries_in_tx_insert_visible_in_tx_raw_scan() {
let server = TestServer::start().await;
setup(&server).await;
server.exec("BEGIN").await.unwrap();
let msgs = server
.client
.simple_query("INSERT INTO sensors (id, ts, v) VALUES ('s4', 4000, 40), ('s5', 5000, 50)")
.await
.expect("in-tx timeseries insert should succeed at the statement");
assert_eq!(
command_count(&msgs),
Some(2),
"in-tx timeseries INSERT must report the real row count at statement time"
);
let all = server
.client
.simple_query("SELECT v FROM sensors")
.await
.unwrap();
assert_eq!(
sorted_col(&all, "v"),
vec!["10", "20", "30", "40", "50"],
"staged timeseries inserts must be visible in the same transaction \
alongside the committed base rows"
);
server.client.simple_query("COMMIT").await.unwrap();
let committed = server
.client
.simple_query("SELECT v FROM sensors")
.await
.unwrap();
assert_eq!(
sorted_col(&committed, "v"),
vec!["10", "20", "30", "40", "50"],
"committed timeseries inserts must persist"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn timeseries_in_tx_insert_narrow_range_sees_staged_row() {
let server = TestServer::start().await;
setup(&server).await;
server.exec("BEGIN").await.unwrap();
let msgs = server
.client
.simple_query("INSERT INTO sensors (id, ts, v) VALUES ('s4', 4000, 40)")
.await
.unwrap();
assert_eq!(command_count(&msgs), Some(1));
let windowed = server
.client
.simple_query("SELECT v FROM sensors WHERE ts >= 3500 AND ts <= 4500")
.await
.unwrap();
assert_eq!(
sorted_col(&windowed, "v"),
vec!["40"],
"a narrow time-range scan must surface the staged row inside the window"
);
let empty = server
.client
.simple_query("SELECT v FROM sensors WHERE ts >= 6000 AND ts <= 7000")
.await
.unwrap();
assert!(
sorted_col(&empty, "v").is_empty(),
"a window matching no staged or base row must return nothing"
);
server.client.simple_query("ROLLBACK").await.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn timeseries_in_tx_insert_rollback_discards_staged_rows() {
let server = TestServer::start().await;
setup(&server).await;
server.exec("BEGIN").await.unwrap();
let msgs = server
.client
.simple_query("INSERT INTO sensors (id, ts, v) VALUES ('s9', 9000, 90)")
.await
.unwrap();
assert_eq!(command_count(&msgs), Some(1));
let in_tx = server
.client
.simple_query("SELECT v FROM sensors WHERE ts >= 8500 AND ts <= 9500")
.await
.unwrap();
assert_eq!(sorted_col(&in_tx, "v"), vec!["90"]);
server.client.simple_query("ROLLBACK").await.unwrap();
let after = server
.client
.simple_query("SELECT v FROM sensors")
.await
.unwrap();
assert_eq!(
sorted_col(&after, "v"),
vec!["10", "20", "30"],
"rolled-back timeseries inserts must not persist; base rows untouched"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn timeseries_in_tx_aggregate_is_committed_only() {
let server = TestServer::start().await;
setup(&server).await;
server.exec("BEGIN").await.unwrap();
let msgs = server
.client
.simple_query("INSERT INTO sensors (id, ts, v) VALUES ('s4', 4000, 40), ('s5', 5000, 50)")
.await
.unwrap();
assert_eq!(command_count(&msgs), Some(2));
let raw = server
.client
.simple_query("SELECT v FROM sensors")
.await
.unwrap();
assert_eq!(sorted_col(&raw, "v").len(), 5);
let count = server
.query_rows("SELECT COUNT(*) FROM sensors")
.await
.unwrap();
assert_eq!(
count[0][0], "3",
"mid-transaction COUNT(*) must exclude staged rows (committed-only)"
);
let sum = server
.query_rows("SELECT SUM(v) FROM sensors")
.await
.unwrap();
assert_eq!(
sum[0][0].parse::<f64>().ok(),
Some(60.0),
"mid-transaction SUM must exclude staged rows (committed-only)"
);
server.client.simple_query("COMMIT").await.unwrap();
let count_after = server
.query_rows("SELECT COUNT(*) FROM sensors")
.await
.unwrap();
assert_eq!(
count_after[0][0], "5",
"committed timeseries inserts must be counted after COMMIT"
);
}