From 13b6d850a32b495b83293dbc5990b366fadbb50f Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Mon, 28 Sep 2026 13:50:23 +0530 Subject: [PATCH 1/6] fix: query execution at the time of hottier eviction query takes a guard if time range matches hottier eviction flow skips the hottier sync cycle so query can be served from hottier add check to skip stream related tasks for suspended tenants like - - hottier sync - retention task --- src/alerts/alerts_utils.rs | 9 ++ src/hottier.rs | 249 +++++++++++++++++++++++++++++++++++-- src/hottier/local_state.rs | 22 +++- src/query/mod.rs | 48 ++++++- src/storage/retention.rs | 17 ++- 5 files changed, 328 insertions(+), 17 deletions(-) diff --git a/src/alerts/alerts_utils.rs b/src/alerts/alerts_utils.rs index 611fd458d..1d95133ec 100644 --- a/src/alerts/alerts_utils.rs +++ b/src/alerts/alerts_utils.rs @@ -40,6 +40,7 @@ use crate::{ option::Mode, parseable::PARSEABLE, query::{QUERY_SESSION, execute, resolve_stream_names}, + tenants::TENANT_METADATA, utils::time::TimeRange, }; @@ -57,6 +58,14 @@ use super::{ALERTS, AlertError, AlertOperator, AlertState}; /// /// check whether notification needs to be triggered or not pub async fn evaluate_alert(alert: &dyn AlertTrait) -> Result<(), AlertError> { + if alert + .get_tenant_id() + .as_deref() + .is_some_and(|tenant| TENANT_METADATA.is_workspace_suspended(tenant)) + { + return Ok(()); + } + trace!("RUNNING EVAL TASK FOR- {alert:?}"); let message = alert.eval_alert().await?; diff --git a/src/hottier.rs b/src/hottier.rs index 95c1361b4..be2f0a3b2 100644 --- a/src/hottier.rs +++ b/src/hottier.rs @@ -22,8 +22,8 @@ use std::{ io, path::{Path, PathBuf}, sync::{ - Arc, OnceLock, - atomic::{AtomicUsize, Ordering}, + Arc, Mutex as StdMutex, OnceLock, + atomic::{AtomicU64, AtomicUsize, Ordering}, }, }; use tokio::sync::{Mutex as AsyncMutex, RwLock as AsyncRwLock, mpsc}; @@ -42,6 +42,7 @@ use crate::{ tenants::TENANT_METADATA, utils::{ disk::{DiskUtil, disk_usage_for_path}, + extract_datetime, human_size::bytes_to_human_size, }, validator::error::HotTierValidationError, @@ -260,6 +261,54 @@ pub struct StreamHotTier { /// Per-stream in-memory bookkeeping. Downloads run outside the lock. struct StreamSyncState { runtime: AsyncMutex, + eviction: Arc>, + query_pins: StdMutex>, + next_query_pin: AtomicU64, +} + +struct QueryPin { + start: DateTime, + end: DateTime, +} + +impl QueryPin { + fn contains(&self, minute: &str) -> bool { + extract_datetime(minute) + .map(|timestamp| timestamp.and_utc()) + .is_none_or(|bucket_start| { + let bucket_end = bucket_start + chrono::Duration::minutes(1); + bucket_end > self.start && bucket_start <= self.end + }) + } +} + +pub(crate) struct HotTierQueryGuard { + state: Arc, + pin_id: u64, +} + +impl Drop for HotTierQueryGuard { + fn drop(&mut self) { + self.state + .query_pins + .lock() + .expect("hot-tier query pins lock poisoned") + .remove(&self.pin_id); + } +} + +impl StreamSyncState { + fn pin_query(self: &Arc, start: DateTime, end: DateTime) -> HotTierQueryGuard { + let pin_id = self.next_query_pin.fetch_add(1, Ordering::Relaxed); + self.query_pins + .lock() + .expect("hot-tier query pins lock poisoned") + .insert(pin_id, QueryPin { start, end }); + HotTierQueryGuard { + state: self.clone(), + pin_id, + } + } } #[derive(Clone)] @@ -539,6 +588,9 @@ impl HotTierManager { let runtime = load_or_rebuild(&stream_root).await?; let state = Arc::new(StreamSyncState { runtime: AsyncMutex::new(runtime), + eviction: Arc::new(AsyncRwLock::new(())), + query_pins: StdMutex::new(HashMap::new()), + next_query_pin: AtomicU64::new(0), }); self.state_cache .write() @@ -560,6 +612,21 @@ impl HotTierManager { hot_tier_disk_path(self.hot_tier_path, manifest_path) } + /// Prevent eviction of buckets in the queried time range while local files are in use. + pub(crate) async fn query_guard( + &self, + stream: &str, + tenant_id: &Option, + start: DateTime, + end: DateTime, + ) -> Result { + let state = self.get_or_load_state(stream, tenant_id).await?; + let registration_guard = state.eviction.read().await; + let query_guard = state.pin_query(start, end); + drop(registration_guard); + Ok(query_guard) + } + /// Drop cached state for a stream (used after delete). pub async fn invalidate_state(&self, stream: &str, tenant_id: &Option) { let key: StreamKey = (tenant_id.clone(), stream.to_owned()); @@ -884,6 +951,13 @@ impl HotTierManager { tenant_id: Option, anchor: DateTime, ) -> Result<(), HotTierError> { + if tenant_id + .as_deref() + .is_some_and(|tenant| TENANT_METADATA.is_workspace_suspended(tenant)) + { + return Ok(()); + } + let stream_start = std::time::Instant::now(); self.process_manifest(&stream, &tenant_id, anchor) .await @@ -1292,11 +1366,27 @@ impl HotTierManager { ); let mut reclaim = ReclaimBudget::new(reclaim_target); let mut evicted = 0_u64; + let _eviction_guard = if reclaim.needs_more() { + // Query registrations hold read access only long enough to publish their time range. + // Waiting here cannot wait for query execution, and writer priority prevents a steady + // stream of new registrations from starving eviction. + Some(state.eviction.write().await) + } else { + None + }; while reclaim.needs_more() { - let Some(oldest) = runtime - .oldest_evictable_bucket_before(minute) - .map(str::to_owned) - else { + let oldest = { + let query_pins = state + .query_pins + .lock() + .expect("hot-tier query pins lock poisoned"); + runtime + .oldest_evictable_bucket_before_with(minute, |candidate| { + query_pins.values().any(|pin| pin.contains(candidate)) + }) + .map(str::to_owned) + }; + let Some(oldest) = oldest else { runtime.unmark_bucket_inflight(minute); return Ok(None); }; @@ -1758,7 +1848,7 @@ mod tests { use super::local_state::{MinuteTotals, RuntimeState}; use super::{ - DiskBudget, DiskUtil, HotTierManager, ReclaimBudget, StreamSyncState, WorkItem, + DiskBudget, DiskUtil, HotTierManager, QueryPin, ReclaimBudget, StreamSyncState, WorkItem, classify_manifest_files_with, hot_tier_disk_path, manifest_check_concurrency, required_reclaim, run_bounded_newest_first, }; @@ -1806,6 +1896,16 @@ mod tests { assert!(hot_tier_disk_path(root.path(), "../outside.parquet").is_err()); } + #[test] + fn query_pin_treats_unparseable_bucket_as_overlapping() { + let pin = QueryPin { + start: Utc.with_ymd_and_hms(2026, 7, 16, 12, 0, 30).unwrap(), + end: Utc.with_ymd_and_hms(2026, 7, 16, 12, 5, 0).unwrap(), + }; + + assert!(pin.contains("invalid-minute-bucket")); + } + fn disk_budget() -> DiskBudget { DiskBudget::new( Some(DiskUtil { @@ -1903,6 +2003,9 @@ mod tests { } let state = Arc::new(StreamSyncState { runtime: tokio::sync::Mutex::new(runtime), + eviction: Arc::new(tokio::sync::RwLock::new(())), + query_pins: std::sync::Mutex::new(std::collections::HashMap::new()), + next_query_pin: std::sync::atomic::AtomicU64::new(0), }); let item = WorkItem { timestamp: Utc::now(), @@ -1931,6 +2034,138 @@ mod tests { assert!(stream_root.join(next).exists()); } + #[tokio::test] + async fn active_query_protects_overlapping_bucket() { + let root = Box::leak(tempfile::tempdir().unwrap().keep().into_boxed_path()); + let stream_root = root.join("logs"); + let oldest = "date=2026-07-16/hour=12/minute=00"; + let oldest_directory = stream_root.join(oldest); + tokio::fs::create_dir_all(&oldest_directory).await.unwrap(); + tokio::fs::write(oldest_directory.join("cached.parquet"), [0_u8; 100]) + .await + .unwrap(); + + let mut runtime = RuntimeState::default(); + runtime.minutes.insert( + oldest.to_owned(), + MinuteTotals { + bytes: 100, + files: 1, + verified: true, + }, + ); + let state = Arc::new(StreamSyncState { + runtime: tokio::sync::Mutex::new(runtime), + eviction: Arc::new(tokio::sync::RwLock::new(())), + query_pins: std::sync::Mutex::new(std::collections::HashMap::new()), + next_query_pin: std::sync::atomic::AtomicU64::new(0), + }); + let query_guard = state.pin_query( + Utc.with_ymd_and_hms(2026, 7, 16, 12, 0, 30).unwrap(), + Utc.with_ymd_and_hms(2026, 7, 16, 12, 0, 59).unwrap(), + ); + let budget = disk_budget(); + let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel(); + let manager = HotTierManager::new(root, sender); + let item = WorkItem { + timestamp: Utc::now(), + minute_path: stream_root.join("date=2026-07-16/hour=12/minute=01"), + local_path: stream_root.join("date=2026-07-16/hour=12/minute=01/new.parquet"), + file: manifest_file("logs/date=2026-07-16/hour=12/minute=01/new.parquet", 100), + }; + + let evicted = manager + .reserve_item( + &state, + &stream_root, + "date=2026-07-16/hour=12/minute=01", + &item, + 100, + &budget, + ) + .await + .unwrap(); + + assert_eq!(evicted, None); + assert!(oldest_directory.exists()); + + drop(query_guard); + let evicted = manager + .reserve_item( + &state, + &stream_root, + "date=2026-07-16/hour=12/minute=01", + &item, + 100, + &budget, + ) + .await + .unwrap(); + assert_eq!(evicted, Some(100)); + assert!(!oldest_directory.exists()); + } + + #[tokio::test] + async fn active_recent_query_does_not_block_old_bucket_eviction() { + let root = Box::leak(tempfile::tempdir().unwrap().keep().into_boxed_path()); + let stream_root = root.join("logs"); + let oldest = "date=2026-07-16/hour=12/minute=00"; + let pinned = "date=2026-07-16/hour=12/minute=01"; + for minute in [oldest, pinned] { + let directory = stream_root.join(minute); + tokio::fs::create_dir_all(&directory).await.unwrap(); + tokio::fs::write(directory.join("cached.parquet"), [0_u8; 100]) + .await + .unwrap(); + } + + let mut runtime = RuntimeState::default(); + for minute in [oldest, pinned] { + runtime.minutes.insert( + minute.to_owned(), + MinuteTotals { + bytes: 100, + files: 1, + verified: true, + }, + ); + } + let state = Arc::new(StreamSyncState { + runtime: tokio::sync::Mutex::new(runtime), + eviction: Arc::new(tokio::sync::RwLock::new(())), + query_pins: std::sync::Mutex::new(std::collections::HashMap::new()), + next_query_pin: std::sync::atomic::AtomicU64::new(0), + }); + let _query_guard = state.pin_query( + Utc.with_ymd_and_hms(2026, 7, 16, 12, 1, 30).unwrap(), + Utc.with_ymd_and_hms(2026, 7, 16, 12, 1, 59).unwrap(), + ); + let item = WorkItem { + timestamp: Utc::now(), + minute_path: stream_root.join("date=2026-07-16/hour=12/minute=02"), + local_path: stream_root.join("date=2026-07-16/hour=12/minute=02/new.parquet"), + file: manifest_file("logs/date=2026-07-16/hour=12/minute=02/new.parquet", 100), + }; + let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel(); + let manager = HotTierManager::new(root, sender); + + let evicted = manager + .reserve_item( + &state, + &stream_root, + "date=2026-07-16/hour=12/minute=02", + &item, + 200, + &disk_budget(), + ) + .await + .unwrap(); + + assert_eq!(evicted, Some(100)); + assert!(!stream_root.join(oldest).exists()); + assert!(stream_root.join(pinned).exists()); + } + async fn wait_until_len(values: &Arc>>, expected: usize) { tokio::time::timeout(std::time::Duration::from_secs(2), async { while values.lock().unwrap().len() < expected { diff --git a/src/hottier/local_state.rs b/src/hottier/local_state.rs index 8f5060ab4..32951ca3d 100644 --- a/src/hottier/local_state.rs +++ b/src/hottier/local_state.rs @@ -85,11 +85,17 @@ impl RuntimeState { } } - pub fn oldest_evictable_bucket_before(&self, target: &str) -> Option<&str> { + pub fn oldest_evictable_bucket_before_with( + &self, + target: &str, + is_pinned: impl Fn(&str) -> bool, + ) -> Option<&str> { self.minutes .keys() .take_while(|minute| minute.as_str() < target) - .find(|minute| !self.inflight_buckets.contains_key(*minute)) + .find(|minute| { + !self.inflight_buckets.contains_key(*minute) && !is_pinned(minute.as_str()) + }) .map(String::as_str) } @@ -308,7 +314,9 @@ mod tests { state.mark_bucket_inflight("date=2026-07-16/hour=12/minute=00"); assert_eq!( - state.oldest_evictable_bucket_before("date=2026-07-16/hour=12/minute=14"), + state.oldest_evictable_bucket_before_with("date=2026-07-16/hour=12/minute=14", |_| { + false + },), Some("date=2026-07-16/hour=12/minute=05") ); } @@ -318,7 +326,9 @@ mod tests { let state = runtime_with_minutes(&["minute=14"]); assert_eq!( - state.oldest_evictable_bucket_before("date=2026-07-16/hour=12/minute=00"), + state.oldest_evictable_bucket_before_with("date=2026-07-16/hour=12/minute=00", |_| { + false + },), None ); } @@ -328,7 +338,9 @@ mod tests { let state = runtime_with_minutes(&["minute=14"]); assert_eq!( - state.oldest_evictable_bucket_before("date=2026-07-16/hour=12/minute=14"), + state.oldest_evictable_bucket_before_with("date=2026-07-16/hour=12/minute=14", |_| { + false + },), None ); } diff --git a/src/query/mod.rs b/src/query/mod.rs index fbd573423..33efeb30a 100644 --- a/src/query/mod.rs +++ b/src/query/mod.rs @@ -26,7 +26,8 @@ use chrono::NaiveDateTime; use chrono::{DateTime, Duration, Utc}; use datafusion::arrow::record_batch::RecordBatch; use datafusion::catalog::SchemaProvider; -use datafusion::common::tree_node::Transformed; +use datafusion::common::tree_node::{Transformed, TreeNodeRecursion}; +use datafusion::error::DataFusionError; use datafusion::execution::disk_manager::DiskManager; use datafusion::execution::{ RecordBatchStream, SendableRecordBatchStream, SessionState, SessionStateBuilder, @@ -49,6 +50,7 @@ use itertools::Itertools; use once_cell::sync::{Lazy, OnceCell}; use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; +use std::collections::BTreeSet; use std::ops::Bound; use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -70,6 +72,7 @@ use crate::catalog::manifest::Manifest; use crate::catalog::snapshot::Snapshot; use crate::event::DEFAULT_TIMESTAMP_KEY; use crate::handlers::http::query::QueryError; +use crate::hottier::{GLOBAL_HOTTIER, HotTierQueryGuard}; use crate::metrics::increment_bytes_scanned_in_query_by_date; use crate::option::Mode; use crate::parseable::{DEFAULT_TENANT, PARSEABLE}; @@ -371,9 +374,10 @@ impl Query { )] pub async fn execute(&self, is_streaming: bool, tenant_id: &Option) -> QueryResult { let ctx = QUERY_SESSION.get_ctx(); - let df = ctx - .execute_logical_plan(self.final_logical_plan(tenant_id)) - .await?; + let logical_plan = self.final_logical_plan(tenant_id); + let hot_tier_guards = + hot_tier_query_guards(&logical_plan, tenant_id, &self.time_range).await?; + let df = ctx.execute_logical_plan(logical_plan).await?; let tenant = tenant_id.as_deref().unwrap_or(DEFAULT_TENANT); let fields = df .schema() @@ -404,6 +408,7 @@ impl Query { let current_date = chrono::Utc::now().date_naive().to_string(); increment_bytes_scanned_in_query_by_date(actual_io_bytes, ¤t_date, tenant); + drop(hot_tier_guards); Either::Left(batches) } else { let task_ctx = ctx.task_ctx(); @@ -413,6 +418,7 @@ impl Query { let monitor_state = Arc::new(MonitorState { plan: plan.clone(), active_streams: AtomicUsize::new(output_partitions), + _hot_tier_guards: hot_tier_guards, }); let partition_streams = execute_stream_partitioned(plan.clone(), task_ctx.clone())?; @@ -540,6 +546,39 @@ impl Query { } } +/// Pins the queried time range for every hot-tier stream referenced by the query, including +/// subqueries, from physical planning through execution. Eviction remains free to reclaim buckets +/// outside active query ranges, so long-running or overlapping queries cannot starve hot-tier sync. +async fn hot_tier_query_guards( + logical_plan: &LogicalPlan, + tenant_id: &Option, + time_range: &TimeRange, +) -> Result, DataFusionError> { + let Some(manager) = GLOBAL_HOTTIER.get() else { + return Ok(Vec::new()); + }; + + let mut streams = BTreeSet::new(); + logical_plan.apply_with_subqueries(|plan| { + if let LogicalPlan::TableScan(table) = plan { + streams.insert(table.table_name.table().to_owned()); + } + Ok(TreeNodeRecursion::Continue) + })?; + + let mut guards = Vec::new(); + for stream in streams { + if manager.check_stream_hot_tier_exists(&stream, tenant_id) { + let guard = manager + .query_guard(&stream, tenant_id, time_range.start, time_range.end) + .await + .map_err(|error| DataFusionError::External(Box::new(error)))?; + guards.push(guard); + } + } + Ok(guards) +} + /// Recursively sums up "bytes_scanned" from all nodes in the plan fn get_total_bytes_scanned(plan: &Arc) -> u64 { let mut total_bytes = 0; @@ -1077,6 +1116,7 @@ pub mod error { struct MonitorState { plan: Arc, active_streams: AtomicUsize, + _hot_tier_guards: Vec, } /// A wrapper that monitors the ExecutionPlan and logs metrics when the stream finishes. diff --git a/src/storage/retention.rs b/src/storage/retention.rs index 2fafd6dd0..64faf95a9 100644 --- a/src/storage/retention.rs +++ b/src/storage/retention.rs @@ -29,7 +29,7 @@ use once_cell::sync::Lazy; use tokio::task::JoinHandle; use tracing::{info, warn}; -use crate::parseable::PARSEABLE; +use crate::{parseable::PARSEABLE, tenants::TENANT_METADATA}; type SchedulerHandle = JoinHandle<()>; @@ -51,6 +51,13 @@ pub fn init_scheduler() { vec![None] }; for tenant_id in tenants { + if tenant_id + .as_deref() + .is_some_and(|tenant| TENANT_METADATA.is_workspace_suspended(tenant)) + { + continue; + } + for stream_name in PARSEABLE.streams.list(&tenant_id) { match PARSEABLE.get_stream(&stream_name, &tenant_id) { Ok(stream) => { @@ -202,6 +209,7 @@ impl From for Vec { mod action { use crate::catalog::remove_manifest_from_snapshot; use crate::parseable::PARSEABLE; + use crate::tenants::TENANT_METADATA; use chrono::{Days, NaiveDate, Utc}; use futures::{StreamExt, stream::FuturesUnordered}; use itertools::Itertools; @@ -209,6 +217,13 @@ mod action { use tracing::{error, info}; pub(super) async fn delete(stream_name: String, days: u32, tenant_id: &Option) { + if tenant_id + .as_deref() + .is_some_and(|tenant| TENANT_METADATA.is_workspace_suspended(tenant)) + { + return; + } + info!("running retention task - delete for stream={stream_name}"); let store = PARSEABLE.storage.get_object_store(); From 6f2ee9b330ddbb14bf1fb05788a11d5dc1ab3669 Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Mon, 28 Sep 2026 14:38:52 +0530 Subject: [PATCH 2/6] serve query from object store if query_guard fails --- src/query/mod.rs | 77 +++++++++++++++++++++----- src/query/stream_schema_provider.rs | 86 ++++++++++++++++++++++++++--- 2 files changed, 141 insertions(+), 22 deletions(-) diff --git a/src/query/mod.rs b/src/query/mod.rs index 33efeb30a..2bbc4a6fe 100644 --- a/src/query/mod.rs +++ b/src/query/mod.rs @@ -62,8 +62,10 @@ use tokio_stream::wrappers::UnboundedReceiverStream; use tracing::Instrument; use self::error::ExecuteError; -use self::stream_schema_provider::GlobalSchemaProvider; pub use self::stream_schema_provider::PartialTimeFilter; +use self::stream_schema_provider::{ + GlobalSchemaProvider, HotTierStreamKey, guarded_hot_tier_table_source, hot_tier_stream_key, +}; use crate::alerts::alert_structs::Conditions; use crate::alerts::alerts_utils::get_filter_string; use crate::catalog::Snapshot as CatalogSnapshot; @@ -374,9 +376,12 @@ impl Query { )] pub async fn execute(&self, is_streaming: bool, tenant_id: &Option) -> QueryResult { let ctx = QUERY_SESSION.get_ctx(); - let logical_plan = self.final_logical_plan(tenant_id); - let hot_tier_guards = - hot_tier_query_guards(&logical_plan, tenant_id, &self.time_range).await?; + let mut logical_plan = self.final_logical_plan(tenant_id); + let (hot_tier_guards, guarded_hot_tier_streams) = + hot_tier_query_guards(&logical_plan, &self.time_range).await?; + if !guarded_hot_tier_streams.is_empty() { + logical_plan = enable_hot_tier_reads(logical_plan, &guarded_hot_tier_streams)?; + } let df = ctx.execute_logical_plan(logical_plan).await?; let tenant = tenant_id.as_deref().unwrap_or(DEFAULT_TENANT); let fields = df @@ -551,32 +556,74 @@ impl Query { /// outside active query ranges, so long-running or overlapping queries cannot starve hot-tier sync. async fn hot_tier_query_guards( logical_plan: &LogicalPlan, - tenant_id: &Option, time_range: &TimeRange, -) -> Result, DataFusionError> { +) -> Result<(Vec, BTreeSet), DataFusionError> { let Some(manager) = GLOBAL_HOTTIER.get() else { - return Ok(Vec::new()); + return Ok((Vec::new(), BTreeSet::new())); }; let mut streams = BTreeSet::new(); logical_plan.apply_with_subqueries(|plan| { - if let LogicalPlan::TableScan(table) = plan { - streams.insert(table.table_name.table().to_owned()); + if let LogicalPlan::TableScan(table) = plan + && let Some(key) = hot_tier_stream_key(&table.source) + { + streams.insert(key); } Ok(TreeNodeRecursion::Continue) })?; let mut guards = Vec::new(); + let mut guarded_streams = BTreeSet::new(); for stream in streams { - if manager.check_stream_hot_tier_exists(&stream, tenant_id) { - let guard = manager - .query_guard(&stream, tenant_id, time_range.start, time_range.end) + if manager.check_stream_hot_tier_exists(&stream.stream, &stream.tenant_id) { + match manager + .query_guard( + &stream.stream, + &stream.tenant_id, + time_range.start, + time_range.end, + ) .await - .map_err(|error| DataFusionError::External(Box::new(error)))?; - guards.push(guard); + { + Ok(guard) => { + guards.push(guard); + guarded_streams.insert(stream); + } + Err(error) => { + tracing::warn!( + stream = %stream.stream, + tenant = ?stream.tenant_id, + %error, + "hot-tier query guard unavailable; using object storage" + ); + } + } } } - Ok(guards) + Ok((guards, guarded_streams)) +} + +fn enable_hot_tier_reads( + plan: LogicalPlan, + guarded_streams: &BTreeSet, +) -> Result { + plan.transform_up_with_subqueries(|plan| match plan { + LogicalPlan::TableScan(mut table) => { + let Some(key) = hot_tier_stream_key(&table.source) else { + return Ok(Transformed::no(LogicalPlan::TableScan(table))); + }; + if !guarded_streams.contains(&key) { + return Ok(Transformed::no(LogicalPlan::TableScan(table))); + } + let Some(source) = guarded_hot_tier_table_source(&table.source) else { + return Ok(Transformed::no(LogicalPlan::TableScan(table))); + }; + table.source = source; + Ok(Transformed::yes(LogicalPlan::TableScan(table))) + } + _ => Ok(Transformed::no(plan)), + }) + .map(|transformed| transformed.data) } /// Recursively sums up "bytes_scanned" from all nodes in the plan diff --git a/src/query/stream_schema_provider.rs b/src/query/stream_schema_provider.rs index 5c737f084..a1891cccf 100644 --- a/src/query/stream_schema_provider.rs +++ b/src/query/stream_schema_provider.rs @@ -16,7 +16,7 @@ * */ -use std::{cmp::Reverse, collections::HashMap, ops::Bound, sync::Arc}; +use std::{any::Any, cmp::Reverse, collections::HashMap, ops::Bound, sync::Arc}; use arrow_array::RecordBatch; use arrow_schema::{Schema, SchemaRef, SortOptions}; @@ -29,15 +29,16 @@ use datafusion::{ tree_node::{TreeNode, TreeNodeRecursion}, }, datasource::{ - MemTable, TableProvider, + DefaultTableSource, MemTable, TableProvider, file_format::{FileFormat, parquet::ParquetFormat}, listing::PartitionedFile, physical_plan::{FileGroup, FileScanConfigBuilder, ParquetSource}, + provider_as_source, }, error::{DataFusionError, Result as DataFusionResult}, execution::object_store::ObjectStoreUrl, logical_expr::{ - BinaryExpr, Operator, TableProviderFilterPushDown, TableType, + BinaryExpr, Operator, TableProviderFilterPushDown, TableSource, TableType, physical_planning_context::PhysicalPlanningContext, utils::conjunction, }, physical_expr::{LexOrdering, PhysicalSortExpr, create_physical_expr, expressions::col}, @@ -97,6 +98,7 @@ impl SchemaProvider for GlobalSchemaProvider { .get_schema(), stream: name.to_owned(), tenant_id: self.tenant_id.clone(), + hot_tier_guarded: false, }))) } else { Ok(None) @@ -114,6 +116,39 @@ struct StandardTableProvider { // prefix under which to find snapshot stream: String, tenant_id: Option, + hot_tier_guarded: bool, +} + +#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub(super) struct HotTierStreamKey { + pub tenant_id: Option, + pub stream: String, +} + +fn standard_table_provider(source: &Arc) -> Option<&StandardTableProvider> { + let source = (source.as_ref() as &dyn Any).downcast_ref::()?; + (source.table_provider.as_ref() as &dyn Any).downcast_ref::() +} + +pub(super) fn hot_tier_stream_key(source: &Arc) -> Option { + let provider = standard_table_provider(source)?; + Some(HotTierStreamKey { + tenant_id: provider.tenant_id.clone(), + stream: provider.stream.clone(), + }) +} + +/// Creates a query-specific source that may read hot-tier files because its guard is held. +pub(super) fn guarded_hot_tier_table_source( + source: &Arc, +) -> Option> { + let provider = standard_table_provider(source)?; + Some(provider_as_source(Arc::new(StandardTableProvider { + schema: provider.schema.clone(), + stream: provider.stream.clone(), + tenant_id: provider.tenant_id.clone(), + hot_tier_guarded: true, + }))) } pub fn exact_source_filters(filters: &[Expr]) -> Vec { @@ -830,7 +865,8 @@ impl TableProvider for StandardTableProvider { } // Hot tier data fetch - if let Some(hot_tier_manager) = GLOBAL_HOTTIER.get() + if self.hot_tier_guarded + && let Some(hot_tier_manager) = GLOBAL_HOTTIER.get() && hot_tier_manager.check_stream_hot_tier_exists(&self.stream, &self.tenant_id) { self.get_hottier_exectuion_plan( @@ -1275,8 +1311,10 @@ mod tests { use datafusion::{ datasource::{ TableProvider, + empty::EmptyTable, listing::PartitionedFile, physical_plan::{FileScanConfigBuilder, FileSource}, + provider_as_source, source::DataSourceExec, }, execution::context::{SessionConfig, SessionContext}, @@ -1294,11 +1332,45 @@ mod tests { }; use super::{ - PartialTimeFilter, balanced_file_groups, build_parquet_scan_components, - build_parquet_scan_components_with_full_filters, exact_source_filters, - extract_timestamp_bound, file_groups_are_time_ordered, is_overlapping_query, + HotTierStreamKey, PartialTimeFilter, StandardTableProvider, balanced_file_groups, + build_parquet_scan_components, build_parquet_scan_components_with_full_filters, + exact_source_filters, extract_timestamp_bound, file_groups_are_time_ordered, + guarded_hot_tier_table_source, hot_tier_stream_key, is_overlapping_query, + standard_table_provider, }; + #[test] + fn guarded_hot_tier_source_preserves_provider_identity() { + let source = provider_as_source(Arc::new(StandardTableProvider { + schema: Arc::new(Schema::empty()), + stream: "logs".to_owned(), + tenant_id: Some("other-tenant".to_owned()), + hot_tier_guarded: false, + })); + + assert_eq!( + hot_tier_stream_key(&source), + Some(HotTierStreamKey { + tenant_id: Some("other-tenant".to_owned()), + stream: "logs".to_owned(), + }) + ); + + let guarded_source = guarded_hot_tier_table_source(&source).unwrap(); + let guarded_provider = standard_table_provider(&guarded_source).unwrap(); + assert_eq!(guarded_provider.stream, "logs"); + assert_eq!(guarded_provider.tenant_id.as_deref(), Some("other-tenant")); + assert!(guarded_provider.hot_tier_guarded); + } + + #[test] + fn hot_tier_source_rewrite_ignores_non_standard_providers() { + let source = provider_as_source(Arc::new(EmptyTable::new(Arc::new(Schema::empty())))); + + assert!(hot_tier_stream_key(&source).is_none()); + assert!(guarded_hot_tier_table_source(&source).is_none()); + } + fn file(path: &str, size: u64) -> File { File { file_path: path.to_owned(), From 2d0c476ebc874060cf797913963fe349b03b83ed Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Mon, 28 Sep 2026 16:43:13 +0530 Subject: [PATCH 3/6] reuse in enterprise --- src/hottier.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/hottier.rs b/src/hottier.rs index be2f0a3b2..744ab9267 100644 --- a/src/hottier.rs +++ b/src/hottier.rs @@ -282,7 +282,8 @@ impl QueryPin { } } -pub(crate) struct HotTierQueryGuard { +/// Keeps hot-tier buckets needed by an active query from being evicted. +pub struct HotTierQueryGuard { state: Arc, pin_id: u64, } @@ -613,7 +614,7 @@ impl HotTierManager { } /// Prevent eviction of buckets in the queried time range while local files are in use. - pub(crate) async fn query_guard( + pub async fn query_guard( &self, stream: &str, tenant_id: &Option, From 287ff52fd7ee9418909d1c862dc2ffa5f13fbcad Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Mon, 28 Sep 2026 17:20:19 +0530 Subject: [PATCH 4/6] reconciliation and failed-download cleanup to respect query pins --- src/hottier.rs | 110 +++++++++++++++++++++++++++++++++++++++-- src/hottier/planner.rs | 19 +++---- 2 files changed, 113 insertions(+), 16 deletions(-) diff --git a/src/hottier.rs b/src/hottier.rs index 744ab9267..02608f024 100644 --- a/src/hottier.rs +++ b/src/hottier.rs @@ -64,7 +64,7 @@ mod planner; use local_state::{ RuntimeState, cleanup_stale_partials, load_or_rebuild, persist_checkpoint, verify_bucket, }; -use planner::{WorkItem, build_work, reconcile_local_file}; +use planner::{WorkItem, build_work, local_file_matches_expected_size}; async fn run_bounded_newest_first( work: Vec, @@ -299,6 +299,14 @@ impl Drop for HotTierQueryGuard { } impl StreamSyncState { + fn bucket_is_pinned(&self, minute: &str) -> bool { + self.query_pins + .lock() + .expect("hot-tier query pins lock poisoned") + .values() + .any(|pin| pin.contains(minute)) + } + fn pin_query(self: &Arc, start: DateTime, end: DateTime) -> HotTierQueryGuard { let pin_id = self.next_query_pin.fetch_add(1, Ordering::Relaxed); self.query_pins @@ -1019,7 +1027,7 @@ impl HotTierManager { &s3_manifests, latest_minutes, cache_root, - reconcile_local_file, + local_file_matches_expected_size, ) }) .await @@ -1328,7 +1336,8 @@ impl HotTierManager { .await .is_ok_and(|metadata| metadata.len() == item.file.file_size); if !valid { - let _ = fs::remove_file(&item.local_path).await; + self.remove_file_if_unpinned(&context.state, &minute, &item.local_path) + .await; } self.finish_reservation(&context.state, &context.disk_budget, &item, &minute, valid) .await; @@ -1358,7 +1367,32 @@ impl HotTierManager { if item.file.file_size > quota { return Ok(None); } + let removed_wrong_sized_file = match fs::metadata(&item.local_path).await { + Ok(metadata) if metadata.len() == item.file.file_size => return Ok(None), + Ok(_) => { + let _reconciliation_guard = state.eviction.write().await; + if state.bucket_is_pinned(minute) { + return Ok(None); + } + match fs::metadata(&item.local_path).await { + Ok(metadata) if metadata.len() == item.file.file_size => return Ok(None), + Ok(_) => { + fs::remove_file(&item.local_path).await?; + true + } + Err(error) if error.kind() == io::ErrorKind::NotFound => false, + Err(error) => return Err(error.into()), + } + } + Err(error) if error.kind() == io::ErrorKind::NotFound => false, + Err(error) => return Err(error.into()), + }; let mut runtime = state.runtime.lock().await; + if removed_wrong_sized_file { + runtime + .minutes + .insert(minute.to_owned(), verify_bucket(stream_root, minute).await?); + } runtime.mark_bucket_inflight(minute); let reclaim_target = required_reclaim( runtime.free_bytes(quota), @@ -1426,6 +1460,18 @@ impl HotTierManager { Ok(None) } + async fn remove_file_if_unpinned( + &self, + state: &Arc, + minute: &str, + path: &Path, + ) { + let _reconciliation_guard = state.eviction.write().await; + if !state.bucket_is_pinned(minute) { + let _ = fs::remove_file(path).await; + } + } + async fn finish_reservation( &self, state: &Arc, @@ -2106,6 +2152,64 @@ mod tests { assert!(!oldest_directory.exists()); } + #[tokio::test] + async fn active_query_protects_wrong_sized_file_during_reconciliation() { + let root = Box::leak(tempfile::tempdir().unwrap().keep().into_boxed_path()); + let stream_root = root.join("logs"); + let minute = "date=2026-07-16/hour=12/minute=00"; + let local_path = stream_root.join(minute).join("cached.parquet"); + tokio::fs::create_dir_all(local_path.parent().unwrap()) + .await + .unwrap(); + tokio::fs::write(&local_path, [0_u8; 50]).await.unwrap(); + + let mut runtime = RuntimeState::default(); + runtime.minutes.insert( + minute.to_owned(), + MinuteTotals { + bytes: 50, + files: 1, + verified: true, + }, + ); + let state = Arc::new(StreamSyncState { + runtime: tokio::sync::Mutex::new(runtime), + eviction: Arc::new(tokio::sync::RwLock::new(())), + query_pins: std::sync::Mutex::new(std::collections::HashMap::new()), + next_query_pin: std::sync::atomic::AtomicU64::new(0), + }); + let query_guard = state.pin_query( + Utc.with_ymd_and_hms(2026, 7, 16, 12, 0, 0).unwrap(), + Utc.with_ymd_and_hms(2026, 7, 16, 12, 0, 59).unwrap(), + ); + let item = WorkItem { + timestamp: Utc::now(), + minute_path: stream_root.join(minute), + local_path: local_path.clone(), + file: manifest_file("logs/date=2026-07-16/hour=12/minute=00/cached.parquet", 100), + }; + let budget = disk_budget(); + let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel(); + let manager = HotTierManager::new(root, sender); + + let reservation = manager + .reserve_item(&state, &stream_root, minute, &item, 1_000, &budget) + .await + .unwrap(); + + assert_eq!(reservation, None); + assert_eq!(tokio::fs::metadata(&local_path).await.unwrap().len(), 50); + + drop(query_guard); + let reservation = manager + .reserve_item(&state, &stream_root, minute, &item, 1_000, &budget) + .await + .unwrap(); + + assert_eq!(reservation, Some(0)); + assert!(!local_path.exists()); + } + #[tokio::test] async fn active_recent_query_does_not_block_old_bucket_eviction() { let root = Box::leak(tempfile::tempdir().unwrap().keep().into_boxed_path()); diff --git a/src/hottier/planner.rs b/src/hottier/planner.rs index 9d753a286..4901d1940 100644 --- a/src/hottier/planner.rs +++ b/src/hottier/planner.rs @@ -118,15 +118,8 @@ pub(super) fn minute_ancestor(path: &Path) -> Option<&Path> { None } -pub(super) fn reconcile_local_file(path: &Path, expected_size: u64) -> bool { - match std::fs::metadata(path) { - Ok(metadata) if metadata.len() == expected_size => true, - Ok(_) => { - let _ = std::fs::remove_file(path); - false - } - Err(_) => false, - } +pub(super) fn local_file_matches_expected_size(path: &Path, expected_size: u64) -> bool { + std::fs::metadata(path).is_ok_and(|metadata| metadata.len() == expected_size) } #[cfg(test)] @@ -137,7 +130,7 @@ mod tests { use crate::catalog::manifest::{File, Manifest}; - use super::{build_work, minute_ancestor, reconcile_local_file}; + use super::{build_work, local_file_matches_expected_size, minute_ancestor}; fn file(path: &str, size: u64) -> File { File { @@ -206,12 +199,12 @@ mod tests { } #[test] - fn wrong_sized_local_file_is_removed_before_reservation() { + fn planning_does_not_remove_wrong_sized_local_file() { let temp = tempfile::tempdir().unwrap(); let path = temp.path().join("wrong.parquet"); std::fs::write(&path, [0_u8; 3]).unwrap(); - assert!(!reconcile_local_file(&path, 7)); - assert!(!path.exists()); + assert!(!local_file_matches_expected_size(&path, 7)); + assert!(path.exists()); } } From 7bd1b614720f4e4f25fa49ac4ee763c7d1a7ca55 Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Mon, 28 Sep 2026 17:40:02 +0530 Subject: [PATCH 5/6] coderabbit comments and deepsource fix --- src/hottier.rs | 73 +++++++++++++++++++++++++++++++++----------------- 1 file changed, 49 insertions(+), 24 deletions(-) diff --git a/src/hottier.rs b/src/hottier.rs index 02608f024..c483deffa 100644 --- a/src/hottier.rs +++ b/src/hottier.rs @@ -1367,25 +1367,11 @@ impl HotTierManager { if item.file.file_size > quota { return Ok(None); } - let removed_wrong_sized_file = match fs::metadata(&item.local_path).await { - Ok(metadata) if metadata.len() == item.file.file_size => return Ok(None), - Ok(_) => { - let _reconciliation_guard = state.eviction.write().await; - if state.bucket_is_pinned(minute) { - return Ok(None); - } - match fs::metadata(&item.local_path).await { - Ok(metadata) if metadata.len() == item.file.file_size => return Ok(None), - Ok(_) => { - fs::remove_file(&item.local_path).await?; - true - } - Err(error) if error.kind() == io::ErrorKind::NotFound => false, - Err(error) => return Err(error.into()), - } - } - Err(error) if error.kind() == io::ErrorKind::NotFound => false, - Err(error) => return Err(error.into()), + let Some(removed_wrong_sized_file) = self + .prepare_local_file_for_download(state, minute, item, disk_budget) + .await? + else { + return Ok(None); }; let mut runtime = state.runtime.lock().await; if removed_wrong_sized_file { @@ -1460,6 +1446,38 @@ impl HotTierManager { Ok(None) } + async fn prepare_local_file_for_download( + &self, + state: &Arc, + minute: &str, + item: &WorkItem, + disk_budget: &DiskBudget, + ) -> Result, HotTierError> { + let removed_wrong_sized_file = match fs::metadata(&item.local_path).await { + Ok(metadata) if metadata.len() == item.file.file_size => return Ok(None), + Ok(_) => { + let _reconciliation_guard = state.eviction.write().await; + if state.bucket_is_pinned(minute) { + return Ok(None); + } + match fs::metadata(&item.local_path).await { + Ok(metadata) if metadata.len() == item.file.file_size => return Ok(None), + Ok(metadata) => { + let removed_bytes = metadata.len(); + fs::remove_file(&item.local_path).await?; + disk_budget.credit_eviction(removed_bytes).await; + true + } + Err(error) if error.kind() == io::ErrorKind::NotFound => false, + Err(error) => return Err(error.into()), + } + } + Err(error) if error.kind() == io::ErrorKind::NotFound => false, + Err(error) => return Err(error.into()), + }; + Ok(Some(removed_wrong_sized_file)) + } + async fn remove_file_if_unpinned( &self, state: &Arc, @@ -2161,13 +2179,13 @@ mod tests { tokio::fs::create_dir_all(local_path.parent().unwrap()) .await .unwrap(); - tokio::fs::write(&local_path, [0_u8; 50]).await.unwrap(); + tokio::fs::write(&local_path, [0_u8; 100]).await.unwrap(); let mut runtime = RuntimeState::default(); runtime.minutes.insert( minute.to_owned(), MinuteTotals { - bytes: 50, + bytes: 100, files: 1, verified: true, }, @@ -2186,9 +2204,16 @@ mod tests { timestamp: Utc::now(), minute_path: stream_root.join(minute), local_path: local_path.clone(), - file: manifest_file("logs/date=2026-07-16/hour=12/minute=00/cached.parquet", 100), + file: manifest_file("logs/date=2026-07-16/hour=12/minute=00/cached.parquet", 50), }; - let budget = disk_budget(); + let budget = DiskBudget::new( + Some(DiskUtil { + total_space: 1_000, + available_space: 0, + used_space: 1_000, + }), + 100.0, + ); let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel(); let manager = HotTierManager::new(root, sender); @@ -2198,7 +2223,7 @@ mod tests { .unwrap(); assert_eq!(reservation, None); - assert_eq!(tokio::fs::metadata(&local_path).await.unwrap().len(), 50); + assert_eq!(tokio::fs::metadata(&local_path).await.unwrap().len(), 100); drop(query_guard); let reservation = manager From 797b6319b6bd3d2743fafd82dcc515e5cf202bba Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Mon, 28 Sep 2026 20:51:49 +0530 Subject: [PATCH 6/6] fix issue with unguarded queries --- src/query/mod.rs | 32 +++++++++++++++++++++++++++-- src/query/stream_schema_provider.rs | 2 +- 2 files changed, 31 insertions(+), 3 deletions(-) diff --git a/src/query/mod.rs b/src/query/mod.rs index 2bbc4a6fe..377931bb1 100644 --- a/src/query/mod.rs +++ b/src/query/mod.rs @@ -34,7 +34,8 @@ use datafusion::execution::{ }; use datafusion::logical_expr::expr::Alias; use datafusion::logical_expr::{ - Aggregate, Explain, Filter, LogicalPlan, PlanType, Projection, ScalarUDF, ToStringifiedPlan, + Aggregate, Explain, Filter, LogicalPlan, PlanType, Projection, ScalarUDF, TableSource, + ToStringifiedPlan, }; use datafusion::physical_plan::stream::RecordBatchStreamAdapter; use datafusion::physical_plan::{ @@ -64,7 +65,9 @@ use tracing::Instrument; use self::error::ExecuteError; pub use self::stream_schema_provider::PartialTimeFilter; use self::stream_schema_provider::{ - GlobalSchemaProvider, HotTierStreamKey, guarded_hot_tier_table_source, hot_tier_stream_key, + GlobalSchemaProvider, HotTierStreamKey, + guarded_hot_tier_table_source as default_guarded_hot_tier_table_source, + hot_tier_stream_key as default_hot_tier_stream_key, }; use crate::alerts::alert_structs::Conditions; use crate::alerts::alerts_utils::get_filter_string; @@ -117,6 +120,31 @@ pub trait ParseableSchemaProvider: Send + Sync { storage: Option>, tenant_id: &Option, ) -> Box; + + fn hot_tier_stream_key(&self, _source: &Arc) -> Option { + None + } + + fn guarded_hot_tier_table_source( + &self, + _source: &Arc, + ) -> Option> { + None + } +} + +fn hot_tier_stream_key(source: &Arc) -> Option { + SCHEMA_PROVIDER + .get() + .and_then(|provider| provider.hot_tier_stream_key(source)) + .or_else(|| default_hot_tier_stream_key(source)) +} + +fn guarded_hot_tier_table_source(source: &Arc) -> Option> { + SCHEMA_PROVIDER + .get() + .and_then(|provider| provider.guarded_hot_tier_table_source(source)) + .or_else(|| default_guarded_hot_tier_table_source(source)) } fn get_schema_provider(tenant_id: &Option) -> Box { diff --git a/src/query/stream_schema_provider.rs b/src/query/stream_schema_provider.rs index a1891cccf..c6d673d68 100644 --- a/src/query/stream_schema_provider.rs +++ b/src/query/stream_schema_provider.rs @@ -120,7 +120,7 @@ struct StandardTableProvider { } #[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)] -pub(super) struct HotTierStreamKey { +pub struct HotTierStreamKey { pub tenant_id: Option, pub stream: String, }