#![cfg(all(feature = "testkit", feature = "flight"))]
use arrow_flight::sql::client::FlightSqlServiceClient;
use meterstore::serve::FlightSqlServer;
use meterstore::testkit::{MeteringWorkload, Oracle, TestHarness};
use time::macros::datetime;
use time::{Duration, OffsetDateTime};
use tonic::transport::Channel;
const START: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
struct Served {
client: FlightSqlServiceClient<Channel>,
_shutdown: tokio::sync::oneshot::Sender<()>,
}
async fn serve(store: meterstore::MeterStore) -> Served {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind an ephemeral port");
let address = listener.local_addr().expect("local address");
let (shutdown, rx) = tokio::sync::oneshot::channel::<()>();
tokio::spawn(async move {
tonic::transport::Server::builder()
.add_service(FlightSqlServer::new(store).into_service())
.serve_with_incoming_shutdown(
tokio_stream::wrappers::TcpListenerStream::new(listener),
async {
let _ = rx.await;
},
)
.await
.expect("flight server");
});
let channel = Channel::from_shared(format!("http://{address}"))
.expect("endpoint")
.connect()
.await
.expect("connect to the flight server");
Served {
client: FlightSqlServiceClient::new(channel),
_shutdown: shutdown,
}
}
async fn query(served: &mut Served, sql: &str) -> Vec<datafusion::arrow::array::RecordBatch> {
use futures::TryStreamExt;
let info = served
.client
.execute(sql.to_string(), None)
.await
.expect("get_flight_info");
let ticket = info.endpoint[0]
.ticket
.clone()
.expect("an endpoint carries a ticket");
served
.client
.do_get(ticket)
.await
.expect("do_get")
.try_collect()
.await
.expect("collect batches")
}
async fn split_store() -> (TestHarness, meterstore::MeterStore, Oracle) {
let workload = MeteringWorkload::new(START)
.seed(0xF117)
.malo_ids(4)
.days(3);
let (from, to) = workload.range();
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(from, to + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(from).await.expect("watermark");
let store = harness.store().await.expect("store");
let series = workload.generate().expect("workload");
let mut oracle = Oracle::new();
oracle.record(&series).expect("oracle");
harness.ingest(&store, &series).await.expect("ingest");
store
.archive(from + Duration::days(3), 2)
.await
.expect("archive");
assert_eq!(
store.watermark().await.unwrap().get(),
from + Duration::days(2),
"the fixture must actually straddle the watermark"
);
(harness, store, oracle)
}
fn expect_status(error: arrow_flight::error::FlightError) -> tonic::Status {
match error {
arrow_flight::error::FlightError::Tonic(status) => *status,
other => panic!("expected a gRPC status, got {other:?}"),
}
}
fn count(batches: &[datafusion::arrow::array::RecordBatch]) -> i64 {
use datafusion::arrow::array::AsArray;
batches[0]
.column(0)
.as_primitive::<datafusion::arrow::datatypes::Int64Type>()
.value(0)
}
#[tokio::test]
async fn a_flight_client_gets_the_unified_view() {
let (_h, store, oracle) = split_store().await;
let (from, to) = (START, START + Duration::days(3));
let mut served = serve(store).await;
let batches = query(&mut served, "SELECT COUNT(*) FROM readings").await;
assert_eq!(
count(&batches) as u64,
oracle.row_count(from, to),
"a Flight client must see both tiers"
);
}
#[tokio::test]
async fn results_carry_the_boundary_they_were_computed_against() {
let (_h, store, _oracle) = split_store().await;
let expected = store.watermark().await.unwrap();
let mut served = serve(store).await;
let batches = query(&mut served, "SELECT COUNT(*) FROM readings").await;
let metadata = batches[0].schema();
let metadata = metadata.metadata();
assert_eq!(
metadata
.get(meterstore::watermark::WATERMARK_PROPERTY)
.map(String::as_str),
Some(expected.to_string().as_str()),
"the watermark must travel with the rows"
);
let tiers = metadata
.get("meterstore.tiers_scanned")
.expect("tiers scanned");
assert!(
tiers.contains("cold") && tiers.contains("hot"),
"this query spans both tiers, and the metadata must say so: {tiers}"
);
}
#[tokio::test]
async fn a_corrected_interval_is_counted_once_over_flight() {
let workload = MeteringWorkload::new(START)
.seed(0xC0DE)
.malo_ids(3)
.days(2)
.with_corrections(0.25);
let (from, to) = workload.range();
let harness = TestHarness::start().await.expect("harness");
harness
.ensure_partitions(from, to + Duration::days(1))
.await
.expect("partitions");
harness.seed_watermark(from).await.expect("watermark");
let store = harness.store().await.expect("store");
let series = workload.generate().expect("workload");
let mut oracle = Oracle::new();
oracle.record(&series).expect("oracle");
harness.ingest(&store, &series).await.expect("ingest");
let mut served = serve(store).await;
let resolved = query(&mut served, "SELECT COUNT(*) FROM readings").await;
assert_eq!(count(&resolved) as u64, oracle.row_count(from, to));
let raw = query(&mut served, "SELECT COUNT(*) FROM readings_versions").await;
assert!(
count(&raw) as u64 > oracle.row_count(from, to),
"the fixture must produce corrections for this to mean anything"
);
}
#[tokio::test]
async fn a_client_can_browse_the_tables() {
let (_h, store, _oracle) = split_store().await;
let mut served = serve(store).await;
let batches = query(
&mut served,
"SELECT table_name FROM information_schema.tables ORDER BY 1",
)
.await;
use datafusion::arrow::array::AsArray;
let mut names = Vec::new();
for batch in &batches {
let column = batch.column(0).as_string::<i32>();
for i in 0..batch.num_rows() {
names.push(column.value(i).to_string());
}
}
assert!(names.iter().any(|n| n == "readings"), "{names:?}");
assert!(names.iter().any(|n| n == "readings_versions"), "{names:?}");
}
#[tokio::test]
async fn a_write_is_refused_with_the_reason() {
let (_h, store, _oracle) = split_store().await;
let mut served = serve(store).await;
let error = served
.client
.execute_update(
"INSERT INTO readings_versions (malo_id) VALUES ('x')".to_string(),
None,
)
.await
.expect_err("a write must be refused");
let status = expect_status(error);
assert_eq!(status.code(), tonic::Code::PermissionDenied);
let message = status.message();
assert!(message.contains("read-only"), "{message}");
assert!(message.contains("MeterStore::append"), "{message}");
}
#[tokio::test]
async fn a_prepared_statement_round_trips() {
let (_h, store, oracle) = split_store().await;
let (from, to) = (START, START + Duration::days(3));
let mut served = serve(store).await;
let mut prepared = served
.client
.prepare("SELECT COUNT(*) FROM readings".to_string(), None)
.await
.expect("prepare");
let info = prepared.execute().await.expect("execute");
let ticket = info.endpoint[0].ticket.clone().expect("ticket");
use futures::TryStreamExt;
let batches: Vec<_> = served
.client
.do_get(ticket)
.await
.expect("do_get")
.try_collect()
.await
.expect("collect");
assert_eq!(count(&batches) as u64, oracle.row_count(from, to));
prepared.close().await.expect("close");
}
#[tokio::test]
async fn a_bad_statement_fails_at_plan_time() {
let (_h, store, _oracle) = split_store().await;
let mut served = serve(store).await;
let error = served
.client
.execute("SELECT * FROM no_such_table".to_string(), None)
.await
.expect_err("an unknown table must fail");
assert_eq!(expect_status(error).code(), tonic::Code::InvalidArgument);
}