List all registered aggregates with status.
(&self)
| 238 | |
| 239 | /// List all registered aggregates with status. |
| 240 | pub fn list_aggregates(&self) -> Vec<AggregateInfo> { |
| 241 | self.definitions |
| 242 | .values() |
| 243 | .map(|def| { |
| 244 | let wm = self.watermarks.get(&def.name); |
| 245 | let bucket_count = self |
| 246 | .materialized |
| 247 | .get(&def.name) |
| 248 | .map_or(0, |m| m.len() as u64); |
| 249 | AggregateInfo { |
| 250 | name: def.name.clone(), |
| 251 | source: def.source.clone(), |
| 252 | bucket_interval: def.bucket_interval.clone(), |
| 253 | refresh_policy: def.refresh_policy.clone(), |
| 254 | watermark_ts: wm.map_or(i64::MIN, |w| w.watermark_ts), |
| 255 | rows_aggregated: wm.map_or(0, |w| w.rows_aggregated), |
| 256 | materialized_buckets: bucket_count, |
| 257 | stale: def.stale, |
| 258 | } |
| 259 | }) |
| 260 | .collect() |
| 261 | } |
| 262 | } |
| 263 | |
| 264 | impl Default for ContinuousAggregateManager { |
no test coverage detected