use std::sync::Arc;
use datafusion::physical_plan::ExecutionPlan;
use crate::arrow::array::RecordBatch;
use crate::arrow::datatypes::SchemaRef;
use crate::error::{Error, Result};
use crate::planner::ReadMode;
use crate::watermark::{Tier, TieringWatermark};
#[derive(Debug, Clone)]
pub struct QueryResult {
batches: Vec<RecordBatch>,
schema: SchemaRef,
watermark: TieringWatermark,
watermarks: Vec<(String, TieringWatermark)>,
tiers: Vec<Tier>,
read_mode: ReadMode,
}
impl QueryResult {
pub(crate) fn new(
batches: Vec<RecordBatch>,
schema: SchemaRef,
watermarks: Vec<(String, TieringWatermark)>,
tiers: Vec<Tier>,
read_mode: ReadMode,
) -> Self {
let watermark = watermarks
.iter()
.map(|(_, w)| *w)
.min()
.unwrap_or_else(TieringWatermark::empty);
Self {
batches,
schema,
watermark,
watermarks,
tiers,
read_mode,
}
}
pub fn batches(&self) -> &[RecordBatch] {
&self.batches
}
pub fn into_batches(self) -> Vec<RecordBatch> {
self.batches
}
pub fn schema(&self) -> SchemaRef {
self.schema.clone()
}
pub fn num_rows(&self) -> usize {
self.batches.iter().map(RecordBatch::num_rows).sum()
}
pub fn to_json(&self) -> Result<Vec<serde_json::Value>> {
if self.num_rows() == 0 {
return Ok(Vec::new());
}
let mut buf = Vec::new();
let mut writer = crate::arrow::json::ArrayWriter::new(&mut buf);
for batch in &self.batches {
writer
.write(batch)
.map_err(|e| Error::Storage(e.to_string()))?;
}
writer.finish().map_err(|e| Error::Storage(e.to_string()))?;
match serde_json::from_slice(&buf).map_err(|e| Error::Storage(e.to_string()))? {
serde_json::Value::Array(rows) => Ok(rows),
other => Ok(vec![other]),
}
}
pub fn watermark(&self) -> TieringWatermark {
self.watermark
}
pub fn watermarks(&self) -> &[(String, TieringWatermark)] {
&self.watermarks
}
pub fn tiers_scanned(&self) -> &[Tier] {
&self.tiers
}
pub fn spans_tiers(&self) -> bool {
self.tiers.contains(&Tier::Cold) && self.tiers.contains(&Tier::Hot)
}
pub fn touched_hot_tier(&self) -> bool {
self.tiers.contains(&Tier::Hot)
}
pub fn read_mode(&self) -> ReadMode {
self.read_mode
}
}
#[derive(Debug, Clone)]
pub struct QueryDescription {
schema: SchemaRef,
watermark: TieringWatermark,
watermarks: Vec<(String, TieringWatermark)>,
tiers: Vec<Tier>,
read_mode: ReadMode,
}
impl QueryDescription {
pub(crate) fn new(
schema: SchemaRef,
watermarks: Vec<(String, TieringWatermark)>,
tiers: Vec<Tier>,
read_mode: ReadMode,
) -> Self {
let watermark = watermarks
.iter()
.map(|(_, w)| *w)
.min()
.unwrap_or_else(TieringWatermark::empty);
Self {
schema,
watermark,
watermarks,
tiers,
read_mode,
}
}
pub fn schema(&self) -> SchemaRef {
self.schema.clone()
}
pub fn watermark(&self) -> TieringWatermark {
self.watermark
}
pub fn watermarks(&self) -> &[(String, TieringWatermark)] {
&self.watermarks
}
pub fn tiers_scanned(&self) -> &[Tier] {
&self.tiers
}
pub fn spans_tiers(&self) -> bool {
self.tiers.contains(&Tier::Cold) && self.tiers.contains(&Tier::Hot)
}
pub fn touched_hot_tier(&self) -> bool {
self.tiers.contains(&Tier::Hot)
}
pub fn read_mode(&self) -> ReadMode {
self.read_mode
}
}
pub(crate) fn tiers_of(plan: &Arc<dyn ExecutionPlan>) -> Vec<Tier> {
let mut found = Vec::new();
walk(plan, &mut found);
let mut tiers = Vec::with_capacity(2);
if found.contains(&Tier::Cold) {
tiers.push(Tier::Cold);
}
if found.contains(&Tier::Hot) {
tiers.push(Tier::Hot);
}
tiers
}
fn walk(plan: &Arc<dyn ExecutionPlan>, found: &mut Vec<Tier>) {
let any = plan.as_any();
if any
.downcast_ref::<crate::planner::provider::HotScanExec>()
.is_some()
{
found.push(Tier::Hot);
} else if any
.downcast_ref::<iceberg_datafusion::physical_plan::IcebergTableScan>()
.is_some()
{
found.push(Tier::Cold);
}
for child in plan.children() {
walk(child, found);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::arrow::datatypes::{DataType, Field, Schema};
fn empty(tiers: Vec<Tier>) -> QueryResult {
let schema = Arc::new(Schema::new(vec![Field::new("n", DataType::Int64, false)]));
QueryResult::new(
Vec::new(),
schema,
vec![(
"readings_versions".to_string(),
TieringWatermark::new(time::macros::datetime!(2026-07-20 00:00 UTC)),
)],
tiers,
ReadMode::Unified,
)
}
#[test]
fn spanning_requires_both_tiers() {
assert!(empty(vec![Tier::Cold, Tier::Hot]).spans_tiers());
assert!(!empty(vec![Tier::Cold]).spans_tiers());
assert!(!empty(vec![Tier::Hot]).spans_tiers());
assert!(!empty(vec![]).spans_tiers());
}
#[test]
fn a_reporting_query_can_tell_it_read_the_mutable_tier() {
assert!(empty(vec![Tier::Cold, Tier::Hot]).touched_hot_tier());
assert!(!empty(vec![Tier::Cold]).touched_hot_tier());
}
#[test]
fn the_watermark_travels_with_the_rows() {
let r = empty(vec![Tier::Cold]);
assert_eq!(
r.watermark().get(),
time::macros::datetime!(2026-07-20 00:00 UTC)
);
}
#[test]
fn an_empty_result_still_reports_row_count_zero() {
assert_eq!(empty(vec![]).num_rows(), 0);
}
#[test]
fn an_empty_result_serialises_to_an_empty_array() {
assert_eq!(
empty(vec![]).to_json().unwrap(),
Vec::<serde_json::Value>::new()
);
}
#[test]
fn to_json_maps_each_row_to_an_object() {
use crate::arrow::array::{Int64Array, StringArray};
let schema = Arc::new(Schema::new(vec![
Field::new("malo_id", DataType::Utf8, false),
Field::new("total_kwh", DataType::Int64, false),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(StringArray::from(vec!["11111111111", "22222222222"])),
Arc::new(Int64Array::from(vec![42, 7])),
],
)
.unwrap();
let result = QueryResult::new(
vec![batch],
schema,
vec![(
"readings_versions".to_string(),
TieringWatermark::new(time::macros::datetime!(2026-07-20 00:00 UTC)),
)],
vec![Tier::Cold],
ReadMode::Unified,
);
let rows = result.to_json().unwrap();
assert_eq!(rows.len(), 2);
assert_eq!(rows[0]["malo_id"], "11111111111");
assert_eq!(rows[0]["total_kwh"], 42);
assert_eq!(rows[1]["malo_id"], "22222222222");
}
}