diff --git a/crates/graphforge-api/src/algorithm_writeback.rs b/crates/graphforge-api/src/algorithm_writeback.rs index e010623c1..e7c064594 100644 --- a/crates/graphforge-api/src/algorithm_writeback.rs +++ b/crates/graphforge-api/src/algorithm_writeback.rs @@ -97,7 +97,7 @@ impl GraphForge { self.ontology_mode, self.adjacency_provider_for_session(), Some(std::sync::Arc::clone(&self.ordinal_identities)), - &self.session_resource_config(), + &self.session_resource_config()?, )?; let session = if self.read_only { session.restrict_to_reads() diff --git a/crates/graphforge-api/src/embedding_refresh.rs b/crates/graphforge-api/src/embedding_refresh.rs index 978ced91e..225676a8b 100644 --- a/crates/graphforge-api/src/embedding_refresh.rs +++ b/crates/graphforge-api/src/embedding_refresh.rs @@ -292,6 +292,7 @@ impl GraphForge { compute_pool: Arc::clone(&self.compute_pool), heavy_query_admission: Arc::clone(&self.heavy_query_admission), construction_cpu_admission: Arc::clone(&self.construction_cpu_admission), + query_spill: Arc::clone(&self.query_spill), provider_refresh_driver_active: Arc::clone(&self.provider_refresh_driver_active), provider_refresh_runtimes: Arc::clone(&self.provider_refresh_runtimes), provider_find_runtimes: Arc::clone(&self.provider_find_runtimes), diff --git a/crates/graphforge-api/src/explanation.rs b/crates/graphforge-api/src/explanation.rs index 7f73b7d6d..51077f669 100644 --- a/crates/graphforge-api/src/explanation.rs +++ b/crates/graphforge-api/src/explanation.rs @@ -132,7 +132,7 @@ impl GraphForge { execution_mode, adjacency_provider, Some(Arc::clone(&self.ordinal_identities)), - &self.session_resource_config(), + &self.planning_resource_config(), )?; // No executable plan or write capability escapes this rendering boundary. let physical = self.block_on(async move { session.explain_physical(&plan).await })?; diff --git a/crates/graphforge-api/src/lib.rs b/crates/graphforge-api/src/lib.rs index 1b2e7ba81..367558ea5 100644 --- a/crates/graphforge-api/src/lib.rs +++ b/crates/graphforge-api/src/lib.rs @@ -599,6 +599,10 @@ pub struct GraphForge { /// Instance-wide limit on parallel construction lanes (#1586, ADR 0047): /// `compute_threads - construction_cpu_reserve`, shared by every import. construction_cpu_admission: Arc, + /// This instance's project scratch directory for query spill, acquired on + /// the first query that may spill and released when the last facade + /// sharing it drops (#1595). + query_spill: Arc>>, /// Ensures mutation bursts share one bounded process-local driver thread. provider_refresh_driver_active: Arc, /// Runtime-only provider recipes capable of refreshing exact lineages. @@ -789,6 +793,7 @@ impl GraphForge { resource_policy.compute_threads, )?), construction_cpu_admission: resource_policy.construction_cpu_admission(), + query_spill: Arc::new(Mutex::new(None)), provider_refresh_driver_active: Arc::new(AtomicBool::new(false)), provider_refresh_runtimes: Arc::new(Mutex::new(Vec::new())), provider_find_runtimes: Arc::new(Mutex::new(Vec::new())), @@ -1023,6 +1028,7 @@ impl GraphForge { compute_pool, heavy_query_admission, construction_cpu_admission, + query_spill: Arc::new(Mutex::new(None)), provider_refresh_driver_active: Arc::new(AtomicBool::new(false)), provider_refresh_runtimes: Arc::new(Mutex::new(Vec::new())), provider_find_runtimes: Arc::new(Mutex::new(Vec::new())), diff --git a/crates/graphforge-api/src/query_execution.rs b/crates/graphforge-api/src/query_execution.rs index a7c584294..060d64e5c 100644 --- a/crates/graphforge-api/src/query_execution.rs +++ b/crates/graphforge-api/src/query_execution.rs @@ -444,7 +444,7 @@ impl GraphForge { execution_mode, adjacency_provider, Some(Arc::clone(&self.ordinal_identities)), - &self.session_resource_config(), + &self.session_resource_config()?, )?; let session = if self.read_only { session.restrict_to_reads() @@ -639,7 +639,7 @@ impl GraphForge { execution_mode, adjacency_provider, Some(Arc::clone(&self.ordinal_identities)), - &self.session_resource_config(), + &self.session_resource_config()?, )?; let session = if self.read_only { session.restrict_to_reads() diff --git a/crates/graphforge-api/src/resource_policy.rs b/crates/graphforge-api/src/resource_policy.rs index ebbf8b985..1e9a6c503 100644 --- a/crates/graphforge-api/src/resource_policy.rs +++ b/crates/graphforge-api/src/resource_policy.rs @@ -25,17 +25,35 @@ pub enum ResourcePolicyMode { } /// Fail-closed spill configuration. -#[derive(Clone, Debug, Default, Eq, PartialEq)] +/// +/// The default (#1595) is `enabled` with no `directory`: a durable project's +/// queries spill into its own scratch directory +/// (`graphforge_storage::query_spill`), capped at `max_bytes` or +/// `DEFAULT_QUERY_SPILL_MAX_BYTES` per query, and an in-memory instance does +/// not spill at all. `enabled: false` never spills: a query over its memory +/// budget fails with a resource error. +#[derive(Clone, Debug, Eq, PartialEq)] pub struct SpillPolicy { - /// When false, DataFusion must not spill to disk. + /// When false, queries never spill to disk. pub enabled: bool, - /// Optional spill directory. Relative paths are rejected; must be absolute - /// when enabled. Symlinks and non-directories fail closed at normalize time. + /// Optional spill directory. Relative paths are rejected. Symlinks and + /// non-directories fail closed at normalize time. `None` with `enabled` + /// selects the project scratch directory. pub directory: Option, - /// Optional upper bound on temporary spill bytes. + /// Optional upper bound on one query's temporary spill bytes. pub max_bytes: Option, } +impl Default for SpillPolicy { + fn default() -> Self { + Self { + enabled: true, + directory: None, + max_bytes: None, + } + } +} + /// Requested (pre-normalization) execution resource policy. #[derive(Clone, Debug, Eq, PartialEq)] pub struct ExecutionResourcePolicy { @@ -109,7 +127,8 @@ pub struct NormalizedResourcePolicy { pub memory_budget_bytes: u64, /// Whether spill is enabled. pub spill_enabled: bool, - /// Absolute spill directory when spill is enabled. + /// Caller-configured absolute spill directory. `None` with spill enabled + /// selects the project scratch directory (#1595). pub spill_directory: Option, /// Optional spill byte cap. pub spill_max_bytes: Option, @@ -335,18 +354,20 @@ impl ExecutionResourcePolicy { } let (spill_enabled, spill_directory, spill_max_bytes) = if self.spill.enabled { - let Some(dir) = self.spill.directory.as_ref() else { - return Err(validation( - "spill.directory is required when spill is enabled", - )); - }; - let dir = validate_spill_directory(dir)?; + // No directory selects the project scratch directory (#1595), + // resolved when an instance opens. + let dir = self + .spill + .directory + .as_ref() + .map(|dir| validate_spill_directory(dir)) + .transpose()?; if let Some(max) = self.spill.max_bytes && max == 0 { return Err(validation("spill.max_bytes must be greater than zero")); } - (true, Some(dir), self.spill.max_bytes) + (true, dir, self.spill.max_bytes) } else { if self.spill.directory.is_some() || self.spill.max_bytes.is_some() { return Err(validation( @@ -538,7 +559,10 @@ mod tests { assert_eq!(normalized.io_concurrency, expected); assert_eq!(normalized.compute_threads, expected); assert_eq!(normalized.mode, ResourcePolicyMode::Automatic); - assert!(!normalized.spill_enabled); + // #1595: spill into the project scratch directory by default. + assert!(normalized.spill_enabled); + assert_eq!(normalized.spill_directory, None); + assert_eq!(normalized.spill_max_bytes, None); assert_eq!(normalized.batch_size, DEFAULT_BATCH_SIZE); assert_eq!(normalized.memory_budget_bytes, DEFAULT_MEMORY_BUDGET_BYTES); assert_eq!( @@ -566,9 +590,12 @@ mod tests { assert!(matches!(err, GfError::Validation(_))); } + /// #1595: spill enabled with no directory selects the project scratch + /// directory, with the caller's cap; a zero cap and a cap or directory + /// without spill are still refused. #[test] - fn spill_without_directory_fails() { - let err = ExecutionResourcePolicy { + fn spill_without_directory_selects_project_scratch() { + let normalized = ExecutionResourcePolicy { spill: SpillPolicy { enabled: true, directory: None, @@ -577,8 +604,30 @@ mod tests { ..ExecutionResourcePolicy::default() } .normalize() - .expect_err("spill needs directory"); - assert!(matches!(err, GfError::Validation(_))); + .unwrap(); + assert!(normalized.spill_enabled); + assert_eq!(normalized.spill_directory, None); + assert_eq!(normalized.spill_max_bytes, Some(1024)); + for spill in [ + SpillPolicy { + enabled: true, + directory: None, + max_bytes: Some(0), + }, + SpillPolicy { + enabled: false, + directory: None, + max_bytes: Some(1024), + }, + ] { + let err = ExecutionResourcePolicy { + spill, + ..ExecutionResourcePolicy::default() + } + .normalize() + .expect_err("invalid spill"); + assert!(matches!(err, GfError::Validation(_))); + } } #[test] diff --git a/crates/graphforge-api/src/runtime_ownership.rs b/crates/graphforge-api/src/runtime_ownership.rs index 640e3001e..6be79ca98 100644 --- a/crates/graphforge-api/src/runtime_ownership.rs +++ b/crates/graphforge-api/src/runtime_ownership.rs @@ -32,16 +32,81 @@ impl GraphForge { } } - pub(super) fn session_resource_config(&self) -> graphforge_exec::SessionResourceConfig { + /// Session resources for one query. Spill follows the policy (#1595): + /// disabled means no disk at all; a configured directory is used as given; + /// otherwise a durable project spills into its scratch directory, capped + /// per query, and any other instance does not spill. + /// + /// # Errors + /// Returns the error from acquiring the project scratch directory. + pub(super) fn session_resource_config( + &self, + ) -> Result { + let policy = &self.resource_policy; + let (spill_enabled, spill_directory, spill_max_bytes) = if !policy.spill_enabled { + (false, None, None) + } else if let Some(directory) = &policy.spill_directory { + (true, Some(directory.clone()), policy.spill_max_bytes) + } else if let Some(scratch) = self.query_spill_directory()? { + ( + true, + Some(scratch), + Some( + policy + .spill_max_bytes + .unwrap_or(graphforge_storage::query_spill::DEFAULT_QUERY_SPILL_MAX_BYTES), + ), + ) + } else { + (false, None, None) + }; + Ok(graphforge_exec::SessionResourceConfig { + target_partitions: policy.target_partitions, + batch_size: policy.batch_size, + memory_budget_bytes: policy.memory_budget_bytes, + spill_enabled, + spill_directory, + spill_max_bytes, + io_concurrency: policy.io_concurrency, + }) + } + + /// Session resources for a plan that is rendered, never executed + /// (EXPLAIN): the same knobs with no spill, so planning never acquires + /// scratch it cannot use. + pub(super) fn planning_resource_config(&self) -> graphforge_exec::SessionResourceConfig { + let policy = &self.resource_policy; graphforge_exec::SessionResourceConfig { - target_partitions: self.resource_policy.target_partitions, - batch_size: self.resource_policy.batch_size, - memory_budget_bytes: self.resource_policy.memory_budget_bytes, - spill_enabled: self.resource_policy.spill_enabled, - spill_directory: self.resource_policy.spill_directory.clone(), - spill_max_bytes: self.resource_policy.spill_max_bytes, - io_concurrency: self.resource_policy.io_concurrency, + target_partitions: policy.target_partitions, + batch_size: policy.batch_size, + memory_budget_bytes: policy.memory_budget_bytes, + spill_enabled: false, + spill_directory: None, + spill_max_bytes: None, + io_concurrency: policy.io_concurrency, + } + } + + /// The project scratch directory for query spill, acquired on first use. + /// `None` unless this facade is a durable project, including read-only + /// views of one (checkpoints, inspection), which spill like any query. + fn query_spill_directory(&self) -> Result, GfError> { + let Some(project) = &self.path else { + return Ok(None); + }; + if self.lifecycle_mode + != graphforge_storage::filesystem_admission::ProjectLifecycleMode::Durable + { + return Ok(None); + } + let mut slot = self + .query_spill + .lock() + .map_err(|_| GfError::Storage("query spill directory lock poisoned".into()))?; + if slot.is_none() { + *slot = Some(graphforge_storage::query_spill::QuerySpillDirectory::acquire(project)?); } + Ok(slot.as_ref().map(|scratch| scratch.path().to_path_buf())) } pub(super) fn admit_heavy_query(&self) -> Result, GfError> { @@ -179,3 +244,6 @@ pub(super) fn build_runtime( #[cfg(test)] mod tests; + +#[cfg(test)] +mod spill_tests; diff --git a/crates/graphforge-api/src/runtime_ownership/spill_tests.rs b/crates/graphforge-api/src/runtime_ownership/spill_tests.rs new file mode 100644 index 000000000..7c571839e --- /dev/null +++ b/crates/graphforge-api/src/runtime_ownership/spill_tests.rs @@ -0,0 +1,217 @@ +//! #1595: where each instance's queries may spill. +//! +//! What DataFusion does with a resolved configuration (spill into the given +//! directory under its cap, or refuse when disabled) is proven in +//! `graphforge-exec`'s session tests. These tests prove which configuration +//! each instance and policy resolves to. + +use crate::{ + ExecutionResourcePolicy, GraphForge, GraphForgeOptions, ResourcePolicyMode, SpillPolicy, +}; +use graphforge_storage::query_spill::{DEFAULT_QUERY_SPILL_MAX_BYTES, QUERY_SPILL_DIR}; + +fn options(spill: SpillPolicy) -> GraphForgeOptions { + GraphForgeOptions { + resource: ExecutionResourcePolicy { + mode: ResourcePolicyMode::Explicit, + spill, + ..ExecutionResourcePolicy::default() + }, + ..GraphForgeOptions::default() + } +} + +fn scratch_entries(project: &std::path::Path) -> Vec { + let Ok(entries) = std::fs::read_dir(project.join(QUERY_SPILL_DIR)) else { + return Vec::new(); + }; + let mut names = entries + .map(|entry| entry.unwrap().file_name().to_string_lossy().into_owned()) + .collect::>(); + names.sort(); + names +} + +#[test] +fn a_durable_project_spills_into_its_own_capped_scratch_by_default() { + let project = tempfile::tempdir().unwrap(); + let graph = GraphForge::new_with_options( + Some(project.path().to_str().unwrap()), + options(SpillPolicy::default()), + ) + .unwrap(); + assert!( + scratch_entries(project.path()).is_empty(), + "scratch is acquired by the first query, not at open" + ); + let resources = graph.session_resource_config().unwrap(); + assert!(resources.spill_enabled); + assert_eq!( + resources.spill_max_bytes, + Some(DEFAULT_QUERY_SPILL_MAX_BYTES) + ); + let directory = resources.spill_directory.unwrap(); + assert!(directory.is_dir()); + assert_eq!( + directory.canonicalize().unwrap().parent().unwrap(), + project.path().join(QUERY_SPILL_DIR).canonicalize().unwrap() + ); + // Later queries reuse the same directory. + assert_eq!( + graph.session_resource_config().unwrap().spill_directory, + Some(directory.clone()) + ); + // A query really runs with it. + assert_eq!( + graph + .execute("RETURN 1 AS one") + .unwrap() + .stats + .rows_produced, + 1 + ); + let entries = scratch_entries(project.path()); + assert_eq!(entries.len(), 2, "one directory and its lock: {entries:?}"); + drop(graph); + assert!( + scratch_entries(project.path()).is_empty(), + "the instance's scratch outlived it" + ); +} + +#[test] +fn the_callers_cap_applies_to_project_scratch() { + let project = tempfile::tempdir().unwrap(); + let graph = GraphForge::new_with_options( + Some(project.path().to_str().unwrap()), + options(SpillPolicy { + enabled: true, + directory: None, + max_bytes: Some(4096), + }), + ) + .unwrap(); + let resources = graph.session_resource_config().unwrap(); + assert!(resources.spill_enabled); + assert_eq!(resources.spill_max_bytes, Some(4096)); +} + +#[test] +fn disabled_spill_resolves_to_no_directory_and_creates_no_scratch() { + let project = tempfile::tempdir().unwrap(); + let graph = GraphForge::new_with_options( + Some(project.path().to_str().unwrap()), + options(SpillPolicy { + enabled: false, + directory: None, + max_bytes: None, + }), + ) + .unwrap(); + let resources = graph.session_resource_config().unwrap(); + assert!(!resources.spill_enabled); + assert_eq!(resources.spill_directory, None); + graph.execute("RETURN 1 AS one").unwrap(); + assert!(!project.path().join(QUERY_SPILL_DIR).exists()); +} + +#[test] +fn an_in_memory_instance_does_not_spill_by_default() { + let graph = GraphForge::new_with_options(None, options(SpillPolicy::default())).unwrap(); + let resources = graph.session_resource_config().unwrap(); + assert!(!resources.spill_enabled); + assert_eq!(resources.spill_directory, None); +} + +#[test] +fn a_configured_directory_is_used_instead_of_project_scratch() { + let project = tempfile::tempdir().unwrap(); + let spill = tempfile::tempdir().unwrap(); + let graph = GraphForge::new_with_options( + Some(project.path().to_str().unwrap()), + options(SpillPolicy { + enabled: true, + directory: Some(spill.path().to_path_buf()), + max_bytes: None, + }), + ) + .unwrap(); + let resources = graph.session_resource_config().unwrap(); + assert!(resources.spill_enabled); + assert_eq!(resources.spill_max_bytes, None); + assert_eq!( + resources.spill_directory.unwrap().canonicalize().unwrap(), + spill.path().canonicalize().unwrap() + ); + graph.execute("RETURN 1 AS one").unwrap(); + assert!(!project.path().join(QUERY_SPILL_DIR).exists()); +} + +/// Two instances of one project each own a scratch directory; dropping one +/// leaves the other's in place. +#[test] +fn two_instances_of_one_project_do_not_share_scratch() { + let project = tempfile::tempdir().unwrap(); + let path = project.path().to_str().unwrap(); + let first = GraphForge::new_with_options(Some(path), options(SpillPolicy::default())).unwrap(); + let first_dir = first + .session_resource_config() + .unwrap() + .spill_directory + .unwrap(); + let second = GraphForge::new_with_options(Some(path), options(SpillPolicy::default())).unwrap(); + let second_dir = second + .session_resource_config() + .unwrap() + .spill_directory + .unwrap(); + assert_ne!(first_dir, second_dir); + drop(second); + assert!( + first_dir.is_dir(), + "a live instance's scratch was reclaimed" + ); + assert!(!second_dir.exists()); +} + +/// A read-only view of a durable project (here a checkpoint view) spills like +/// any other query of that project: it acquires its own scratch. +#[test] +fn a_checkpoint_view_of_a_durable_project_gets_scratch() { + let project = tempfile::tempdir().unwrap(); + let graph = GraphForge::new_with_options( + Some(project.path().to_str().unwrap()), + options(SpillPolicy::default()), + ) + .unwrap(); + graph.execute("CREATE (:Row {v: 1})").unwrap(); + graph + .checkpoint(crate::CheckpointRequest { + name: "Spill".into(), + description: None, + idempotency_key: crate::OperationId(uuid::Uuid::from_u128(1595)), + actor_uuid: None, + }) + .unwrap(); + let before = scratch_entries(project.path()); + let view = graph.open_checkpoint("Spill").unwrap(); + view.execute("MATCH (n:Row) RETURN n.v AS v ORDER BY v") + .unwrap(); + let during = scratch_entries(project.path()); + assert_eq!(during.len(), before.len() + 2, "{before:?} -> {during:?}"); + drop(view); + assert_eq!(scratch_entries(project.path()), before); +} + +/// EXPLAIN renders a plan and never executes it, so it acquires no scratch. +#[test] +fn explain_acquires_no_scratch() { + let project = tempfile::tempdir().unwrap(); + let graph = GraphForge::new_with_options( + Some(project.path().to_str().unwrap()), + options(SpillPolicy::default()), + ) + .unwrap(); + graph.explain("MATCH (n) RETURN n ORDER BY n").unwrap(); + assert!(scratch_entries(project.path()).is_empty()); +} diff --git a/crates/graphforge-api/src/search_index.rs b/crates/graphforge-api/src/search_index.rs index f521d7229..c4c68e60a 100644 --- a/crates/graphforge-api/src/search_index.rs +++ b/crates/graphforge-api/src/search_index.rs @@ -355,13 +355,15 @@ impl GraphForge { staged_project_dir: &std::path::Path, ) -> graphforge_storage::adjacency::AdjacencyBuildOptions { let policy = &self.resource_policy; - let spill_dir = if policy.spill_enabled { - policy.spill_directory.clone() - } else { - Some( + // Only a caller-configured directory replaces the stage's own spill + // root; the default project scratch directory is for query spill + // (#1595), and its per-query cap is not this build's bound. + let spill_dir = match (&policy.spill_enabled, &policy.spill_directory) { + (true, Some(directory)) => Some(directory.clone()), + _ => Some( graphforge_storage::adjacency::adjacency_dir(staged_project_dir) .join(graphforge_storage::adjacency::ADJACENCY_SPILL_DIR_NAME), - ) + ), }; let mut options = graphforge_storage::adjacency::AdjacencyBuildOptions::default(); options.chunk_rows = { @@ -378,7 +380,7 @@ impl GraphForge { }; options.batch_size = policy.batch_size.max(1); options.spill_dir = spill_dir; - options.spill_max_bytes = policy.spill_max_bytes; + options.spill_max_bytes = policy.spill_directory.as_ref().and(policy.spill_max_bytes); options.memory_budget_bytes = Some(policy.memory_budget_bytes); options.shard_max_edges = { let budget_entries = policy diff --git a/crates/graphforge-exec/src/session.rs b/crates/graphforge-exec/src/session.rs index f384cfeb4..938646fb7 100644 --- a/crates/graphforge-exec/src/session.rs +++ b/crates/graphforge-exec/src/session.rs @@ -344,6 +344,38 @@ impl Default for SessionResourceConfig { } } +/// The DataFusion runtime for one session: its memory pool, and its disk +/// manager as the resources say (#1595). Spill goes only to the configured +/// directory under its cap; disabled means no disk manager at all, not +/// DataFusion's default of the OS temporary directory, which is unbounded and +/// RAM-backed on some hosts. +fn session_runtime( + resources: &SessionResourceConfig, +) -> Result, GfError> { + let memory_budget = usize::try_from(resources.memory_budget_bytes).unwrap_or(usize::MAX); + let mut runtime_builder = datafusion::execution::runtime_env::RuntimeEnvBuilder::new() + .with_memory_limit(memory_budget, 1.0); + if let (true, Some(dir)) = (resources.spill_enabled, &resources.spill_directory) { + std::fs::create_dir_all(dir) + .map_err(|error| GfError::Storage(format!("query spill directory: {error}")))?; + runtime_builder = runtime_builder.with_temp_file_path(dir.clone()); + if let Some(max) = resources.spill_max_bytes { + runtime_builder = runtime_builder.with_max_temp_directory_size(max); + } + } else { + runtime_builder = runtime_builder.with_disk_manager_builder( + datafusion::execution::disk_manager::DiskManagerBuilder::default() + .with_mode(datafusion::execution::disk_manager::DiskManagerMode::Disabled), + ); + } + // A spilling runtime creates its temporary directory here, so a full or + // read-only volume is an error for this query, not a panic. + runtime_builder + .build() + .map(Arc::new) + .map_err(|error| GfError::Storage(format!("query runtime: {error}"))) +} + struct QueryEvidenceStream { inner: Option, physical: Arc, @@ -479,7 +511,7 @@ impl ExecutionSession { /// # Errors /// Returns [`GfError`] if session construction fails. pub fn new(catalog: GraphCatalog, ontology: Option) -> Result { - Ok(Self::build( + Self::build( catalog, ontology, PathBuf::new(), @@ -487,7 +519,7 @@ impl ExecutionSession { None, OrdinalIdentityConfig::default(), &SessionResourceConfig::default(), - )) + ) } /// Create a session that can execute writes against `dir`. @@ -500,7 +532,7 @@ impl ExecutionSession { dir: PathBuf, mode: OntologyMode, ) -> Result { - Ok(Self::build( + Self::build( catalog, ontology, dir, @@ -508,7 +540,7 @@ impl ExecutionSession { None, OrdinalIdentityConfig::default(), &SessionResourceConfig::default(), - )) + ) } /// Like [`new_with_target`](Self::new_with_target) but reusing a @@ -565,6 +597,11 @@ impl ExecutionSession { ordinal_identities: Option>, resources: &SessionResourceConfig, ) -> Result { + if resources.spill_enabled && resources.spill_directory.is_none() { + return Err(GfError::Validation( + "query spill is enabled without a spill directory".into(), + )); + } let identity = match ordinal_identities { Some(resolver) => { let pin = resolver.pin()?; @@ -579,7 +616,7 @@ impl ExecutionSession { } None => OrdinalIdentityConfig::default(), }; - Ok(Self::build( + Self::build( catalog, ontology, dir, @@ -587,7 +624,7 @@ impl ExecutionSession { Some(provider), identity, resources, - )) + ) } fn build( @@ -598,7 +635,7 @@ impl ExecutionSession { shared_provider: Option>, identity: OrdinalIdentityConfig, resources: &SessionResourceConfig, - ) -> Self { + ) -> Result { // The session-scoped adjacency provider (#761), threaded to the // extension planner via SessionConfig extension. Read-only sessions // (empty dir) need no special case: with no `indexes/adjacency/` @@ -665,22 +702,7 @@ impl ExecutionSession { .use_row_number_estimates_to_optimize_partitioning = true; let memory_budget = usize::try_from(resources.memory_budget_bytes).unwrap_or(usize::MAX); - let mut runtime_builder = datafusion::execution::runtime_env::RuntimeEnvBuilder::new() - .with_memory_limit(memory_budget, 1.0); - if resources.spill_enabled - && let Some(dir) = &resources.spill_directory - { - let _ = std::fs::create_dir_all(dir); - runtime_builder = runtime_builder.with_temp_file_path(dir.clone()); - if let Some(max) = resources.spill_max_bytes { - runtime_builder = runtime_builder.with_max_temp_directory_size(max); - } - } - let runtime_env = Arc::new( - runtime_builder - .build() - .expect("DataFusion RuntimeEnv construction"), - ); + let runtime_env = session_runtime(resources)?; let state = SessionStateBuilder::new() .with_default_features() @@ -703,7 +725,7 @@ impl ExecutionSession { .semantic_composition_fingerprint() .map(str::to_owned); ctx.register_catalog("graph", catalog.clone()); - Self { + Ok(Self { mutation_health, ctx, catalog, @@ -715,7 +737,7 @@ impl ExecutionSession { owned_adjacency, #[cfg(feature = "differential-testing")] relational_fixed_hop_reference: false, - } + }) } pub(super) fn refresh_adjacency_after_mutation(&self) -> Result<(), GfError> { diff --git a/crates/graphforge-exec/src/session/tests.rs b/crates/graphforge-exec/src/session/tests.rs index ec8f74422..32009ed8f 100644 --- a/crates/graphforge-exec/src/session/tests.rs +++ b/crates/graphforge-exec/src/session/tests.rs @@ -696,7 +696,8 @@ async fn order_by_spills_when_input_exceeds_query_memory_budget() { None, OrdinalIdentityConfig::default(), &resources, - ); + ) + .unwrap(); let batches = uuid_pair_batches(ROWS, resources.batch_size); let table = MemTable::try_new(batches[0].schema(), vec![batches]).unwrap(); session.ctx.register_table("t", Arc::new(table)).unwrap(); @@ -742,6 +743,158 @@ async fn order_by_spills_when_input_exceeds_query_memory_budget() { } } +/// #1595: with spill disabled, a sort over the budget fails with a resource +/// error. It must not fall back to DataFusion's default disk manager, which +/// spills into the OS temporary directory without any limit. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn disabled_spill_refuses_instead_of_spilling_to_the_os_temp_directory() { + use datafusion::datasource::MemTable; + + const ROWS: usize = 2_500_000; + let dir = TempDir::new().unwrap(); + let catalog = GraphCatalog::open(dir.path(), None, &RuntimeCatalog::new()).unwrap(); + let resources = SessionResourceConfig { + target_partitions: 1, + memory_budget_bytes: 128 * 1024 * 1024, + spill_enabled: false, + ..SessionResourceConfig::default() + }; + let session = ExecutionSession::build( + catalog, + None, + PathBuf::new(), + OntologyMode::Exploratory, + None, + OrdinalIdentityConfig::default(), + &resources, + ) + .unwrap(); + assert!( + !session.ctx.runtime_env().disk_manager.tmp_files_enabled(), + "a disabled spill policy must disable the disk manager" + ); + let batches = uuid_pair_batches(ROWS, resources.batch_size); + let table = MemTable::try_new(batches[0].schema(), vec![batches]).unwrap(); + session.ctx.register_table("t", Arc::new(table)).unwrap(); + let plan = session + .ctx + .sql("SELECT s, d FROM t ORDER BY s, d") + .await + .unwrap() + .create_physical_plan() + .await + .unwrap(); + let error = datafusion::physical_plan::collect(Arc::clone(&plan), session.ctx.task_ctx()) + .await + .expect_err("a sort over the budget cannot complete without spill") + .to_string(); + assert!(error.contains("Resources exhausted"), "{error}"); + assert_eq!(sort_spill_count(&plan), 0); +} + +/// #1595: the spill cap bounds what one query may write to its spill +/// directory. A sort that needs more fails instead of growing past it. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn the_spill_cap_bounds_one_query() { + use datafusion::datasource::MemTable; + + const ROWS: usize = 2_500_000; + let spill = TempDir::new().unwrap(); + let dir = TempDir::new().unwrap(); + let catalog = GraphCatalog::open(dir.path(), None, &RuntimeCatalog::new()).unwrap(); + let resources = SessionResourceConfig { + target_partitions: 1, + memory_budget_bytes: 128 * 1024 * 1024, + spill_enabled: true, + spill_directory: Some(spill.path().to_path_buf()), + spill_max_bytes: Some(4096), + ..SessionResourceConfig::default() + }; + let session = ExecutionSession::build( + catalog, + None, + PathBuf::new(), + OntologyMode::Exploratory, + None, + OrdinalIdentityConfig::default(), + &resources, + ) + .unwrap(); + let batches = uuid_pair_batches(ROWS, resources.batch_size); + let table = MemTable::try_new(batches[0].schema(), vec![batches]).unwrap(); + session.ctx.register_table("t", Arc::new(table)).unwrap(); + let plan = session + .ctx + .sql("SELECT s, d FROM t ORDER BY s, d") + .await + .unwrap() + .create_physical_plan() + .await + .unwrap(); + let error = datafusion::physical_plan::collect(Arc::clone(&plan), session.ctx.task_ctx()) + .await + .expect_err("the sort needs more spill than the cap allows") + .to_string(); + assert!(error.contains("exceeded"), "{error}"); +} + +/// #1595: a spill directory that cannot be used is an error for the session +/// being built, not a panic. +#[test] +fn an_unusable_spill_directory_is_an_error() { + let dir = TempDir::new().unwrap(); + let blocked = dir.path().join("not-a-directory"); + std::fs::write(&blocked, b"a file").unwrap(); + let catalog = GraphCatalog::open(dir.path(), None, &RuntimeCatalog::new()).unwrap(); + let resources = SessionResourceConfig { + spill_enabled: true, + spill_directory: Some(blocked.join("spill")), + ..SessionResourceConfig::default() + }; + let error = ExecutionSession::build( + catalog, + None, + PathBuf::new(), + OntologyMode::Exploratory, + None, + OrdinalIdentityConfig::default(), + &resources, + ) + .err() + .expect("an unusable spill directory must be refused") + .to_string(); + assert!(error.contains("query spill directory"), "{error}"); +} + +/// #1595: spill enabled without a directory has nowhere to go and is refused +/// when the session is built. +#[test] +fn enabled_spill_without_a_directory_is_refused() { + let dir = TempDir::new().unwrap(); + let catalog = GraphCatalog::open(dir.path(), None, &RuntimeCatalog::new()).unwrap(); + let resources = SessionResourceConfig { + spill_enabled: true, + spill_directory: None, + ..SessionResourceConfig::default() + }; + let error = ExecutionSession::new_with_target_provider_resources_and_identity( + catalog, + None, + dir.path().to_path_buf(), + OntologyMode::Exploratory, + Arc::new(PersistentAdjacencyProvider::new( + dir.path().to_path_buf(), + OntologyMode::Exploratory, + )), + None, + &resources, + ) + .err() + .expect("spill without a directory must be refused") + .to_string(); + assert!(error.contains("without a spill directory"), "{error}"); +} + /// Every full sort a Cypher query plans gets coalesced input runs; top-k sorts /// (`LIMIT`) do not, because they never spill and the ordered fast paths match /// their shape. diff --git a/crates/graphforge-storage/src/lib.rs b/crates/graphforge-storage/src/lib.rs index abe099b16..386e77a82 100644 --- a/crates/graphforge-storage/src/lib.rs +++ b/crates/graphforge-storage/src/lib.rs @@ -20,6 +20,7 @@ mod route_component; pub use durable_rewrite::AuxiliaryReceipt; #[doc(hidden)] pub mod filesystem_admission; +pub mod query_spill; pub mod adjacency; pub mod adjacency_delta; diff --git a/crates/graphforge-storage/src/query_spill.rs b/crates/graphforge-storage/src/query_spill.rs new file mode 100644 index 000000000..ae15d8c04 --- /dev/null +++ b/crates/graphforge-storage/src/query_spill.rs @@ -0,0 +1,233 @@ +//! Project-owned scratch space for query spill (#1595). +//! +//! A durable project's queries spill into `/.graphforge-query-spill/` +//! instead of the operating system's temporary directory, which is RAM-backed +//! on some hosts and outside every budget. Several processes may open one +//! project at once, so each open instance owns one subdirectory, named by a +//! random token, and holds an exclusive lock on a sibling `.lock` file +//! for as long as it lives: +//! +//! ```text +//! .graphforge-query-spill/ +//! 3f2c…e9.lock held by the live instance that owns 3f2c…e9/ +//! 3f2c…e9/ DataFusion's spill files for that instance's queries +//! ``` +//! +//! An owner creates its lock file under a temporary name (`.lock.new`), +//! locks it, and only then renames it to `.lock`, so every visible +//! `.lock` is already held by its owner if the owner is alive. It +//! creates its subdirectory after the rename and, when dropped, removes the +//! subdirectory before the lock. +//! +//! A crashed process cannot clean up, so acquiring a new directory first +//! reclaims every entry whose lock is free: its owner is gone, because the +//! operating system releases a process's locks when it exits. A subdirectory +//! is also stale when its lock file is absent at the moment it is examined, +//! because an owner's lock exists for as long as its subdirectory does. +//! Reclaiming other owners' leftovers is housekeeping: an entry that cannot be +//! removed is skipped, and never stops this instance acquiring its own +//! scratch. Entries that do not match this grammar are left alone. +//! +//! Scratch is transient: it is never recovery authority and holds nothing a +//! query needs after it finishes. + +use crate::file_lock::try_lock_exclusive; +use graphforge_core::GfError; +use std::fs::{File, OpenOptions}; +use std::path::{Path, PathBuf}; + +/// The scratch root under a project directory. +pub const QUERY_SPILL_DIR: &str = ".graphforge-query-spill"; + +/// Default cap on one query's spill bytes in the project scratch directory. +pub const DEFAULT_QUERY_SPILL_MAX_BYTES: u64 = 8 << 30; + +/// One instance's scratch subdirectory, held for the instance's lifetime. +#[derive(Debug)] +pub struct QuerySpillDirectory { + directory: PathBuf, + lock_path: PathBuf, + // Held open, and so locked, until drop. + _lock: File, +} + +fn storage(error: impl std::fmt::Display) -> GfError { + GfError::Storage(format!("query spill directory: {error}")) +} + +/// Whether `name` is an owner token: 32 lowercase hex digits. +fn is_token(name: &str) -> bool { + name.len() == 32 + && name + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) +} + +impl QuerySpillDirectory { + /// Reclaim abandoned scratch under `project`, then create and lock this + /// instance's own subdirectory. + /// + /// # Errors + /// Refuses a scratch root that is not a real directory, and returns any + /// filesystem error from reclaiming or creating scratch. + pub fn acquire(project: &Path) -> Result { + let root = project.join(QUERY_SPILL_DIR); + match std::fs::symlink_metadata(&root) { + Ok(metadata) if metadata.is_dir() => {} + Ok(_) => return Err(storage("the scratch root is not a directory")), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + match std::fs::create_dir(&root) { + Ok(()) => {} + Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => { + if !std::fs::symlink_metadata(&root).map_err(storage)?.is_dir() { + return Err(storage("the scratch root is not a directory")); + } + } + Err(error) => return Err(storage(error)), + } + } + Err(error) => return Err(storage(error)), + } + reclaim_abandoned(&root); + // A reclaimer can remove a temporary lock file between its creation + // and its lock; the rename then fails and a fresh token is tried. + let mut attempt = 0; + let (token, lock_path, lock) = loop { + attempt += 1; + match claim_lock(&root)? { + Some(claimed) => break claimed, + None if attempt < 8 => {} + None => return Err(storage("could not claim a scratch lock")), + } + }; + let directory = root.join(&token); + if let Err(error) = std::fs::create_dir(&directory) { + let _ = std::fs::remove_file(&lock_path); + return Err(storage(error)); + } + Ok(Self { + directory, + lock_path, + _lock: lock, + }) + } + + /// The directory this instance's queries spill into. + #[must_use] + pub fn path(&self) -> &Path { + &self.directory + } +} + +impl Drop for QuerySpillDirectory { + fn drop(&mut self) { + // The subdirectory first: while the lock file exists and is held, no + // other process reclaims it. + let _ = std::fs::remove_dir_all(&self.directory); + let _ = std::fs::remove_file(&self.lock_path); + } +} + +/// Create, lock and publish one `.lock`. `None` when another process +/// reclaimed the temporary file before it was locked or published. +fn claim_lock(root: &Path) -> Result, GfError> { + let token = uuid::Uuid::new_v4().simple().to_string(); + let temporary = root.join(format!("{token}.lock.new")); + let lock_path = root.join(format!("{token}.lock")); + let lock = OpenOptions::new() + .write(true) + .create_new(true) + .open(&temporary) + .map_err(storage)?; + #[cfg(test)] + tests::between_create_and_lock(root, &temporary); + if !try_lock_exclusive(&lock).map_err(storage)? { + return Ok(None); + } + // The temporary must still be the file this handle locked. + let identity = graphforge_filesystem::file_identity(&lock).map_err(storage)?; + match graphforge_filesystem::path_identity(&temporary) { + Ok(named) if named == identity => {} + Ok(_) => return Ok(None), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None), + Err(error) => return Err(storage(error)), + } + match std::fs::rename(&temporary, &lock_path) { + Ok(()) => Ok(Some((token, lock_path, lock))), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None), + Err(error) => Err(storage(error)), + } +} + +/// Remove every scratch entry whose owner is gone. Best effort: an entry that +/// cannot be examined or removed is skipped. +fn reclaim_abandoned(root: &Path) { + let _ = reclaim_abandoned_entries(root); +} + +fn reclaim_abandoned_entries(root: &Path) -> Result<(), GfError> { + let mut locks = Vec::new(); + let mut directories = Vec::new(); + for entry in std::fs::read_dir(root).map_err(storage)? { + let entry = entry.map_err(storage)?; + // Not followed: a link is never an owner's entry. + let kind = entry.file_type().map_err(storage)?; + let Some(name) = entry.file_name().to_str().map(str::to_owned) else { + continue; + }; + if kind.is_file() { + let lock = name + .strip_suffix(".lock") + .or_else(|| name.strip_suffix(".lock.new")) + .filter(|token| is_token(token)); + if lock.is_some() { + locks.push(name); + } + } else if kind.is_dir() && is_token(&name) { + directories.push(name); + } + } + for name in &locks { + let lock_path = root.join(name); + let Ok(lock) = File::open(&lock_path) else { + continue; // Reclaimed, released or unreadable meanwhile. + }; + if !matches!(try_lock_exclusive(&lock), Ok(true)) { + continue; // Its owner is alive. + } + if let Some(token) = name.strip_suffix(".lock") { + let _ = remove_directory(&root.join(token)); + } + let _ = remove_file(&lock_path); + } + for token in &directories { + // Examined now, not from the listing: an owner's lock exists for as + // long as its subdirectory does. + match std::fs::symlink_metadata(root.join(format!("{token}.lock"))) { + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + let _ = remove_directory(&root.join(token)); + } + _ => {} + } + } + Ok(()) +} + +fn remove_directory(path: &Path) -> Result<(), GfError> { + match std::fs::remove_dir_all(path) { + Ok(()) => Ok(()), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(error) => Err(storage(error)), + } +} + +fn remove_file(path: &Path) -> Result<(), GfError> { + match std::fs::remove_file(path) { + Ok(()) => Ok(()), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(error) => Err(storage(error)), + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/graphforge-storage/src/query_spill/tests.rs b/crates/graphforge-storage/src/query_spill/tests.rs new file mode 100644 index 000000000..163569369 --- /dev/null +++ b/crates/graphforge-storage/src/query_spill/tests.rs @@ -0,0 +1,252 @@ +use super::*; +use std::cell::{Cell, RefCell}; + +thread_local! { + /// What another process does, in a test, between an owner creating its + /// temporary lock file and locking it. + static INTERLEAVE: Cell> = const { Cell::new(None) }; + /// Locks the simulated other process holds. + static HELD: RefCell> = const { RefCell::new(Vec::new()) }; +} + +#[derive(Clone, Copy)] +enum Interleave { + /// Reclaim the unheld temporary, as an opener's reclamation would. + Reclaim, + /// Lock the temporary first and keep holding it. + Hold, +} + +pub(super) fn between_create_and_lock(root: &Path, temporary: &Path) { + match INTERLEAVE.with(Cell::take) { + Some(Interleave::Reclaim) => reclaim_abandoned(root), + Some(Interleave::Hold) => { + let other = File::open(temporary).unwrap(); + assert!(try_lock_exclusive(&other).unwrap()); + HELD.with(|held| held.borrow_mut().push(other)); + } + None => {} + } +} + +fn entries(root: &Path) -> Vec { + let mut names = std::fs::read_dir(root) + .unwrap() + .map(|entry| entry.unwrap().file_name().to_string_lossy().into_owned()) + .collect::>(); + names.sort(); + names +} + +#[test] +fn an_instance_owns_a_locked_subdirectory_until_dropped() { + let project = tempfile::TempDir::new().unwrap(); + let root = project.path().join(QUERY_SPILL_DIR); + let scratch = QuerySpillDirectory::acquire(project.path()).unwrap(); + assert!(scratch.path().is_dir()); + assert_eq!(scratch.path().parent(), Some(root.as_path())); + std::fs::write(scratch.path().join("spill-0"), b"rows").unwrap(); + assert_eq!(entries(&root).len(), 2, "{:?}", entries(&root)); + drop(scratch); + assert!(entries(&root).is_empty(), "{:?}", entries(&root)); +} + +#[test] +fn a_live_instance_is_never_reclaimed() { + let project = tempfile::TempDir::new().unwrap(); + let first = QuerySpillDirectory::acquire(project.path()).unwrap(); + std::fs::write(first.path().join("spill-0"), b"rows").unwrap(); + // A second instance, even in the same process, holds a separate lock and + // must leave the first one's files in place. + let second = QuerySpillDirectory::acquire(project.path()).unwrap(); + assert_ne!(first.path(), second.path()); + assert!(first.path().join("spill-0").is_file()); + drop(second); + assert!(first.path().join("spill-0").is_file()); +} + +#[cfg(unix)] +#[test] +fn abandoned_scratch_is_reclaimed_and_other_entries_are_left_alone() { + let project = tempfile::TempDir::new().unwrap(); + let root = project.path().join(QUERY_SPILL_DIR); + std::fs::create_dir(&root).unwrap(); + // An owner that crashed: its lock file exists but nobody holds it. + let crashed = "0123456789abcdef0123456789abcdef"; + std::fs::write(root.join(format!("{crashed}.lock")), b"").unwrap(); + std::fs::create_dir(root.join(crashed)).unwrap(); + std::fs::write(root.join(crashed).join("spill-0"), b"rows").unwrap(); + // A subdirectory whose lock file is already gone. + let orphan = "fedcba9876543210fedcba9876543210"; + std::fs::create_dir(root.join(orphan)).unwrap(); + // A lock file that crashed before it was published. + let unpublished = "abcdefabcdefabcdefabcdefabcdefab"; + std::fs::write(root.join(format!("{unpublished}.lock.new")), b"").unwrap(); + // Names outside the grammar, and a link shaped like a token. + std::fs::write(root.join("README"), b"not ours").unwrap(); + std::fs::create_dir(root.join("keep")).unwrap(); + let outside = tempfile::TempDir::new().unwrap(); + std::fs::write(outside.path().join("precious"), b"data").unwrap(); + let link = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; + std::os::unix::fs::symlink(outside.path(), root.join(link)).unwrap(); + + let scratch = QuerySpillDirectory::acquire(project.path()).unwrap(); + let own = scratch + .path() + .file_name() + .unwrap() + .to_string_lossy() + .into_owned(); + let mut expected = vec![ + "README".to_owned(), + link.to_owned(), + own.clone(), + format!("{own}.lock"), + "keep".to_owned(), + ]; + expected.sort(); + assert_eq!(entries(&root), expected); + assert!( + outside.path().join("precious").is_file(), + "a link was followed" + ); +} + +#[cfg(unix)] +#[test] +fn a_scratch_root_that_is_not_a_directory_is_refused() { + let project = tempfile::TempDir::new().unwrap(); + std::fs::write(project.path().join(QUERY_SPILL_DIR), b"a file").unwrap(); + let error = QuerySpillDirectory::acquire(project.path()).unwrap_err(); + assert!(error.to_string().contains("not a directory"), "{error}"); + + let linked = tempfile::TempDir::new().unwrap(); + let target = tempfile::TempDir::new().unwrap(); + std::os::unix::fs::symlink(target.path(), linked.path().join(QUERY_SPILL_DIR)).unwrap(); + let error = QuerySpillDirectory::acquire(linked.path()).unwrap_err(); + assert!(error.to_string().contains("not a directory"), "{error}"); +} + +/// Reclaiming other owners' leftovers is housekeeping: an abandoned entry that +/// cannot be removed is skipped, and this instance still gets its scratch. +#[cfg(unix)] +#[test] +fn an_unremovable_abandoned_entry_does_not_block_acquisition() { + use std::os::unix::fs::PermissionsExt; + let project = tempfile::TempDir::new().unwrap(); + let root = project.path().join(QUERY_SPILL_DIR); + std::fs::create_dir(&root).unwrap(); + let stuck = "0123456789abcdef0123456789abcdef"; + std::fs::write(root.join(format!("{stuck}.lock")), b"").unwrap(); + std::fs::create_dir(root.join(stuck)).unwrap(); + std::fs::create_dir(root.join(stuck).join("inner")).unwrap(); + std::fs::write(root.join(stuck).join("inner").join("spill-0"), b"rows").unwrap(); + // Its contents cannot be unlinked. + std::fs::set_permissions( + root.join(stuck).join("inner"), + std::fs::Permissions::from_mode(0o500), + ) + .unwrap(); + let scratch = QuerySpillDirectory::acquire(project.path()); + std::fs::set_permissions( + root.join(stuck).join("inner"), + std::fs::Permissions::from_mode(0o700), + ) + .unwrap(); + let scratch = scratch.expect("a stuck leftover must not block acquisition"); + assert!(scratch.path().is_dir()); +} + +/// A lock file is only ever visible under its published name once its owner +/// holds it, so a reclaimer never sees an unheld live lock. +#[test] +fn a_published_lock_is_already_held() { + let project = tempfile::TempDir::new().unwrap(); + let scratch = QuerySpillDirectory::acquire(project.path()).unwrap(); + let root = project.path().join(QUERY_SPILL_DIR); + let token = scratch + .path() + .file_name() + .unwrap() + .to_string_lossy() + .into_owned(); + let lock = File::open(root.join(format!("{token}.lock"))).unwrap(); + assert!( + !try_lock_exclusive(&lock).unwrap(), + "the owner does not hold its lock" + ); + assert!(!root.join(format!("{token}.lock.new")).exists()); +} + +/// Concurrent acquirers of one project: every acquisition succeeds and no +/// instance's scratch disappears while it is held. A smoke test only: the race +/// windows are too narrow to hit reliably here, so the two interleaving tests +/// above are what prove a raced claim is retried. +#[test] +fn concurrent_acquisitions_never_fail_or_reclaim_a_live_instance() { + let project = tempfile::TempDir::new().unwrap(); + std::thread::scope(|scope| { + for _ in 0..16 { + scope.spawn(|| { + for _ in 0..50 { + let scratch = QuerySpillDirectory::acquire(project.path()) + .expect("a concurrent acquisition failed"); + std::fs::write(scratch.path().join("spill-0"), b"rows") + .expect("a live instance's scratch was reclaimed"); + assert!(scratch.path().join("spill-0").is_file()); + } + }); + } + }); + assert!(entries(&project.path().join(QUERY_SPILL_DIR)).is_empty()); +} + +/// Another process reclaims an owner's temporary lock file between its +/// creation and its lock: the owner must notice and claim a fresh one, and +/// end up with a published lock it holds. +#[test] +fn a_claim_reclaimed_before_its_lock_is_retried() { + let project = tempfile::TempDir::new().unwrap(); + INTERLEAVE.with(|step| step.set(Some(Interleave::Reclaim))); + let scratch = QuerySpillDirectory::acquire(project.path()).expect("the claim is retried"); + let root = project.path().join(QUERY_SPILL_DIR); + let token = scratch + .path() + .file_name() + .unwrap() + .to_string_lossy() + .into_owned(); + let lock = File::open(root.join(format!("{token}.lock"))).unwrap(); + assert!( + !try_lock_exclusive(&lock).unwrap(), + "the published lock is not held" + ); + let entries = entries(&root); + assert_eq!( + entries, + vec![token.clone(), format!("{token}.lock")], + "{entries:?}" + ); +} + +/// Another process locks an owner's temporary lock file first: the owner +/// must claim a fresh one rather than fail the query. +#[test] +fn a_claim_whose_lock_is_taken_is_retried() { + let project = tempfile::TempDir::new().unwrap(); + INTERLEAVE.with(|step| step.set(Some(Interleave::Hold))); + let scratch = QuerySpillDirectory::acquire(project.path()).expect("the claim is retried"); + let root = project.path().join(QUERY_SPILL_DIR); + let token = scratch + .path() + .file_name() + .unwrap() + .to_string_lossy() + .into_owned(); + let lock = File::open(root.join(format!("{token}.lock"))).unwrap(); + assert!( + !try_lock_exclusive(&lock).unwrap(), + "the published lock is not held" + ); + HELD.with(|held| held.borrow_mut().clear()); +} diff --git a/docs/development/execution-resource-policy.md b/docs/development/execution-resource-policy.md index 22a22cfd1..a47cccb32 100644 --- a/docs/development/execution-resource-policy.md +++ b/docs/development/execution-resource-policy.md @@ -14,7 +14,7 @@ Tokio runtime or any DataFusion execution session. | `target_partitions` | `2` | DataFusion `SessionConfig` | | `batch_size` | `8192` | DataFusion `SessionConfig`; analyst Arrow shaping / property enrichment (#341) | | `memory_budget_bytes` | `512 MiB` | DataFusion `RuntimeEnv` memory pool | -| `spill` | disabled | Optional absolute spill directory + byte cap | +| `spill` | project scratch, capped (see below) | DataFusion disk manager for query spill | | `io_concurrency` | `2` | Reserved I/O concurrency budget | | `max_concurrent_heavy_queries` | `64` | Instance-owned admission semaphore | | `compute_threads` | `2` | Instance-owned private CPU pool (#342 cosine KNN; #343 PageRank; #344 Node2Vec walks; #501 betweenness; #503 closeness BFS; #504 clustering coefficient; #506 Degree; #508 harmonic closeness BFS; #510 HITS hub; #513 resource allocation; #514 total-neighbors aggregate; #515 triangles; #518 Components; #534 filtered Jaccard; #535 Jaccard similarity; #542 Dijkstra APSP sources) | @@ -44,13 +44,48 @@ Normalization returns structured [`GfError::Validation`](../../crates/graphforge (`min(max(4, 2×cpus), 512)`) - reserved `io_concurrency` / `compute_threads` exceed `max(tokio_workers, observed_cpus)` -- spill is enabled without an absolute, non-symlink directory -- spill directory/max_bytes are set while spill is disabled +- a configured spill directory is not absolute, is a symlink, or is not a directory +- spill directory/max_bytes are set while spill is disabled, or `max_bytes` is zero - memory / batch / heavy-query bounds are out of range Invalid settings never partially construct a runtime or authoritative spill state. +## Query spill (#1595) + +`SpillPolicy` decides where a query whose sort or aggregation exceeds +`memory_budget_bytes` may spill: + +| Policy | Durable project | In-memory instance | +|---|---|---| +| default: `enabled`, no `directory` | the project scratch directory `/.graphforge-query-spill/`, at most `max_bytes` (default 8 GiB) per query | no spill | +| `enabled` with an absolute `directory` | that directory, at most `max_bytes` (DataFusion's 100 GiB default when unset) per query | same | +| `enabled: false` | no spill | no spill | + +"No spill" is enforced: DataFusion's disk manager is disabled, so a query over +its budget fails with a resource error rather than writing to the operating +system's temporary directory, which is unbounded and RAM-backed on some hosts. +Read-only views of a durable project, such as checkpoint and inspection views, +spill into the project scratch like any other query. Views that bound their own +memory (branch, private checkpoint-materialization and slice views) set +`enabled: false` and so fail closed the same way. EXPLAIN renders a plan +without executing it, so it never acquires scratch. + +Each open instance of a durable project owns one scratch subdirectory, created +on its first query and held by a lock on a sibling `.lock` file. The lock is +published under its final name only once it is held, so a live instance's +scratch is never taken for abandoned. The instance removes both when it is +dropped. Acquiring scratch reclaims any whose lock is free, which means the +owning process has exited. That is best effort: an entry that cannot be removed +is skipped, and never stops the new instance acquiring its own +(`graphforge_storage::query_spill`). Scratch holds only a running query's spill +files and is never recovery authority. A spill directory that cannot be created +or used makes that query fail with a storage error. + +The search-index adjacency build keeps its own spill root inside its +unpublished stage unless the caller configures a `directory`. The query scratch +cap does not apply to it. + ## Application 1. `GraphForgeOptions::validate` → `(options, NormalizedResourcePolicy)` @@ -58,7 +93,8 @@ state. 3. `build_runtime(&policy)` sets Tokio worker threads 4. Every DataFusion `ExecutionSession` receives a [`SessionResourceConfig`](../../crates/graphforge-exec/src/session.rs) with - partitions, batch size, memory, optional spill, and `io_concurrency` + partitions, batch size, memory, the resolved spill directory and cap (or + none), and `io_concurrency` 5. Heavy ops (`run_query`, streams construction, `rank`, `similar`, `analyze_embedding`) take an admission permit 6. Query-facing Parquet scans (`GraphForgeParquetExec`, #339) defer file I/O to