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..c483deffa 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, @@ -63,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, @@ -260,6 +261,63 @@ 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 + }) + } +} + +/// Keeps hot-tier buckets needed by an active query from being evicted. +pub 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 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 + .lock() + .expect("hot-tier query pins lock poisoned") + .insert(pin_id, QueryPin { start, end }); + HotTierQueryGuard { + state: self.clone(), + pin_id, + } + } } #[derive(Clone)] @@ -539,6 +597,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 +621,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 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 +960,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 @@ -944,7 +1027,7 @@ impl HotTierManager { &s3_manifests, latest_minutes, cache_root, - reconcile_local_file, + local_file_matches_expected_size, ) }) .await @@ -1253,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; @@ -1283,7 +1367,18 @@ impl HotTierManager { if item.file.file_size > quota { return Ok(None); } + 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 { + 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), @@ -1292,11 +1387,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); }; @@ -1335,6 +1446,50 @@ 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, + 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, @@ -1758,7 +1913,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 +1961,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 +2068,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 +2099,203 @@ 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_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; 100]).await.unwrap(); + + let mut runtime = RuntimeState::default(); + 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, 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", 50), + }; + 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); + + 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(), 100); + + 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()); + 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/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()); } } diff --git a/src/query/mod.rs b/src/query/mod.rs index fbd573423..377931bb1 100644 --- a/src/query/mod.rs +++ b/src/query/mod.rs @@ -26,14 +26,16 @@ 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, }; 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::{ @@ -49,6 +51,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}; @@ -60,8 +63,12 @@ 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 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; use crate::catalog::Snapshot as CatalogSnapshot; @@ -70,6 +77,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}; @@ -112,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 { @@ -371,9 +404,13 @@ 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 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 .schema() @@ -404,6 +441,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 +451,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 +579,81 @@ 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, + time_range: &TimeRange, +) -> Result<(Vec, BTreeSet), DataFusionError> { + let Some(manager) = GLOBAL_HOTTIER.get() else { + return Ok((Vec::new(), BTreeSet::new())); + }; + + let mut streams = BTreeSet::new(); + logical_plan.apply_with_subqueries(|plan| { + 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.stream, &stream.tenant_id) { + match manager + .query_guard( + &stream.stream, + &stream.tenant_id, + time_range.start, + time_range.end, + ) + .await + { + 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, 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 fn get_total_bytes_scanned(plan: &Arc) -> u64 { let mut total_bytes = 0; @@ -1077,6 +1191,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/query/stream_schema_provider.rs b/src/query/stream_schema_provider.rs index 5c737f084..c6d673d68 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 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(), 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();