use std::collections::{BTreeMap, BTreeSet};
use std::fs::{File, OpenOptions};
use std::path::PathBuf;
use std::sync::{Arc, Condvar, Mutex, mpsc};
use std::thread;
#[cfg(windows)]
use windows_sys::Win32::System::SystemInformation::{GlobalMemoryStatusEx, MEMORYSTATUSEX};
use fs2::FileExt;
pub(crate) use crate::cancellation::CancellationToken;
use crate::error::ForgeError;
use crate::events::{self as event, LifecycleEvent};
use crate::fsutil::is_lock_contended;
use crate::paths::app_home;
use crate::planning::{ExecutionNode, ExecutionPlan, ResourceClaim};
use crate::util::fnv1a;
const MEMORY_TOKEN_MIB: u32 = 512;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum NodeOutcome {
Completed,
Skipped,
Failed,
Blocked,
Cancelled,
}
#[derive(Debug, Clone)]
pub(crate) struct NodeReport {
pub(crate) node: String,
pub(crate) outcome: NodeOutcome,
}
pub(crate) trait NodeRunner: Send + Sync + 'static {
fn run(
&self,
node: &ExecutionNode,
cancellation: &CancellationToken,
) -> Result<NodeOutcome, ForgeError>;
fn manages_resources(&self, _: &ExecutionNode) -> bool {
false
}
fn run_with_resources(
&self,
node: &ExecutionNode,
cancellation: &CancellationToken,
_: &ResourceCoordinator,
) -> Result<NodeOutcome, ForgeError> {
self.run(node, cancellation)
}
}
#[derive(Debug, Clone)]
pub(crate) struct ResourceBudget {
capacities: BTreeMap<String, u32>,
host_capacities: BTreeMap<String, u32>,
}
impl ResourceBudget {
pub(crate) fn for_plan(plan: &ExecutionPlan) -> Self {
let cpu = thread::available_parallelism().map_or(1, |value| value.get() as u32);
let detected_memory = available_memory_mib().unwrap_or(4096);
let automatic_memory = automatic_memory_budget(detected_memory);
let host_memory = total_memory_mib()
.map(automatic_memory_budget)
.unwrap_or(automatic_memory);
let memory = plan
.policy
.max_memory_mib
.map_or(automatic_memory, |limit| automatic_memory.min(limit))
.min(u64::from(u32::MAX)) as u32;
let host_memory = quantize_memory_capacity(host_memory.min(u64::from(u32::MAX)) as u32);
let memory = quantize_memory_capacity(memory);
let host_disk_io = automatic_disk_io_capacity(cpu, host_memory);
let disk_io = automatic_disk_io_capacity(cpu, memory);
let host_capacities = BTreeMap::from([
("network".into(), plan.policy.max_downloads as u32),
("cpu".into(), cpu),
("disk-io".into(), host_disk_io),
("memory-mib".into(), host_memory),
]);
let mut capacities = BTreeMap::from([
("network".into(), plan.policy.max_downloads as u32),
("cpu".into(), cpu),
("disk-io".into(), disk_io),
("memory-mib".into(), memory),
]);
for node in &plan.nodes {
for claim in &node.resources {
capacities.entry(claim.key.clone()).or_insert(1);
}
}
Self {
capacities,
host_capacities,
}
}
pub(crate) fn capacity(&self, resource: &str) -> u32 {
self.capacities.get(resource).copied().unwrap_or(0)
}
}
fn quantize_memory_capacity(memory_mib: u32) -> u32 {
memory_mib
.max(MEMORY_TOKEN_MIB)
.div_euclid(MEMORY_TOKEN_MIB)
.saturating_mul(MEMORY_TOKEN_MIB)
}
fn automatic_memory_budget(available_mib: u64) -> u64 {
available_mib.saturating_mul(60).div_ceil(100).max(512)
}
fn automatic_disk_io_capacity(cpu: u32, memory_mib: u32) -> u32 {
(cpu / 4).max(1).min((memory_mib / 1024).max(1)).min(8)
}
#[cfg(target_os = "linux")]
fn available_memory_mib() -> Option<u64> {
let contents = std::fs::read_to_string("/proc/meminfo").ok()?;
let kib = contents.lines().find_map(|line| {
line.strip_prefix("MemAvailable:")?
.split_whitespace()
.next()?
.parse::<u64>()
.ok()
})?;
Some(kib / 1024)
}
#[cfg(target_os = "linux")]
fn total_memory_mib() -> Option<u64> {
memory_info_mib("MemTotal:")
}
#[cfg(target_os = "linux")]
fn memory_info_mib(key: &str) -> Option<u64> {
let contents = std::fs::read_to_string("/proc/meminfo").ok()?;
let kib = contents.lines().find_map(|line| {
line.strip_prefix(key)?
.split_whitespace()
.next()?
.parse::<u64>()
.ok()
})?;
Some(kib / 1024)
}
#[cfg(target_os = "macos")]
fn available_memory_mib() -> Option<u64> {
let output = std::process::Command::new("/usr/bin/vm_stat")
.output()
.ok()?;
if output.status.success() {
let text = String::from_utf8(output.stdout).ok()?;
let page_size = text
.lines()
.next()?
.split("page size of ")
.nth(1)?
.split_whitespace()
.next()?
.parse::<u64>()
.ok()?;
let pages = text.lines().skip(1).filter_map(|line| {
let (name, value) = line.split_once(':')?;
matches!(
name,
"Pages free" | "Pages inactive" | "Pages speculative" | "Pages purgeable"
)
.then(|| value.trim().trim_end_matches('.').parse::<u64>().ok())
.flatten()
});
let available = pages.sum::<u64>().saturating_mul(page_size) / 1024 / 1024;
if available > 0 {
return Some(available);
}
}
let output = std::process::Command::new("sysctl")
.args(["-n", "hw.memsize"])
.output()
.ok()?;
if !output.status.success() {
return None;
}
String::from_utf8(output.stdout)
.ok()?
.trim()
.parse::<u64>()
.ok()
.map(|bytes| bytes / 1024 / 1024)
}
#[cfg(target_os = "macos")]
fn total_memory_mib() -> Option<u64> {
let output = std::process::Command::new("/usr/sbin/sysctl")
.args(["-n", "hw.memsize"])
.output()
.ok()?;
output
.status
.success()
.then(|| String::from_utf8(output.stdout).ok())
.flatten()?
.trim()
.parse::<u64>()
.ok()
.map(|bytes| bytes / 1024 / 1024)
}
#[cfg(windows)]
fn available_memory_mib() -> Option<u64> {
let mut status = MEMORYSTATUSEX {
dwLength: std::mem::size_of::<MEMORYSTATUSEX>() as u32,
..unsafe { std::mem::zeroed() }
};
(unsafe { GlobalMemoryStatusEx(&mut status) } != 0).then_some(status.ullAvailPhys / 1024 / 1024)
}
#[cfg(windows)]
fn total_memory_mib() -> Option<u64> {
let mut status = MEMORYSTATUSEX {
dwLength: std::mem::size_of::<MEMORYSTATUSEX>() as u32,
..unsafe { std::mem::zeroed() }
};
(unsafe { GlobalMemoryStatusEx(&mut status) } != 0).then_some(status.ullTotalPhys / 1024 / 1024)
}
#[cfg(not(any(target_os = "linux", target_os = "macos", windows)))]
fn available_memory_mib() -> Option<u64> {
None
}
#[cfg(not(any(target_os = "linux", target_os = "macos", windows)))]
fn total_memory_mib() -> Option<u64> {
None
}
pub(crate) struct Scheduler {
jobs: usize,
resources: Arc<ResourcePool>,
lock_root: PathBuf,
host_capacities: BTreeMap<String, u32>,
priorities: BTreeMap<String, u64>,
dynamic_allocation: Arc<Mutex<()>>,
}
#[derive(Clone)]
pub(crate) struct ResourceCoordinator {
resources: Arc<ResourcePool>,
lock_root: PathBuf,
host_capacities: BTreeMap<String, u32>,
dynamic_allocation: Arc<Mutex<()>>,
}
impl ResourceCoordinator {
pub(crate) fn capacity(&self, resource: &str) -> u32 {
self.resources.capacity(resource)
}
pub(crate) fn available(&self, resource: &str) -> u32 {
self.resources.available(resource)
}
pub(crate) fn host_available(&self, resource: &str) -> u32 {
let capacity = self.host_capacities.get(resource).copied().unwrap_or(0);
available_process_slots(resource, capacity, &self.lock_root).unwrap_or(capacity)
}
pub(crate) fn acquire<'a>(
&'a self,
claims: &[ResourceClaim],
cancellation: &CancellationToken,
) -> Result<ResourceLeases<'a>, ForgeError> {
self.resources
.acquire_all(claims, &self.host_capacities, cancellation, &self.lock_root)
}
pub(crate) fn acquire_cargo_build<'a>(
&'a self,
mut claims: Vec<ResourceClaim>,
remaining_builds: usize,
baseline_jobs: usize,
max_builds: usize,
cancellation: &CancellationToken,
) -> Result<(ResourceLeases<'a>, usize), ForgeError> {
let _allocation = self
.dynamic_allocation
.lock()
.map_err(|_| ForgeError::Command("dynamic resource allocator poisoned".into()))?;
while self.available("disk-io") == 0 || self.host_available("disk-io") == 0 {
if cancellation.is_cancelled() {
return Err(ForgeError::Command(
"dynamic resource wait cancelled".into(),
));
}
thread::sleep(std::time::Duration::from_millis(25));
}
let available_cpu = self.available("cpu").min(self.host_available("cpu")).max(1) as usize;
let available_slots = self
.available("disk-io")
.min(self.host_available("disk-io"))
.max(1) as usize;
let slots_to_fill = available_slots.min(max_builds.max(1));
let jobs = if remaining_builds <= slots_to_fill {
(available_cpu / remaining_builds.max(1)).max(1)
} else {
baseline_jobs.max(1).min(available_cpu)
};
claims.push(ResourceClaim {
key: "cpu".into(),
units: jobs as u32,
});
let leases = self.acquire(&claims, cancellation)?;
Ok((leases, jobs))
}
}
impl Scheduler {
pub(crate) fn new(jobs: usize, budget: ResourceBudget) -> Self {
Self {
jobs: jobs.max(1),
resources: Arc::new(ResourcePool::new(budget.capacities)),
lock_root: app_home().join("locks").join("resources"),
host_capacities: budget.host_capacities,
priorities: BTreeMap::new(),
dynamic_allocation: Arc::new(Mutex::new(())),
}
}
pub(crate) fn with_priorities(mut self, priorities: BTreeMap<String, u64>) -> Self {
self.priorities = priorities;
self
}
pub(crate) fn execute<R: NodeRunner>(
&self,
plan: &ExecutionPlan,
runner: Arc<R>,
cancellation: CancellationToken,
) -> Result<Vec<NodeReport>, ForgeError> {
validate_plan(plan)?;
let coordinator = ResourceCoordinator {
resources: Arc::clone(&self.resources),
lock_root: self.lock_root.clone(),
host_capacities: self.host_capacities.clone(),
dynamic_allocation: Arc::clone(&self.dynamic_allocation),
};
let nodes = plan
.nodes
.iter()
.map(|node| (node.id.clone(), node.clone()))
.collect::<BTreeMap<_, _>>();
let mut pending = nodes.keys().cloned().collect::<BTreeSet<_>>();
let mut outcomes = BTreeMap::new();
let mut reports = Vec::new();
thread::scope(|scope| -> Result<(), ForgeError> {
let (sender, receiver) = mpsc::channel();
let mut running = BTreeSet::<String>::new();
while !pending.is_empty() || !running.is_empty() {
let blocked = pending
.iter()
.filter(|id| {
nodes[*id].dependencies.iter().any(|dependency| {
outcomes.get(dependency).is_some_and(|outcome| {
matches!(
outcome,
NodeOutcome::Failed
| NodeOutcome::Blocked
| NodeOutcome::Cancelled
)
})
})
})
.cloned()
.collect::<Vec<_>>();
for id in blocked {
pending.remove(&id);
outcomes.insert(id.clone(), NodeOutcome::Blocked);
event::emit(
Some(&id),
Some(&nodes[&id].component),
None,
LifecycleEvent::Blocked {
dependency: nodes[&id]
.dependencies
.iter()
.find(|dependency| {
matches!(
outcomes.get(*dependency),
Some(
NodeOutcome::Failed
| NodeOutcome::Blocked
| NodeOutcome::Cancelled
)
)
})
.cloned()
.unwrap_or_else(|| "unknown".into()),
},
);
reports.push(NodeReport {
node: id,
outcome: NodeOutcome::Blocked,
});
}
if cancellation.is_cancelled() {
for id in std::mem::take(&mut pending) {
outcomes.insert(id.clone(), NodeOutcome::Cancelled);
event::emit(
Some(&id),
Some(&nodes[&id].component),
None,
LifecycleEvent::Cancelled {
forced: cancellation.is_forced(),
},
);
reports.push(NodeReport {
node: id,
outcome: NodeOutcome::Cancelled,
});
}
}
let free_slots = self.jobs.saturating_sub(running.len());
let mut ready = pending
.iter()
.filter(|id| {
nodes[*id].dependencies.iter().all(|dependency| {
matches!(
outcomes.get(dependency),
Some(NodeOutcome::Completed | NodeOutcome::Skipped)
)
})
})
.cloned()
.collect::<Vec<_>>();
ready.sort_by(|left, right| {
self.priorities
.get(right)
.copied()
.unwrap_or(0)
.cmp(&self.priorities.get(left).copied().unwrap_or(0))
.then_with(|| left.cmp(right))
});
ready.truncate(free_slots);
for id in ready {
pending.remove(&id);
let node = nodes[&id].clone();
running.insert(id.clone());
let runner = Arc::clone(&runner);
let resources = Arc::clone(&self.resources);
let cancellation = cancellation.clone();
let coordinator = coordinator.clone();
let sender = sender.clone();
scope.spawn(move || {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let started = std::time::Instant::now();
event::emit(
Some(&id),
Some(&node.component),
None,
LifecycleEvent::Queued,
);
let result = if runner.manages_resources(&node) {
runner.run_with_resources(&node, &cancellation, &coordinator)
} else if node.resources.is_empty() {
runner.run(&node, &cancellation)
} else {
let resource = node
.resources
.iter()
.map(|claim| claim.key.as_str())
.collect::<Vec<_>>()
.join(",");
event::emit(
Some(&id),
Some(&node.component),
None,
LifecycleEvent::ResourceWait {
resource: resource.clone(),
},
);
let wait_started = std::time::Instant::now();
resources
.acquire_all(
&node.resources,
&self.host_capacities,
&cancellation,
&self.lock_root,
)
.and_then(|_leases| {
event::emit(
Some(&id),
Some(&node.component),
None,
LifecycleEvent::ResourceAcquired {
resource,
wait_ms: wait_started.elapsed().as_millis(),
},
);
runner.run(&node, &cancellation)
})
};
(id, started.elapsed(), result)
}))
.map_err(|_| ForgeError::Command("scheduler worker panicked".into()));
let _ = sender.send(result);
});
}
if running.is_empty() {
if pending.is_empty() {
break;
}
return Err(ForgeError::Config("scheduler dependency deadlock".into()));
}
let result = receiver
.recv()
.map_err(|_| ForgeError::Command("scheduler result channel closed".into()))?;
match result {
Ok((id, elapsed, Ok(outcome))) => {
running.remove(&id);
outcomes.insert(id.clone(), outcome);
let payload = match outcome {
NodeOutcome::Completed => LifecycleEvent::Completed {
elapsed_ms: elapsed.as_millis(),
},
NodeOutcome::Skipped => LifecycleEvent::Skipped {
reason: "node does not need to run".into(),
},
NodeOutcome::Blocked => LifecycleEvent::Blocked {
dependency: "runner".into(),
},
NodeOutcome::Cancelled => LifecycleEvent::Cancelled {
forced: cancellation.is_forced(),
},
NodeOutcome::Failed => LifecycleEvent::Failed {
code: "BF2001".into(),
message: "node execution failed".into(),
},
};
event::emit(Some(&id), Some(&nodes[&id].component), None, payload);
reports.push(NodeReport { node: id, outcome });
}
Ok((id, _, Err(error))) => {
running.remove(&id);
outcomes.insert(id.clone(), NodeOutcome::Failed);
event::emit(
Some(&id),
Some(&nodes[&id].component),
None,
LifecycleEvent::Failed {
code: "BF2001".into(),
message: error.to_string(),
},
);
reports.push(NodeReport {
node: id,
outcome: NodeOutcome::Failed,
});
}
Err(error) => return Err(error),
}
}
Ok(())
})?;
Ok(reports)
}
}
struct ResourcePool {
state: Mutex<BTreeMap<String, u32>>,
capacities: BTreeMap<String, u32>,
wake: Condvar,
}
impl ResourcePool {
fn new(capacities: BTreeMap<String, u32>) -> Self {
Self {
state: Mutex::new(capacities.clone()),
capacities,
wake: Condvar::new(),
}
}
fn capacity(&self, resource: &str) -> u32 {
self.capacities.get(resource).copied().unwrap_or(0)
}
fn available(&self, resource: &str) -> u32 {
self.state
.lock()
.ok()
.and_then(|state| state.get(resource).copied())
.unwrap_or(0)
}
fn acquire_all<'a>(
&'a self,
claims: &[ResourceClaim],
host_capacities: &BTreeMap<String, u32>,
cancellation: &CancellationToken,
lock_root: &std::path::Path,
) -> Result<ResourceLeases<'a>, ForgeError> {
let claims = normalize_claims(claims, &self.capacities, host_capacities)?;
let mut state = self
.state
.lock()
.map_err(|_| ForgeError::Command("resource pool poisoned".into()))?;
for claim in &claims {
if !is_capacity_resource(&claim.key) {
state.entry(claim.key.clone()).or_insert(claim.units);
}
}
if let Some(claim) = claims.iter().find(|claim| {
self.capacities
.get(&claim.key)
.copied()
.unwrap_or(claim.units)
< claim.units
}) {
return Err(ForgeError::Config(format!(
"resource request exceeds the budget: {} requires {}",
claim.key, claim.units
)));
}
drop(state);
let files = acquire_process_locks(&claims, host_capacities, cancellation, lock_root)?;
let mut state = self
.state
.lock()
.map_err(|_| ForgeError::Command("resource pool poisoned".into()))?;
loop {
if cancellation.is_cancelled() {
return Err(ForgeError::Command("resource wait cancelled".into()));
}
if claims
.iter()
.all(|claim| state.get(&claim.key).copied().unwrap_or(0) >= claim.units)
{
for claim in &claims {
let available = state.get_mut(&claim.key).ok_or_else(|| {
ForgeError::Config(format!(
"resource budget is missing resource: {}",
claim.key
))
})?;
*available -= claim.units;
}
return Ok(ResourceLeases {
pool: self,
claims,
files,
});
}
let (next, _) = self
.wake
.wait_timeout(state, std::time::Duration::from_millis(50))
.map_err(|_| ForgeError::Command("resource pool poisoned".into()))?;
state = next;
}
}
}
fn normalize_claims(
claims: &[ResourceClaim],
local_capacities: &BTreeMap<String, u32>,
host_capacities: &BTreeMap<String, u32>,
) -> Result<Vec<ResourceClaim>, ForgeError> {
let mut merged = BTreeMap::<String, u32>::new();
for claim in claims {
if claim.units == 0 {
return Err(ForgeError::Config(format!(
"resource request must be greater than zero: {}",
claim.key
)));
}
let units = merged.entry(claim.key.clone()).or_default();
*units = units.saturating_add(claim.units);
}
merged
.into_iter()
.map(|(key, units)| {
let units = if key == "memory-mib" {
let obtainable = local_capacities
.get(&key)
.copied()
.unwrap_or(MEMORY_TOKEN_MIB)
.min(
host_capacities
.get(&key)
.copied()
.unwrap_or(MEMORY_TOKEN_MIB),
);
units
.div_ceil(MEMORY_TOKEN_MIB)
.saturating_mul(MEMORY_TOKEN_MIB)
.min(obtainable)
} else {
units
};
if units == 0 {
return Err(ForgeError::Config(format!(
"resource capacity is zero: {key}"
)));
}
Ok(ResourceClaim { key, units })
})
.collect()
}
pub(crate) struct ResourceLeases<'a> {
pool: &'a ResourcePool,
claims: Vec<ResourceClaim>,
files: Vec<File>,
}
impl Drop for ResourceLeases<'_> {
fn drop(&mut self) {
for file in &self.files {
let _ = FileExt::unlock(file);
}
if let Ok(mut state) = self.pool.state.lock() {
for claim in &self.claims {
*state.entry(claim.key.clone()).or_default() += claim.units;
}
self.pool.wake.notify_all();
}
}
}
fn acquire_process_locks(
claims: &[ResourceClaim],
capacities: &BTreeMap<String, u32>,
cancellation: &CancellationToken,
root: &std::path::Path,
) -> Result<Vec<File>, ForgeError> {
if claims.is_empty() {
return Ok(Vec::new());
}
std::fs::create_dir_all(root).map_err(|source| ForgeError::Io {
path: root.to_path_buf(),
source,
})?;
let requests = process_lock_requests(claims, capacities)?;
let mut attempt = 0_u64;
loop {
if cancellation.is_cancelled() {
return Err(ForgeError::Command(
"cancelled while waiting for cross-process resource capacity".into(),
));
}
let mut files = Vec::new();
let mut complete = true;
for request in &requests {
let mut acquired = 0;
for slot in 0..request.slots {
let path = root.join(format!("{}-{slot}.lock", request.key));
let file = open_lock_file(&path)?;
match file.try_lock_exclusive() {
Ok(()) => {
files.push(file);
acquired += 1;
if acquired == request.units {
break;
}
}
Err(error) if is_lock_contended(&error) => {}
Err(source) => return Err(ForgeError::Io { path, source }),
}
}
if acquired != request.units {
complete = false;
break;
}
}
if complete {
return Ok(files);
}
for file in &files {
let _ = FileExt::unlock(file);
}
attempt = attempt.wrapping_add(1);
let jitter = (u64::from(std::process::id()) + attempt * 17) % 23;
thread::sleep(std::time::Duration::from_millis(27 + jitter));
}
}
fn available_process_slots(
resource: &str,
capacity: u32,
root: &std::path::Path,
) -> Result<u32, ForgeError> {
if capacity == 0 {
return Ok(0);
}
std::fs::create_dir_all(root).map_err(|source| ForgeError::Io {
path: root.to_path_buf(),
source,
})?;
let key = format!("capacity-{}", safe_lock_key(resource));
let mut available = 0_u32;
for slot in 0..capacity {
let path = root.join(format!("{key}-{slot}.lock"));
let file = open_lock_file(&path)?;
match file.try_lock_exclusive() {
Ok(()) => {
available = available.saturating_add(1);
let _ = FileExt::unlock(&file);
}
Err(error) if is_lock_contended(&error) => {}
Err(source) => return Err(ForgeError::Io { path, source }),
}
}
Ok(available)
}
struct ProcessLockRequest {
key: String,
units: u32,
slots: u32,
}
fn process_lock_requests(
claims: &[ResourceClaim],
capacities: &BTreeMap<String, u32>,
) -> Result<Vec<ProcessLockRequest>, ForgeError> {
let mut requests = Vec::new();
for claim in claims {
let safe = safe_lock_key(&claim.key);
if is_capacity_resource(&claim.key) {
let capacity = capacities.get(&claim.key).copied().ok_or_else(|| {
ForgeError::Config(format!(
"resource budget is missing resource: {}",
claim.key
))
})?;
let granularity = if claim.key == "memory-mib" {
MEMORY_TOKEN_MIB
} else {
1
};
requests.push(ProcessLockRequest {
key: format!("capacity-{safe}"),
units: claim.units.div_ceil(granularity),
slots: capacity / granularity,
});
} else {
requests.push(ProcessLockRequest {
key: format!("named-{safe}"),
units: 1,
slots: 1,
});
}
}
requests.sort_by(|left, right| left.key.cmp(&right.key));
Ok(requests)
}
fn safe_lock_key(key: &str) -> String {
let readable = key
.chars()
.map(|character| {
if character.is_ascii_alphanumeric() {
character
} else {
'-'
}
})
.collect::<String>();
format!("{readable}-{:016x}", fnv1a(key))
}
fn open_lock_file(path: &std::path::Path) -> Result<File, ForgeError> {
OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(path)
.map_err(|source| ForgeError::Io {
path: path.to_path_buf(),
source,
})
}
fn is_capacity_resource(key: &str) -> bool {
matches!(key, "network" | "cpu" | "memory-mib" | "disk-io")
}
fn validate_plan(plan: &ExecutionPlan) -> Result<(), ForgeError> {
let ids = plan
.nodes
.iter()
.map(|node| node.id.as_str())
.collect::<BTreeSet<_>>();
if ids.len() != plan.nodes.len() {
return Err(ForgeError::Config(
"execution plan has duplicate node ids".into(),
));
}
for node in &plan.nodes {
if node
.dependencies
.iter()
.any(|dependency| !ids.contains(dependency.as_str()))
{
return Err(ForgeError::Config(format!(
"node {} has missing dependency",
node.id
)));
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use crate::cancellation::CancellationToken;
use crate::config::schema::NetworkPolicy;
use crate::execution::scheduler::{
ExecutionNode, ExecutionPlan, ForgeError, NodeOutcome, NodeRunner, ResourceBudget,
ResourceClaim, ResourceCoordinator, ResourcePool, Scheduler, acquire_process_locks,
automatic_disk_io_capacity, automatic_memory_budget, available_process_slots,
process_lock_requests, quantize_memory_capacity,
};
use crate::planning::{NodeKind, PlanPolicy, TargetPlatform};
use crate::util::now_secs;
use std::thread;
struct TrackingRunner {
active: AtomicUsize,
peak: AtomicUsize,
fail: Option<String>,
}
struct TimingRunner {
started: std::time::Instant,
intervals: Mutex<BTreeMap<String, (Duration, Duration)>>,
}
impl NodeRunner for TimingRunner {
fn run(
&self,
node: &ExecutionNode,
_: &CancellationToken,
) -> Result<NodeOutcome, ForgeError> {
let start = self.started.elapsed();
thread::sleep(if node.id == "b" {
Duration::from_millis(300)
} else {
Duration::from_millis(60)
});
self.intervals
.lock()
.unwrap()
.insert(node.id.clone(), (start, self.started.elapsed()));
Ok(NodeOutcome::Completed)
}
}
impl NodeRunner for TrackingRunner {
fn run(
&self,
node: &ExecutionNode,
_: &CancellationToken,
) -> Result<NodeOutcome, ForgeError> {
if self.fail.as_deref() == Some(&node.id) {
return Ok(NodeOutcome::Failed);
}
let active = self.active.fetch_add(1, Ordering::SeqCst) + 1;
self.peak.fetch_max(active, Ordering::SeqCst);
thread::sleep(Duration::from_millis(30));
self.active.fetch_sub(1, Ordering::SeqCst);
Ok(NodeOutcome::Completed)
}
}
fn node(id: &str, dependencies: &[&str], resource: &str) -> ExecutionNode {
ExecutionNode {
id: id.into(),
component: id.into(),
kind: NodeKind::Acquire,
dependencies: dependencies.iter().map(|v| (*v).into()).collect(),
resources: vec![ResourceClaim {
key: resource.into(),
units: 1,
}],
}
}
fn plan(nodes: Vec<ExecutionNode>) -> ExecutionPlan {
ExecutionPlan {
profile: "test".into(),
target: TargetPlatform::host(),
config_hash: String::new(),
plan_hash: String::new(),
policy: PlanPolicy {
network: NetworkPolicy::Online,
max_parallel: 4,
max_downloads: 2,
max_memory_mib: Some(1024),
},
certificate_preflight: None,
components: Vec::new(),
unsupported_components: Vec::new(),
nodes,
environment: Vec::new(),
apt_mirror: None,
origins: BTreeMap::new(),
}
}
#[test]
fn independent_nodes_run_concurrently_but_shared_lock_serializes() {
let first_plan = plan(vec![node("a", &[], "lock-a"), node("b", &[], "lock-b")]);
let runner = Arc::new(TrackingRunner {
active: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
fail: None,
});
Scheduler::new(4, ResourceBudget::for_plan(&first_plan))
.execute(
&first_plan,
Arc::clone(&runner),
CancellationToken::default(),
)
.unwrap();
assert_eq!(runner.peak.load(Ordering::SeqCst), 2);
let plan = plan(vec![node("a", &[], "shared"), node("b", &[], "shared")]);
let runner = Arc::new(TrackingRunner {
active: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
fail: None,
});
Scheduler::new(4, ResourceBudget::for_plan(&plan))
.execute(&plan, Arc::clone(&runner), CancellationToken::default())
.unwrap();
assert_eq!(runner.peak.load(Ordering::SeqCst), 1);
}
#[test]
fn automatic_memory_budget_keeps_host_headroom_and_honors_floor() {
assert_eq!(automatic_memory_budget(16_000), 9_600);
assert_eq!(automatic_memory_budget(256), 512);
assert_eq!(quantize_memory_capacity(513), 512);
}
#[test]
fn staged_resource_lease_releases_capacity_before_transaction_finishes() {
struct StagedRunner {
activation_active: AtomicUsize,
build_overlapped_activation: AtomicUsize,
}
impl NodeRunner for StagedRunner {
fn run(
&self,
_: &ExecutionNode,
_: &CancellationToken,
) -> Result<NodeOutcome, ForgeError> {
unreachable!()
}
fn manages_resources(&self, _: &ExecutionNode) -> bool {
true
}
fn run_with_resources(
&self,
node: &ExecutionNode,
cancellation: &CancellationToken,
coordinator: &ResourceCoordinator,
) -> Result<NodeOutcome, ForgeError> {
if node.id == "b" {
while self.activation_active.load(Ordering::SeqCst) == 0 {
thread::sleep(Duration::from_millis(5));
}
}
{
let _build = coordinator.acquire(
&[ResourceClaim {
key: "cpu".into(),
units: 1,
}],
cancellation,
)?;
thread::sleep(Duration::from_millis(60));
}
self.activation_active.fetch_add(1, Ordering::SeqCst);
if node.id == "a" {
let deadline = std::time::Instant::now() + Duration::from_millis(200);
while self.build_overlapped_activation.load(Ordering::SeqCst) == 0
&& std::time::Instant::now() < deadline
{
thread::sleep(Duration::from_millis(5));
}
} else if self.activation_active.load(Ordering::SeqCst) > 1 {
self.build_overlapped_activation
.fetch_add(1, Ordering::SeqCst);
}
thread::sleep(Duration::from_millis(20));
self.activation_active.fetch_sub(1, Ordering::SeqCst);
Ok(NodeOutcome::Completed)
}
}
let mut plan = plan(vec![node("a", &[], "cpu"), node("b", &[], "cpu")]);
plan.policy.max_parallel = 2;
let mut budget = ResourceBudget::for_plan(&plan);
budget.capacities.insert("cpu".into(), 1);
let runner = Arc::new(StagedRunner {
activation_active: AtomicUsize::new(0),
build_overlapped_activation: AtomicUsize::new(0),
});
Scheduler::new(2, budget)
.execute(&plan, Arc::clone(&runner), CancellationToken::default())
.unwrap();
assert_eq!(runner.build_overlapped_activation.load(Ordering::SeqCst), 1);
}
#[test]
fn completed_worker_is_backfilled_without_waiting_for_the_slow_peer() {
let plan = plan(vec![
node("a", &[], "backfill-lock-a"),
node("b", &[], "backfill-lock-b"),
node("c", &[], "backfill-lock-c"),
]);
let runner = Arc::new(TimingRunner {
started: std::time::Instant::now(),
intervals: Mutex::new(BTreeMap::new()),
});
Scheduler::new(2, ResourceBudget::for_plan(&plan))
.execute(&plan, Arc::clone(&runner), CancellationToken::default())
.unwrap();
let intervals = runner.intervals.lock().unwrap();
assert!(intervals["c"].0 < intervals["b"].1);
}
#[test]
fn historical_priority_starts_heavier_ready_work_first() {
let plan = plan(vec![
node("a", &[], "priority-a"),
node("b", &[], "priority-b"),
]);
let runner = Arc::new(TimingRunner {
started: std::time::Instant::now(),
intervals: Mutex::new(BTreeMap::new()),
});
Scheduler::new(1, ResourceBudget::for_plan(&plan))
.with_priorities(BTreeMap::from([("b".into(), 10_000)]))
.execute(&plan, Arc::clone(&runner), CancellationToken::default())
.unwrap();
let intervals = runner.intervals.lock().unwrap();
assert!(intervals["b"].0 < intervals["a"].0);
}
#[test]
fn mixed_large_and_short_claims_complete_without_starvation() {
struct MixedRunner;
impl NodeRunner for MixedRunner {
fn run(
&self,
node: &ExecutionNode,
_: &CancellationToken,
) -> Result<NodeOutcome, ForgeError> {
thread::sleep(if node.id == "large" {
Duration::from_millis(40)
} else {
Duration::from_millis(2)
});
Ok(NodeOutcome::Completed)
}
}
let mut large = node("large", &[], "cpu");
large.resources[0].units = 3;
let mut nodes = vec![large];
nodes.extend((0..16).map(|index| node(&format!("short-{index:02}"), &[], "cpu")));
let mut plan = plan(nodes);
plan.policy.max_parallel = 8;
let mut budget = ResourceBudget::for_plan(&plan);
budget.capacities.insert("cpu".into(), 4);
budget.host_capacities.insert("cpu".into(), 4);
let mut scheduler = Scheduler::new(8, budget);
scheduler.lock_root = std::env::temp_dir().join(format!(
"bot-forge-fairness-{}-{}",
std::process::id(),
now_secs()
));
let reports = scheduler
.execute(&plan, Arc::new(MixedRunner), CancellationToken::default())
.unwrap();
assert_eq!(reports.len(), 17);
assert!(
reports
.iter()
.all(|report| report.outcome == NodeOutcome::Completed)
);
let _ = std::fs::remove_dir_all(scheduler.lock_root);
}
#[test]
fn failure_blocks_dependents() {
let plan = plan(vec![node("a", &[], "a"), node("b", &["a"], "b")]);
let runner = Arc::new(TrackingRunner {
active: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
fail: Some("a".into()),
});
let reports = Scheduler::new(2, ResourceBudget::for_plan(&plan))
.execute(&plan, runner, CancellationToken::default())
.unwrap();
assert_eq!(
reports.iter().find(|r| r.node == "b").unwrap().outcome,
NodeOutcome::Blocked
);
}
#[test]
fn worker_errors_are_attributed_to_their_own_node() {
struct ErrorRunner;
impl NodeRunner for ErrorRunner {
fn run(
&self,
node: &ExecutionNode,
_: &CancellationToken,
) -> Result<NodeOutcome, ForgeError> {
if node.id == "a" {
return Err(ForgeError::Command("expected failure".into()));
}
Ok(NodeOutcome::Completed)
}
}
let plan = plan(vec![
node("a", &[], "a"),
node("b", &[], "b"),
node("c", &["a"], "c"),
]);
let reports = Scheduler::new(2, ResourceBudget::for_plan(&plan))
.execute(&plan, Arc::new(ErrorRunner), CancellationToken::default())
.unwrap();
assert_eq!(
reports
.iter()
.find(|report| report.node == "a")
.unwrap()
.outcome,
NodeOutcome::Failed
);
assert_eq!(
reports
.iter()
.find(|report| report.node == "b")
.unwrap()
.outcome,
NodeOutcome::Completed
);
assert_eq!(
reports
.iter()
.find(|report| report.node == "c")
.unwrap()
.outcome,
NodeOutcome::Blocked
);
}
#[test]
fn cancellation_wakes_resource_waiters_quickly() {
let plan = plan(vec![node("a", &[], "shared"), node("b", &[], "shared")]);
let cancellation = CancellationToken::default();
let signal = cancellation.clone();
thread::spawn(move || {
thread::sleep(Duration::from_millis(10));
signal.cancel();
});
let runner = Arc::new(TrackingRunner {
active: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
fail: None,
});
let started = std::time::Instant::now();
let _ =
Scheduler::new(2, ResourceBudget::for_plan(&plan)).execute(&plan, runner, cancellation);
assert!(started.elapsed() < Duration::from_millis(200));
}
#[test]
fn multi_resource_claims_use_stable_key_order() {
let mut claims = [
ResourceClaim {
key: "bin:z".into(),
units: 1,
},
ResourceClaim {
key: "bin:a".into(),
units: 1,
},
];
claims.sort_by(|left, right| left.key.cmp(&right.key));
assert_eq!(
claims
.iter()
.map(|claim| claim.key.as_str())
.collect::<Vec<_>>(),
["bin:a", "bin:z"]
);
}
#[test]
fn capacity_claims_expand_to_host_wide_token_slots() {
let capacities = BTreeMap::from([
("cpu".into(), 8),
("memory-mib".into(), 4096),
("network".into(), 4),
("disk-io".into(), 2),
]);
let requests = process_lock_requests(
&[
ResourceClaim {
key: "cpu".into(),
units: 3,
},
ResourceClaim {
key: "memory-mib".into(),
units: 1024,
},
],
&capacities,
)
.unwrap();
assert_eq!(requests[0].units, 3);
assert_eq!(requests[0].slots, 8);
assert_eq!(requests[1].units, 2);
assert_eq!(requests[1].slots, 8);
}
#[test]
fn memory_claim_is_rounded_and_clamped_to_obtainable_capacity() {
let pool = ResourcePool::new(BTreeMap::from([("memory-mib".into(), 512)]));
let root = std::env::temp_dir().join(format!("bot-forge-oversized-{}", std::process::id()));
let lease = pool
.acquire_all(
&[ResourceClaim {
key: "memory-mib".into(),
units: 513,
}],
&BTreeMap::from([("memory-mib".into(), 4096)]),
&CancellationToken::default(),
&root,
)
.unwrap();
assert_eq!(lease.claims[0].units, 512);
drop(lease);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn duplicate_claims_are_merged_before_capacity_validation() {
let pool = ResourcePool::new(BTreeMap::from([("cpu".into(), 1)]));
let root =
std::env::temp_dir().join(format!("bot-forge-duplicate-claims-{}", std::process::id()));
let result = pool.acquire_all(
&[
ResourceClaim {
key: "cpu".into(),
units: 1,
},
ResourceClaim {
key: "cpu".into(),
units: 1,
},
],
&BTreeMap::from([("cpu".into(), 8)]),
&CancellationToken::default(),
&root,
);
assert!(matches!(result, Err(error) if error.to_string().contains("exceeds the budget")));
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn host_network_capacity_tracks_the_plan_limit() {
let mut plan = plan(Vec::new());
plan.policy.max_downloads = 17;
let budget = ResourceBudget::for_plan(&plan);
assert_eq!(budget.host_capacities["network"], 17);
}
#[test]
fn disk_io_capacity_scales_with_machine_resources() {
assert_eq!(automatic_disk_io_capacity(2, 8192), 1);
assert_eq!(automatic_disk_io_capacity(18, 16384), 4);
assert_eq!(automatic_disk_io_capacity(64, 4096), 4);
assert_eq!(automatic_disk_io_capacity(64, 65536), 8);
}
#[test]
fn cargo_build_allocation_rebalances_only_new_processes() {
let root = std::env::temp_dir().join(format!(
"bot-forge-dynamic-cargo-{}-{}",
std::process::id(),
now_secs()
));
let capacities = BTreeMap::from([
("cpu".into(), 18),
("disk-io".into(), 4),
("network".into(), 8),
("memory-mib".into(), 8192),
]);
let coordinator = ResourceCoordinator {
resources: Arc::new(ResourcePool::new(capacities.clone())),
lock_root: root.clone(),
host_capacities: capacities,
dynamic_allocation: Arc::new(Mutex::new(())),
};
let cancellation = CancellationToken::default();
let claims = || {
vec![ResourceClaim {
key: "disk-io".into(),
units: 1,
}]
};
let (first, first_jobs) = coordinator
.acquire_cargo_build(claims(), 8, 4, 4, &cancellation)
.unwrap();
let (second, second_jobs) = coordinator
.acquire_cargo_build(claims(), 8, 4, 4, &cancellation)
.unwrap();
assert_eq!((first_jobs, second_jobs), (4, 4));
drop(first);
let (tail, tail_jobs) = coordinator
.acquire_cargo_build(claims(), 1, 4, 4, &cancellation)
.unwrap();
assert_eq!(tail_jobs, 14);
drop(second);
drop(tail);
let (_last, last_jobs) = coordinator
.acquire_cargo_build(claims(), 1, 4, 4, &cancellation)
.unwrap();
assert_eq!(last_jobs, 18);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn capacity_slots_are_shared_across_independent_resource_pools() {
let root = std::env::temp_dir().join(format!(
"bot-forge-capacity-locks-{}-{}",
std::process::id(),
now_secs()
));
let claims = [ResourceClaim {
key: "cpu".into(),
units: 1,
}];
let capacities = BTreeMap::from([("cpu".into(), 1)]);
let cancellation = CancellationToken::default();
let first = acquire_process_locks(&claims, &capacities, &cancellation, &root).unwrap();
let cancelled = cancellation.clone();
let signal = cancellation.clone();
let path = root.clone();
let waiter = thread::spawn(move || {
let started = std::time::Instant::now();
let result = acquire_process_locks(&claims, &capacities, &cancelled, &path);
(started.elapsed(), result.is_err())
});
thread::sleep(Duration::from_millis(80));
signal.cancel();
let (elapsed, was_cancelled) = waiter.join().unwrap();
assert!(was_cancelled);
assert!(elapsed >= Duration::from_millis(50));
drop(first);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn host_availability_observes_other_process_capacity_leases() {
let root = std::env::temp_dir().join(format!(
"bot-forge-host-availability-{}-{}",
std::process::id(),
now_secs()
));
let claims = [ResourceClaim {
key: "cpu".into(),
units: 2,
}];
let capacities = BTreeMap::from([("cpu".into(), 4)]);
let leases =
acquire_process_locks(&claims, &capacities, &CancellationToken::default(), &root)
.unwrap();
assert_eq!(available_process_slots("cpu", 4, &root).unwrap(), 2);
drop(leases);
assert_eq!(available_process_slots("cpu", 4, &root).unwrap(), 4);
let _ = std::fs::remove_dir_all(root);
}
#[test]
fn filesystem_lock_names_preserve_distinct_resource_identities() {
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let plan = plan(vec![
node("a", &[], &format!("collision:{nonce}")),
node("b", &[], &format!("collision-{nonce}")),
]);
let runner = Arc::new(TrackingRunner {
active: AtomicUsize::new(0),
peak: AtomicUsize::new(0),
fail: None,
});
let started = std::time::Instant::now();
let reports = Scheduler::new(2, ResourceBudget::for_plan(&plan))
.execute(&plan, Arc::clone(&runner), CancellationToken::default())
.unwrap();
assert_eq!(reports.len(), 2);
assert_eq!(runner.peak.load(Ordering::SeqCst), 2);
assert!(started.elapsed() < Duration::from_secs(2));
}
#[test]
fn randomized_dags_complete_without_deadlock() {
struct ImmediateRunner;
impl NodeRunner for ImmediateRunner {
fn run(
&self,
_: &ExecutionNode,
_: &CancellationToken,
) -> Result<NodeOutcome, ForgeError> {
Ok(NodeOutcome::Completed)
}
}
for seed in 1..=64_u64 {
let mut state = seed;
let mut nodes = Vec::new();
for index in 0..40 {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
let dependency = (index > 0).then(|| (state as usize) % index);
nodes.push(node(
&format!("node-{index}"),
&dependency
.map(|value| format!("node-{value}"))
.iter()
.map(String::as_str)
.collect::<Vec<_>>(),
&format!("lock-{}", state % 7),
));
}
let plan = plan(nodes);
let runner = Arc::new(ImmediateRunner);
let reports = Scheduler::new(8, ResourceBudget::for_plan(&plan))
.execute(&plan, runner, CancellationToken::default())
.unwrap();
assert_eq!(reports.len(), 40, "seed {seed}");
assert!(
reports
.iter()
.all(|report| report.outcome == NodeOutcome::Completed),
"seed {seed}"
);
}
}
}