use std::time::Duration;
use serde_json::{Map, Value as JsonValue};
use crate::bridge::envelope::PhysicalPlan;
use crate::control::security::catalog::StoredContinuousAggregate;
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::server::response_shape::types::{DdlColType, ShapedRows};
use crate::control::server::shared::ddl::sync_dispatch;
use crate::control::state::SharedState;
use crate::engine::timeseries::continuous_agg::{AggregateInfo, ContinuousAggregateDef};
use crate::types::DatabaseId;
use nodedb_physical::physical_plan::MetaOp;
use super::super::super::result::{DdlError, DdlResult};
pub async fn show_continuous_aggregates(
state: &SharedState,
identity: &AuthenticatedIdentity,
database_id: DatabaseId,
parts: &[&str],
) -> Result<Vec<DdlResult>, DdlError> {
let source_filter = if parts.len() >= 5 && parts[3].to_uppercase() == "FOR" {
Some(parts[4].to_lowercase())
} else {
None
};
let tenant_id = identity.tenant_id;
let stored_aggs: Vec<StoredContinuousAggregate> = state
.credentials
.catalog()
.list_continuous_aggregates(database_id.as_u64(), tenant_id.as_u64())
.ok()
.unwrap_or_default();
let runtime_infos: Vec<AggregateInfo> = match sync_dispatch::dispatch_async(
state,
tenant_id,
database_id,
"__system",
PhysicalPlan::Meta(MetaOp::ListContinuousAggregates),
Duration::from_secs(5),
)
.await
{
Ok(payload) => sonic_rs::from_slice(&payload).unwrap_or_default(),
Err(_) => Vec::new(),
};
let columns = vec![
"name".to_string(),
"source".to_string(),
"bucket_interval".to_string(),
"refresh_policy".to_string(),
"watermark_ts".to_string(),
"rows_aggregated".to_string(),
"materialized_buckets".to_string(),
"stale".to_string(),
];
let column_types = vec![
DdlColType::Text,
DdlColType::Text,
DdlColType::Text,
DdlColType::Text,
DdlColType::Int8,
DdlColType::Int8,
DdlColType::Int8,
DdlColType::Text,
];
let mut rows = Vec::new();
for stored in &stored_aggs {
if let Some(ref filter) = source_filter
&& stored.source != *filter
{
continue;
}
let Ok(def) = zerompk::from_msgpack::<ContinuousAggregateDef>(&stored.def_bytes) else {
tracing::warn!(
cagg = %stored.name,
tenant = stored.tenant_id,
"continuous aggregate row has unreadable def_bytes; \
skipping in SHOW (the row is still durable in the catalog)"
);
continue;
};
let runtime = runtime_infos.iter().find(|i| i.name == stored.name);
let watermark = runtime.map(|i| i.watermark_ts).unwrap_or(0);
let rows_agg = runtime.map(|i| i.rows_aggregated as i64).unwrap_or(0);
let buckets = runtime.map(|i| i.materialized_buckets as i64).unwrap_or(0);
let stale = runtime.map(|i| i.stale).unwrap_or(def.stale).to_string();
let mut row = Map::new();
row.insert("name".to_string(), JsonValue::String(stored.name.clone()));
row.insert(
"source".to_string(),
JsonValue::String(stored.source.clone()),
);
row.insert(
"bucket_interval".to_string(),
JsonValue::String(def.bucket_interval.clone()),
);
row.insert(
"refresh_policy".to_string(),
JsonValue::String(format!("{:?}", def.refresh_policy)),
);
row.insert(
"watermark_ts".to_string(),
JsonValue::String(watermark.to_string()),
);
row.insert(
"rows_aggregated".to_string(),
JsonValue::String(rows_agg.to_string()),
);
row.insert(
"materialized_buckets".to_string(),
JsonValue::String(buckets.to_string()),
);
row.insert("stale".to_string(), JsonValue::String(stale));
rows.push(row);
}
Ok(vec![DdlResult::Rows(ShapedRows {
columns,
column_types,
rows,
notice: None,
})])
}