Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion crates/graphforge-api/src/algorithm_writeback.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
1 change: 1 addition & 0 deletions crates/graphforge-api/src/embedding_refresh.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
2 changes: 1 addition & 1 deletion crates/graphforge-api/src/explanation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 })?;
Expand Down
6 changes: 6 additions & 0 deletions crates/graphforge-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<graphforge_storage::ConstructionCpuAdmission>,
/// 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<Mutex<Option<graphforge_storage::query_spill::QuerySpillDirectory>>>,
/// Ensures mutation bursts share one bounded process-local driver thread.
provider_refresh_driver_active: Arc<AtomicBool>,
/// Runtime-only provider recipes capable of refreshing exact lineages.
Expand Down Expand Up @@ -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())),
Expand Down Expand Up @@ -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())),
Expand Down
4 changes: 2 additions & 2 deletions crates/graphforge-api/src/query_execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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()
Expand Down
85 changes: 67 additions & 18 deletions crates/graphforge-api/src/resource_policy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PathBuf>,
/// Optional upper bound on temporary spill bytes.
/// Optional upper bound on one query's temporary spill bytes.
pub max_bytes: Option<u64>,
}

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 {
Expand Down Expand Up @@ -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<PathBuf>,
/// Optional spill byte cap.
pub spill_max_bytes: Option<u64>,
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -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,
Expand All @@ -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]
Expand Down
84 changes: 76 additions & 8 deletions crates/graphforge-api/src/runtime_ownership.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<graphforge_exec::SessionResourceConfig, GfError> {
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<Option<std::path::PathBuf>, 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<tokio::sync::SemaphorePermit<'_>, GfError> {
Expand Down Expand Up @@ -179,3 +244,6 @@ pub(super) fn build_runtime(

#[cfg(test)]
mod tests;

#[cfg(test)]
mod spill_tests;
Loading