From 20688b89572c8a4e8a4b2e9bdae47c3cb4f07d3c Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Mon, 28 Sep 2026 02:47:15 +0000 Subject: [PATCH 1/4] fix(exec,api): disabled spill means no spill; durable projects spill into capped project scratch (#1595) A disabled spill policy used to leave DataFusion's default disk manager in place, which spills into the OS temporary directory with no cap, so queries spilled when the policy said they must not. Now: - spill disabled disables the disk manager: a query over its budget fails with a resource error; - the default policy (enabled, no directory) spills a durable project's queries into /.graphforge-query-spill/, one locked subdirectory per open instance, capped at 8 GiB per query unless the caller sets max_bytes; an in-memory instance does not spill; - a configured directory works as before; - scratch left by an exited process is reclaimed when an instance next acquires scratch; live instances' scratch is never touched; - search-index adjacency builds keep their own stage spill root unless a directory is configured. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../graphforge-api/src/algorithm_writeback.rs | 2 +- .../graphforge-api/src/embedding_refresh.rs | 1 + crates/graphforge-api/src/explanation.rs | 2 +- crates/graphforge-api/src/lib.rs | 6 + crates/graphforge-api/src/query_execution.rs | 4 +- crates/graphforge-api/src/resource_policy.rs | 85 +++++++-- .../graphforge-api/src/runtime_ownership.rs | 70 ++++++- .../src/runtime_ownership/spill_tests.rs | 175 +++++++++++++++++ crates/graphforge-api/src/search_index.rs | 14 +- crates/graphforge-exec/src/session.rs | 17 +- crates/graphforge-exec/src/session/tests.rs | 122 ++++++++++++ crates/graphforge-storage/src/lib.rs | 1 + crates/graphforge-storage/src/query_spill.rs | 178 ++++++++++++++++++ .../src/query_spill/tests.rs | 96 ++++++++++ docs/development/execution-resource-policy.md | 37 +++- 15 files changed, 766 insertions(+), 44 deletions(-) create mode 100644 crates/graphforge-api/src/runtime_ownership/spill_tests.rs create mode 100644 crates/graphforge-storage/src/query_spill.rs create mode 100644 crates/graphforge-storage/src/query_spill/tests.rs 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..71ab1dbb7 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.session_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..0e57ec36e 100644 --- a/crates/graphforge-api/src/runtime_ownership.rs +++ b/crates/graphforge-api/src/runtime_ownership.rs @@ -32,16 +32,65 @@ impl GraphForge { } } - pub(super) fn session_resource_config(&self) -> graphforge_exec::SessionResourceConfig { - 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, + /// 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, + }) + } + + /// The project scratch directory for query spill, acquired on first use. + /// `None` unless this facade is a writable durable project. + fn query_spill_directory(&self) -> Result, GfError> { + let Some(project) = &self.path else { + return Ok(None); + }; + if self.read_only + || 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 +228,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..4215d4866 --- /dev/null +++ b/crates/graphforge-api/src/runtime_ownership/spill_tests.rs @@ -0,0 +1,175 @@ +//! #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()); +} 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..3b984836d 100644 --- a/crates/graphforge-exec/src/session.rs +++ b/crates/graphforge-exec/src/session.rs @@ -565,6 +565,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()?; @@ -667,14 +672,20 @@ impl ExecutionSession { 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 - { + // #1595: spill goes only where the policy says. 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. + if let (true, Some(dir)) = (resources.spill_enabled, &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); } + } else { + runtime_builder = runtime_builder.with_disk_manager_builder( + datafusion::execution::disk_manager::DiskManagerBuilder::default() + .with_mode(datafusion::execution::disk_manager::DiskManagerMode::Disabled), + ); } let runtime_env = Arc::new( runtime_builder diff --git a/crates/graphforge-exec/src/session/tests.rs b/crates/graphforge-exec/src/session/tests.rs index ec8f74422..a755f59e2 100644 --- a/crates/graphforge-exec/src/session/tests.rs +++ b/crates/graphforge-exec/src/session/tests.rs @@ -742,6 +742,128 @@ 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, + ); + 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, + ); + 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: 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..2d1213708 --- /dev/null +++ b/crates/graphforge-storage/src/query_spill.rs @@ -0,0 +1,178 @@ +//! 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 instance removes its own subdirectory and lock when it is dropped. A +//! crashed process cannot, 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 with no lock file +//! is also stale, since an owner creates and locks its lock file before its +//! subdirectory. 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)?; + let token = uuid::Uuid::new_v4().simple().to_string(); + let lock_path = root.join(format!("{token}.lock")); + let lock = OpenOptions::new() + .write(true) + .create_new(true) + .open(&lock_path) + .map_err(storage)?; + if !try_lock_exclusive(&lock).map_err(storage)? { + return Err(storage("a new scratch lock is already held")); + } + 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); + } +} + +/// Remove every scratch entry whose owner is gone. +fn reclaim_abandoned(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() { + if let Some(token) = name.strip_suffix(".lock").filter(|token| is_token(token)) { + locks.push(token.to_owned()); + } + } else if kind.is_dir() && is_token(&name) { + directories.push(name); + } + } + for token in &locks { + let lock_path = root.join(format!("{token}.lock")); + let lock = match File::open(&lock_path) { + Ok(lock) => lock, + // Another process reclaimed or released it meanwhile. + Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue, + Err(error) => return Err(storage(error)), + }; + if !try_lock_exclusive(&lock).map_err(storage)? { + continue; // Its owner is alive. + } + remove_directory(&root.join(token))?; + remove_file(&lock_path)?; + } + for token in directories.iter().filter(|token| !locks.contains(token)) { + 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..28bfd5e6f --- /dev/null +++ b/crates/graphforge-storage/src/query_spill/tests.rs @@ -0,0 +1,96 @@ +use super::*; + +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(); + // 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}"); +} diff --git a/docs/development/execution-resource-policy.md b/docs/development/execution-resource-policy.md index 22a22cfd1..d5964d603 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,41 @@ 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` if set | 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. +Views that bound their own memory (branch, checkpoint and slice views) set +`enabled: false` and so fail closed the same way. + +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 instance +removes both when it is dropped. Opening another instance reclaims any whose +lock is free, which means the owning process has exited, and leaves live +instances' scratch alone (`graphforge_storage::query_spill`). Scratch holds only +a running query's spill files and is never recovery authority. + +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 +86,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 From f6d88ab9d77852f3b9bb19cfa8ba3d7ebc25868b Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Mon, 28 Sep 2026 03:40:01 +0000 Subject: [PATCH 2/4] refactor(exec): build the session runtime in its own function (#1595) Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/graphforge-exec/src/session.rs | 53 ++++++++++++++++----------- 1 file changed, 31 insertions(+), 22 deletions(-) diff --git a/crates/graphforge-exec/src/session.rs b/crates/graphforge-exec/src/session.rs index 3b984836d..a29c85963 100644 --- a/crates/graphforge-exec/src/session.rs +++ b/crates/graphforge-exec/src/session.rs @@ -344,6 +344,36 @@ 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, +) -> Arc { + 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) { + 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); + } + } else { + runtime_builder = runtime_builder.with_disk_manager_builder( + datafusion::execution::disk_manager::DiskManagerBuilder::default() + .with_mode(datafusion::execution::disk_manager::DiskManagerMode::Disabled), + ); + } + Arc::new( + runtime_builder + .build() + .expect("DataFusion RuntimeEnv construction"), + ) +} + struct QueryEvidenceStream { inner: Option, physical: Arc, @@ -670,28 +700,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); - // #1595: spill goes only where the policy says. 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. - if let (true, Some(dir)) = (resources.spill_enabled, &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); - } - } else { - runtime_builder = runtime_builder.with_disk_manager_builder( - datafusion::execution::disk_manager::DiskManagerBuilder::default() - .with_mode(datafusion::execution::disk_manager::DiskManagerMode::Disabled), - ); - } - 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() From f45df116dee1bb40d53f153d4143dfb1decf53ef Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Mon, 28 Sep 2026 04:07:01 +0000 Subject: [PATCH 3/4] fix(storage,exec,api): address review of the query spill scratch (#1595) - A lock file is created under a temporary name, locked, then renamed into place, so a reclaimer never sees an unheld live lock; a claim a reclaimer raced is retried with a fresh token. A subdirectory counts as lockless only if its lock is absent when examined. - Reclaiming other owners' leftovers is best effort: an unremovable entry is skipped and never blocks acquiring this instance's scratch. - Building a spilling session is fallible: an unusable spill directory or runtime error is a storage error for that query, not a panic. - Read-only views of a durable project (checkpoints, inspection) spill into project scratch like any query. - EXPLAIN renders its plan with a no-spill configuration and never acquires scratch. - Docs: DataFusion's 100 GiB default cap for a configured directory without max_bytes; the read-only and EXPLAIN rules. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/graphforge-api/src/explanation.rs | 2 +- .../graphforge-api/src/runtime_ownership.rs | 24 +++- .../src/runtime_ownership/spill_tests.rs | 42 +++++++ crates/graphforge-exec/src/session.rs | 36 +++--- crates/graphforge-exec/src/session/tests.rs | 37 +++++- crates/graphforge-storage/src/query_spill.rs | 119 +++++++++++++----- .../src/query_spill/tests.rs | 54 ++++++++ docs/development/execution-resource-policy.md | 23 ++-- 8 files changed, 271 insertions(+), 66 deletions(-) diff --git a/crates/graphforge-api/src/explanation.rs b/crates/graphforge-api/src/explanation.rs index 71ab1dbb7..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/runtime_ownership.rs b/crates/graphforge-api/src/runtime_ownership.rs index 0e57ec36e..6be79ca98 100644 --- a/crates/graphforge-api/src/runtime_ownership.rs +++ b/crates/graphforge-api/src/runtime_ownership.rs @@ -71,15 +71,31 @@ impl GraphForge { }) } + /// 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: 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 writable durable project. + /// `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.read_only - || self.lifecycle_mode - != graphforge_storage::filesystem_admission::ProjectLifecycleMode::Durable + if self.lifecycle_mode + != graphforge_storage::filesystem_admission::ProjectLifecycleMode::Durable { return Ok(None); } diff --git a/crates/graphforge-api/src/runtime_ownership/spill_tests.rs b/crates/graphforge-api/src/runtime_ownership/spill_tests.rs index 4215d4866..7c571839e 100644 --- a/crates/graphforge-api/src/runtime_ownership/spill_tests.rs +++ b/crates/graphforge-api/src/runtime_ownership/spill_tests.rs @@ -173,3 +173,45 @@ fn two_instances_of_one_project_do_not_share_scratch() { ); 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-exec/src/session.rs b/crates/graphforge-exec/src/session.rs index a29c85963..938646fb7 100644 --- a/crates/graphforge-exec/src/session.rs +++ b/crates/graphforge-exec/src/session.rs @@ -351,12 +351,13 @@ impl Default for SessionResourceConfig { /// RAM-backed on some hosts. fn session_runtime( resources: &SessionResourceConfig, -) -> Arc { +) -> 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) { - let _ = std::fs::create_dir_all(dir); + 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); @@ -367,11 +368,12 @@ fn session_runtime( .with_mode(datafusion::execution::disk_manager::DiskManagerMode::Disabled), ); } - Arc::new( - runtime_builder - .build() - .expect("DataFusion RuntimeEnv construction"), - ) + // 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 { @@ -509,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(), @@ -517,7 +519,7 @@ impl ExecutionSession { None, OrdinalIdentityConfig::default(), &SessionResourceConfig::default(), - )) + ) } /// Create a session that can execute writes against `dir`. @@ -530,7 +532,7 @@ impl ExecutionSession { dir: PathBuf, mode: OntologyMode, ) -> Result { - Ok(Self::build( + Self::build( catalog, ontology, dir, @@ -538,7 +540,7 @@ impl ExecutionSession { None, OrdinalIdentityConfig::default(), &SessionResourceConfig::default(), - )) + ) } /// Like [`new_with_target`](Self::new_with_target) but reusing a @@ -614,7 +616,7 @@ impl ExecutionSession { } None => OrdinalIdentityConfig::default(), }; - Ok(Self::build( + Self::build( catalog, ontology, dir, @@ -622,7 +624,7 @@ impl ExecutionSession { Some(provider), identity, resources, - )) + ) } fn build( @@ -633,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/` @@ -700,7 +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 runtime_env = session_runtime(resources); + let runtime_env = session_runtime(resources)?; let state = SessionStateBuilder::new() .with_default_features() @@ -723,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, @@ -735,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 a755f59e2..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(); @@ -766,7 +767,8 @@ async fn disabled_spill_refuses_instead_of_spilling_to_the_os_temp_directory() { None, OrdinalIdentityConfig::default(), &resources, - ); + ) + .unwrap(); assert!( !session.ctx.runtime_env().disk_manager.tmp_files_enabled(), "a disabled spill policy must disable the disk manager" @@ -816,7 +818,8 @@ async fn the_spill_cap_bounds_one_query() { 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(); @@ -835,6 +838,34 @@ async fn the_spill_cap_bounds_one_query() { 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] diff --git a/crates/graphforge-storage/src/query_spill.rs b/crates/graphforge-storage/src/query_spill.rs index 2d1213708..66492d5bc 100644 --- a/crates/graphforge-storage/src/query_spill.rs +++ b/crates/graphforge-storage/src/query_spill.rs @@ -13,12 +13,20 @@ //! 3f2c…e9/ DataFusion's spill files for that instance's queries //! ``` //! -//! An instance removes its own subdirectory and lock when it is dropped. A -//! crashed process cannot, 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 with no lock file -//! is also stale, since an owner creates and locks its lock file before its -//! subdirectory. Entries that do not match this grammar are left alone. +//! 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. @@ -80,17 +88,18 @@ impl QuerySpillDirectory { } Err(error) => return Err(storage(error)), } - reclaim_abandoned(&root)?; - let token = uuid::Uuid::new_v4().simple().to_string(); - let lock_path = root.join(format!("{token}.lock")); - let lock = OpenOptions::new() - .write(true) - .create_new(true) - .open(&lock_path) - .map_err(storage)?; - if !try_lock_exclusive(&lock).map_err(storage)? { - return Err(storage("a new scratch lock is already held")); - } + 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); @@ -119,8 +128,42 @@ impl Drop for QuerySpillDirectory { } } -/// Remove every scratch entry whose owner is gone. -fn reclaim_abandoned(root: &Path) -> Result<(), GfError> { +/// 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)?; + 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)? { @@ -131,29 +174,39 @@ fn reclaim_abandoned(root: &Path) -> Result<(), GfError> { continue; }; if kind.is_file() { - if let Some(token) = name.strip_suffix(".lock").filter(|token| is_token(token)) { - locks.push(token.to_owned()); + 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 token in &locks { - let lock_path = root.join(format!("{token}.lock")); - let lock = match File::open(&lock_path) { - Ok(lock) => lock, - // Another process reclaimed or released it meanwhile. - Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue, - Err(error) => return Err(storage(error)), + 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 !try_lock_exclusive(&lock).map_err(storage)? { + if !matches!(try_lock_exclusive(&lock), Ok(true)) { continue; // Its owner is alive. } - remove_directory(&root.join(token))?; - remove_file(&lock_path)?; + if let Some(token) = name.strip_suffix(".lock") { + let _ = remove_directory(&root.join(token)); + } + let _ = remove_file(&lock_path); } - for token in directories.iter().filter(|token| !locks.contains(token)) { - remove_directory(&root.join(token))?; + 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(()) } diff --git a/crates/graphforge-storage/src/query_spill/tests.rs b/crates/graphforge-storage/src/query_spill/tests.rs index 28bfd5e6f..4ff8e23d4 100644 --- a/crates/graphforge-storage/src/query_spill/tests.rs +++ b/crates/graphforge-storage/src/query_spill/tests.rs @@ -50,6 +50,9 @@ fn abandoned_scratch_is_reclaimed_and_other_entries_are_left_alone() { // 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(); @@ -94,3 +97,54 @@ fn a_scratch_root_that_is_not_a_directory_is_refused() { 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()); +} diff --git a/docs/development/execution-resource-policy.md b/docs/development/execution-resource-policy.md index d5964d603..a47cccb32 100644 --- a/docs/development/execution-resource-policy.md +++ b/docs/development/execution-resource-policy.md @@ -59,21 +59,28 @@ state. | 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` if set | same | +| `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. -Views that bound their own memory (branch, checkpoint and slice views) set -`enabled: false` and so fail closed the same way. +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 instance -removes both when it is dropped. Opening another instance reclaims any whose -lock is free, which means the owning process has exited, and leaves live -instances' scratch alone (`graphforge_storage::query_spill`). Scratch holds only -a running query's spill files and is never recovery authority. +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 From 812c10765ec59b2ae2a35d863a9ad75a41ac2eaf Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Mon, 28 Sep 2026 04:10:25 +0000 Subject: [PATCH 4/4] test(storage): prove a raced scratch claim is retried (#1595) CI's shared-directory test hit the race the review described: a reader's reclamation locked another reader's freshly created lock file. The claim now publishes its lock only once held and retries a raced claim; two tests drive each interleaving deterministically through a test-only hook, and fail when the retry is removed. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/graphforge-storage/src/query_spill.rs | 2 + .../src/query_spill/tests.rs | 102 ++++++++++++++++++ 2 files changed, 104 insertions(+) diff --git a/crates/graphforge-storage/src/query_spill.rs b/crates/graphforge-storage/src/query_spill.rs index 66492d5bc..ae15d8c04 100644 --- a/crates/graphforge-storage/src/query_spill.rs +++ b/crates/graphforge-storage/src/query_spill.rs @@ -139,6 +139,8 @@ fn claim_lock(root: &Path) -> Result, GfError> { .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); } diff --git a/crates/graphforge-storage/src/query_spill/tests.rs b/crates/graphforge-storage/src/query_spill/tests.rs index 4ff8e23d4..163569369 100644 --- a/crates/graphforge-storage/src/query_spill/tests.rs +++ b/crates/graphforge-storage/src/query_spill/tests.rs @@ -1,4 +1,33 @@ 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) @@ -148,3 +177,76 @@ fn a_published_lock_is_already_held() { ); 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()); +}