diff --git a/README.md b/README.md index 4f96f8aa767..0b316498193 100644 --- a/README.md +++ b/README.md @@ -207,6 +207,7 @@ the current branch. | `IRONCLAW_REBORN_PROFILE` | Boot profile selector. Supported values: `local-dev`, `local-dev-yolo`, `hosted-single-tenant`, `hosted-single-tenant-volume`, `production`, `migration-dry-run`. | | `IRONCLAW_REBORN_POSTGRES_URL` | Production PostgreSQL storage URL when `[storage].backend = "postgres"` and `[storage].url_env` names this variable. Keep it out of `config.toml`; remote providers must use TLS. | | `IRONCLAW_REBORN_POSTGRES_POOL_MAX_SIZE` | Optional override for the Reborn PostgreSQL client pool size. Use this when a managed provider enforces a small session-pool cap. | +| `IRONCLAW_RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH` | Optional `true`/`1`/`yes`/`on` toggle that skips durable resource-governor reserve/reconcile/release writes when no finite limits are configured. Defaults to false so production keeps durable accounting unless this is explicitly enabled. | | `IRONCLAW_FILESYSTEM_POSTGRES_MIGRATION_CONNECT_MAX_WAIT_SECS` | Optional startup wait window for Postgres filesystem migration connection retries. Defaults to 300 seconds. | | `IRONCLAW_REBORN_SECRET_MASTER_KEY` | Production Reborn secret master key when `[storage].secret_master_key_env` names this variable. Keep it independent from the database URL and out of `config.toml`. | | `IRONCLAW_REBORN_LOG` | Tracing filter for the Reborn binary, for example `debug,ironclaw_reborn=trace`. | diff --git a/crates/ironclaw_reborn_composition/src/error.rs b/crates/ironclaw_reborn_composition/src/error.rs index 839b49f273b..9a65a18859b 100644 --- a/crates/ironclaw_reborn_composition/src/error.rs +++ b/crates/ironclaw_reborn_composition/src/error.rs @@ -57,6 +57,9 @@ impl From for RebornBuildError { impl From for RebornBuildError { fn from(error: crate::RebornCompositionError) -> Self { match error { + crate::RebornCompositionError::InvalidConfig { reason } => { + Self::InvalidConfig { reason } + } crate::RebornCompositionError::MissingSecretMasterKey => Self::MissingSecretMasterKey, crate::RebornCompositionError::Mount(error) => Self::Mount(error), crate::RebornCompositionError::Filesystem(error) => Self::Filesystem(error), diff --git a/crates/ironclaw_reborn_composition/src/factory.rs b/crates/ironclaw_reborn_composition/src/factory.rs index 77c1011c054..4d88980edfb 100644 --- a/crates/ironclaw_reborn_composition/src/factory.rs +++ b/crates/ironclaw_reborn_composition/src/factory.rs @@ -1,4 +1,6 @@ // arch-exempt: large_file, needs Reborn composition helper extraction, plan #4469 +#[cfg(test)] +use std::cell::RefCell; use std::{ collections::VecDeque, path::{Path, PathBuf}, @@ -89,11 +91,13 @@ use ironclaw_product_workflow::{ LifecycleProductSurfaceContext, ProductAuthTurnGateResumeDispatcher, ProjectService, }; use ironclaw_projects::ProjectRepository; -use ironclaw_resources::InMemoryResourceGovernor; #[cfg(any(feature = "libsql", feature = "postgres"))] use ironclaw_resources::{ BroadcastBudgetEventSink, BudgetGateStore, FilesystemBudgetGateStore, - FilesystemResourceGovernorStore, PersistentResourceGovernor, ResourceGovernor, + FilesystemResourceGovernorStore, ResourceGovernor, +}; +use ironclaw_resources::{ + InMemoryResourceGovernor, PersistentResourceGovernor, ResourceGovernorStore, }; #[cfg(any(feature = "libsql", feature = "postgres"))] use ironclaw_run_state::{FilesystemApprovalRequestStore, FilesystemRunStateStore}; @@ -236,6 +240,16 @@ type LocalDevResourceGovernor = #[cfg(not(any(feature = "libsql", feature = "postgres")))] type LocalDevResourceGovernor = InMemoryResourceGovernor; +#[cfg(test)] +thread_local! { + static RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV_OVERRIDE: RefCell> = + const { RefCell::new(None) }; +} + +#[cfg(any(feature = "libsql", feature = "postgres", test))] +const RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV: &str = + "IRONCLAW_RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH"; + #[cfg(any(feature = "libsql", feature = "postgres"))] type LocalDevRunStateStore = FilesystemRunStateStore; #[cfg(not(any(feature = "libsql", feature = "postgres")))] @@ -1699,6 +1713,86 @@ fn local_dev_extension_installation_state_path( }) } +#[cfg_attr(not(any(feature = "libsql", feature = "postgres")), allow(dead_code))] +fn apply_resource_governor_unlimited_fast_path( + governor: PersistentResourceGovernor, +) -> Result, String> +where + S: ResourceGovernorStore, +{ + if resource_governor_unlimited_fast_path_enabled_from_env()? { + Ok(governor.with_unlimited_fast_path()) + } else { + Ok(governor) + } +} + +#[cfg_attr(not(any(feature = "libsql", feature = "postgres")), allow(dead_code))] +fn resource_governor_unlimited_fast_path_enabled_from_env() -> Result { + #[cfg(not(any(feature = "libsql", feature = "postgres")))] + { + return Ok(false); + } + + #[cfg(any(feature = "libsql", feature = "postgres"))] + match resource_governor_unlimited_fast_path_env_value() { + Ok(Some(value)) => parse_bool_env_value(&value).ok_or_else(|| { + format!( + "{RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV} must be one of true, false, 1, 0, yes, no, on, or off" + ) + }), + Ok(None) => Ok(false), + Err(reason) => Err(reason), + } +} + +#[cfg(any(feature = "libsql", feature = "postgres"))] +fn resource_governor_unlimited_fast_path_env_value() -> Result, String> { + #[cfg(test)] + if let Some(value) = + RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV_OVERRIDE.with(|value| value.borrow().clone()) + { + return Ok(Some(value)); + } + + match std::env::var(RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV) { + Ok(value) => Ok(Some(value)), + Err(std::env::VarError::NotPresent) => Ok(None), + Err(std::env::VarError::NotUnicode(_)) => Err(format!( + "{RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV} must be valid UTF-8" + )), + } +} + +#[cfg(test)] +fn set_resource_governor_unlimited_fast_path_env_override_for_test( + value: impl Into, +) -> Result { + RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV_OVERRIDE + .with(|override_value| *override_value.borrow_mut() = Some(value.into())); + Ok(ResourceGovernorFastPathEnvOverrideGuard) +} + +#[cfg(test)] +struct ResourceGovernorFastPathEnvOverrideGuard; + +#[cfg(test)] +impl Drop for ResourceGovernorFastPathEnvOverrideGuard { + fn drop(&mut self) { + RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV_OVERRIDE + .with(|override_value| *override_value.borrow_mut() = None); + } +} + +#[cfg(any(feature = "libsql", feature = "postgres"))] +fn parse_bool_env_value(value: &str) -> Option { + match value.trim().to_ascii_lowercase().as_str() { + "" | "0" | "false" | "no" | "off" => Some(false), + "1" | "true" | "yes" | "on" => Some(true), + _ => None, + } +} + #[cfg(any(feature = "libsql", feature = "postgres"))] fn build_local_dev_store_graph( input: RebornLocalDevStoreGraphInput, @@ -1756,12 +1850,13 @@ fn build_local_dev_store_graph( let budget_gate_store: Arc = Arc::new(FilesystemBudgetGateStore::new( Arc::clone(&scoped_filesystem), )); - let resource_governor: Arc = Arc::new( - PersistentResourceGovernor::new(FilesystemResourceGovernorStore::new(Arc::clone( - &scoped_filesystem, - ))) - .with_event_sink(Arc::clone(&budget_event_sink)), - ); + let resource_governor = + apply_resource_governor_unlimited_fast_path(PersistentResourceGovernor::new( + FilesystemResourceGovernorStore::new(Arc::clone(&scoped_filesystem)), + )) + .map_err(|reason| RebornBuildError::InvalidConfig { reason })? + .with_event_sink(Arc::clone(&budget_event_sink)); + let resource_governor: Arc = Arc::new(resource_governor); let skill_mounts = skill_management_mount_view().map_err(|error| RebornBuildError::InvalidConfig { reason: error.to_string(), @@ -3524,7 +3619,11 @@ where ) .await?; let resource_store = FilesystemResourceGovernorStore::new(Arc::clone(&scoped_filesystem)); - let governor = Arc::new(PersistentResourceGovernor::new(resource_store)); + let governor = apply_resource_governor_unlimited_fast_path(PersistentResourceGovernor::new( + resource_store, + )) + .map_err(|reason| crate::RebornCompositionError::InvalidConfig { reason })?; + let governor = Arc::new(governor); let capability_leases = Arc::new(FilesystemCapabilityLeaseStore::new(Arc::clone( &scoped_filesystem, ))); @@ -3784,12 +3883,13 @@ where let thread_service: Arc = Arc::new( FilesystemSessionThreadService::new(Arc::clone(&stores.scoped_filesystem)), ); - let resource_governor = Arc::new( - PersistentResourceGovernor::new(FilesystemResourceGovernorStore::new(Arc::clone( - &stores.scoped_filesystem, - ))) - .with_event_sink(Arc::clone(&budget_event_sink)), - ); + let resource_governor = + apply_resource_governor_unlimited_fast_path(PersistentResourceGovernor::new( + FilesystemResourceGovernorStore::new(Arc::clone(&stores.scoped_filesystem)), + )) + .map_err(|reason| RebornBuildError::InvalidConfig { reason })? + .with_event_sink(Arc::clone(&budget_event_sink)); + let resource_governor = Arc::new(resource_governor); let production_resource_governor: Arc = resource_governor.clone(); let budget_gate_store: Arc = Arc::new(FilesystemBudgetGateStore::new( Arc::clone(&stores.scoped_filesystem), @@ -4100,8 +4200,8 @@ mod tests { CapabilityGrant, CapabilityGrantId, CapabilityId, CapabilitySet, EffectKind, ExecutionContext, ExtensionId, GrantConstraints, InvocationId, MountAlias, MountGrant, MountPermissions, NetworkPolicy, NetworkScheme, NetworkTargetPattern, Principal, - ResourceEstimate, ResourceScope, RuntimeKind, ScopedPath, SecretHandle, TenantId, - TrustClass, UserId, VirtualPath, + ResourceEstimate, ResourceScope, ResourceUsage, RuntimeKind, ScopedPath, SecretHandle, + TenantId, TrustClass, UserId, VirtualPath, }; #[cfg(any(feature = "libsql", feature = "postgres"))] use ironclaw_host_api::{ @@ -4115,6 +4215,8 @@ mod tests { }; use ironclaw_product_workflow::{LifecyclePackageKind, LifecyclePackageRef, LifecyclePhase}; use ironclaw_trust::{AuthorityCeiling, EffectiveTrustClass, TrustDecision, TrustProvenance}; + #[cfg(any(feature = "libsql", feature = "postgres"))] + use rust_decimal_macros::dec; #[cfg(feature = "libsql")] use secrecy::ExposeSecret; @@ -4127,6 +4229,105 @@ mod tests { runtime::SKILL_ACTIVATE_CAPABILITY_ID, }; + #[cfg(any(feature = "libsql", feature = "postgres"))] + #[test] + fn resource_governor_fast_path_env_parser_accepts_documented_values() { + for value in ["true", "TRUE", "1", "yes", "on", " true "] { + assert_eq!(parse_bool_env_value(value), Some(true), "value={value}"); + } + for value in ["", "false", "FALSE", "0", "no", "off", " off "] { + assert_eq!(parse_bool_env_value(value), Some(false), "value={value}"); + } + } + + #[cfg(any(feature = "libsql", feature = "postgres"))] + #[test] + fn resource_governor_fast_path_env_parser_rejects_unknown_values() { + assert_eq!(parse_bool_env_value("maybe"), None); + } + + #[cfg(any(feature = "libsql", feature = "postgres"))] + #[test] + fn build_reborn_services_rejects_invalid_resource_governor_fast_path_env() { + let _override = set_resource_governor_unlimited_fast_path_env_override_for_test("maybe") + .expect("resource governor env override"); + let dir = tempfile::tempdir().expect("tempdir"); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("tokio runtime"); + + let result = runtime.block_on(build_reborn_services(RebornBuildInput::local_dev( + "resource-governor-invalid-env-owner", + dir.path().join("local-dev"), + ))); + + let Err(RebornBuildError::InvalidConfig { reason }) = result else { + panic!("expected invalid config for resource governor fast-path env"); + }; + assert!(reason.contains(RESOURCE_GOVERNOR_UNLIMITED_FAST_PATH_ENV)); + } + + #[cfg(any(feature = "libsql", feature = "postgres"))] + #[test] + fn build_reborn_services_applies_resource_governor_fast_path_env() { + let _override = set_resource_governor_unlimited_fast_path_env_override_for_test("true") + .expect("resource governor env override"); + let dir = tempfile::tempdir().expect("tempdir"); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("tokio runtime"); + + let services = runtime + .block_on(build_reborn_services(RebornBuildInput::local_dev( + "resource-governor-enabled-env-owner", + dir.path().join("local-dev"), + ))) + .expect("local-dev services build"); + let local_runtime = services.local_runtime.as_ref().expect("local runtime"); + let scope = ResourceScope { + tenant_id: TenantId::new("resource-governor-tenant").expect("tenant"), + user_id: UserId::new("resource-governor-user").expect("user"), + agent_id: None, + project_id: None, + mission_id: None, + thread_id: None, + invocation_id: InvocationId::new(), + }; + let account = ironclaw_resources::ResourceAccount::tenant(scope.tenant_id.clone()); + + let reservation = local_runtime + .resource_governor + .reserve( + scope, + ResourceEstimate { + usd: Some(dec!(0.10)), + ..ResourceEstimate::default() + }, + ) + .expect("reservation"); + local_runtime + .resource_governor + .reconcile( + reservation.id, + ResourceUsage { + usd: dec!(0.10), + ..ResourceUsage::default() + }, + ) + .expect("reconcile"); + + assert_eq!( + local_runtime + .resource_governor + .usage_for(&account) + .expect("usage") + .usd, + dec!(0.10) + ); + } + #[test] fn extension_installation_state_path_stays_legacy_for_local_dev() { let path = diff --git a/crates/ironclaw_reborn_composition/src/lib.rs b/crates/ironclaw_reborn_composition/src/lib.rs index 0769dc10d7c..f28afb927ea 100644 --- a/crates/ironclaw_reborn_composition/src/lib.rs +++ b/crates/ironclaw_reborn_composition/src/lib.rs @@ -882,6 +882,8 @@ where #[derive(Debug, Error)] pub enum RebornCompositionError { + #[error("invalid reborn production configuration: {reason}")] + InvalidConfig { reason: String }, #[error( "reborn production composition requires a configured or keychain-resolvable secret master key" )] diff --git a/crates/ironclaw_resources/src/cas_snapshot.rs b/crates/ironclaw_resources/src/cas_snapshot.rs index 991ee8ee4f3..1fb02924d0d 100644 --- a/crates/ironclaw_resources/src/cas_snapshot.rs +++ b/crates/ironclaw_resources/src/cas_snapshot.rs @@ -130,6 +130,44 @@ where self.update_with_scope::(self.scope.clone(), update) } + /// Read the underlying snapshot through the store's default scope without + /// writing it back. + pub(crate) fn inspect(&self, inspect: U) -> Result + where + S: Snapshot, + T: Send + 'static, + E: StorageError, + U: FnOnce(&S) -> Result + Send + 'static, + { + self.inspect_with_scope::(self.scope.clone(), inspect) + } + + /// Read the underlying snapshot through a caller-supplied scope without + /// writing it back. + pub(crate) fn inspect_with_scope( + &self, + scope: ResourceScope, + inspect: U, + ) -> Result + where + S: Snapshot, + T: Send + 'static, + E: StorageError, + U: FnOnce(&S) -> Result + Send + 'static, + { + let filesystem = Arc::clone(&self.filesystem); + let path_str = self.path_str; + let worker_cell = Arc::clone(&self.worker); + let worker_name = self.worker_thread_name; + run_on_worker(&worker_cell, worker_name, move || async move { + let path = ScopedPath::new(path_str.to_string()).map_err(|error| { + E::storage(format!("invalid snapshot path {path_str}: {error}")) + })?; + let snapshot = read_snapshot::(&filesystem, &scope, &path).await?; + inspect(&snapshot) + }) + } + /// Run a read-modify-write transaction against the underlying /// snapshot using a caller-supplied [`ResourceScope`]. /// @@ -187,11 +225,10 @@ where U: FnOnce(&mut S) -> Result, { let (mut snapshot, expectation) = match filesystem.get(scope, path).await { - Ok(Some(versioned)) => { - let snapshot: S = serde_json::from_slice(&versioned.entry.body) - .map_err(|error| E::storage(format!("decode snapshot: {error}")))?; - (snapshot, CasExpectation::Version(versioned.version)) - } + Ok(Some(versioned)) => ( + decode_snapshot::(&versioned.entry.body)?, + CasExpectation::Version(versioned.version), + ), Ok(None) => (S::fresh(), CasExpectation::Absent), Err(error) => return Err(E::storage_from(error)), }; @@ -208,6 +245,31 @@ where } } +async fn read_snapshot( + filesystem: &ScopedFilesystem, + scope: &ResourceScope, + path: &ScopedPath, +) -> Result +where + F: RootFilesystem, + S: Snapshot, + E: StorageError, +{ + match filesystem.get(scope, path).await { + Ok(Some(versioned)) => decode_snapshot::(&versioned.entry.body), + Ok(None) => Ok(S::fresh()), + Err(error) => Err(E::storage_from(error)), + } +} + +fn decode_snapshot(body: &[u8]) -> Result +where + S: Snapshot, + E: StorageError, +{ + serde_json::from_slice(body).map_err(|error| E::storage(format!("decode snapshot: {error}"))) +} + enum PutError { VersionMismatch, Other(E), diff --git a/crates/ironclaw_resources/src/filesystem_store.rs b/crates/ironclaw_resources/src/filesystem_store.rs index 41f84153c55..78021169d35 100644 --- a/crates/ironclaw_resources/src/filesystem_store.rs +++ b/crates/ironclaw_resources/src/filesystem_store.rs @@ -100,6 +100,15 @@ where self.store .update::(update) } + + fn inspect(&self, inspect: U) -> Result + where + T: Send + 'static, + U: FnOnce(&ResourceGovernorSnapshot) -> Result + Send + 'static, + { + self.store + .inspect::(inspect) + } } // --------------------------------------------------------------------------- diff --git a/crates/ironclaw_resources/src/lib.rs b/crates/ironclaw_resources/src/lib.rs index 90759415e88..6d09db1a433 100644 --- a/crates/ironclaw_resources/src/lib.rs +++ b/crates/ironclaw_resources/src/lib.rs @@ -804,6 +804,11 @@ pub trait ResourceGovernorStore: Send + Sync + 'static { where T: Send + 'static, F: FnOnce(&mut ResourceGovernorSnapshot) -> Result + Send + 'static; + + fn inspect(&self, inspect: F) -> Result + where + T: Send + 'static, + F: FnOnce(&ResourceGovernorSnapshot) -> Result + Send + 'static; } /// File-backed resource-governor store using a stable sidecar lock file around @@ -849,6 +854,33 @@ impl ResourceGovernorStore for JsonFileResourceGovernorStore { (Ok(_), Err(error)) => Err(error), } } + + fn inspect(&self, inspect: F) -> Result + where + T: Send + 'static, + F: FnOnce(&ResourceGovernorSnapshot) -> Result + Send + 'static, + { + let lock_path = lock_path_for(&self.path); + if let Some(parent) = lock_path.parent() { + std::fs::create_dir_all(parent).map_err(storage_error)?; + } + let lock_file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .open(lock_path) + .map_err(storage_error)?; + lock_file.lock_shared().map_err(storage_error)?; + + let result = read_file_snapshot(&self.path).and_then(|snapshot| inspect(&snapshot)); + let unlock_result = lock_file.unlock().map_err(storage_error); + match (result, unlock_result) { + (Ok(value), Ok(())) => Ok(value), + (Err(error), _) => Err(error), + (Ok(_), Err(error)) => Err(error), + } + } } fn lock_path_for(path: &Path) -> PathBuf { @@ -1065,6 +1097,8 @@ where { store: S, clock: Arc, + unlimited_fast_path: bool, + unlimited_state: Mutex, /// Optional sink that receives `BudgetEvent`s as reservations, /// reconciliations, warnings, and denials happen. Wired by /// composition; defaults to [`NoOpBudgetEventSink`] so the @@ -1081,6 +1115,8 @@ where Self { store, clock: Arc::new(SystemClock), + unlimited_fast_path: false, + unlimited_state: Mutex::new(UnlimitedFastPathState::default()), event_sink: Arc::new(NoOpBudgetEventSink), } } @@ -1092,10 +1128,25 @@ where Self { store, clock, + unlimited_fast_path: false, + unlimited_state: Mutex::new(UnlimitedFastPathState::default()), event_sink: Arc::new(NoOpBudgetEventSink), } } + /// Keep unlimited/no-quota reservation bookkeeping process-local. + /// + /// When no finite resource limits are configured, the persistent governor + /// does not need durable snapshot writes to enforce limits. This opt-in + /// path preserves same-process reservation lifecycle checks while avoiding + /// synchronous writes on reserve/reconcile/release. As soon as any finite + /// limit is present in the durable snapshot, new reservations use the + /// durable path again. + pub fn with_unlimited_fast_path(mut self) -> Self { + self.unlimited_fast_path = true; + self + } + /// Plug in an audit/SSE sink. Every `reserve`, `reconcile`, /// `release`, warning, and approval/denial emits a [`BudgetEvent`] /// to this sink (parity with @@ -1113,6 +1164,24 @@ where limits: ResourceLimits, ) -> Result<(), ResourceError> { let now = self.clock.now(); + if self.unlimited_fast_path { + let mut local = self.lock_unlimited_state()?; + let local_state = local.state.clone(); + let local_initialized = local.initialized; + let updated_state = self.store.update(move |snapshot| { + if local_initialized + && !resource_state_has_finite_limits(&snapshot.state) + && !resource_state_has_durable_activity(&snapshot.state) + { + snapshot.state = local_state; + } + set_limit_in_state(&mut snapshot.state, account, limits, now); + Ok(snapshot.state.clone()) + })?; + local.state = updated_state; + local.initialized = true; + return Ok(()); + } self.store.update(move |snapshot| { set_limit_in_state(&mut snapshot.state, account, limits, now); Ok(()) @@ -1122,10 +1191,9 @@ where pub fn reserved_for(&self, account: &ResourceAccount) -> Result { let account = account.clone(); let now = self.clock.now(); - self.store.update(move |snapshot| { - advance_period_if_rolled_over(&mut snapshot.state, &account, now); - Ok(snapshot - .state + self.update_active_state(move |state| { + advance_period_if_rolled_over(state, &account, now); + Ok(state .reserved_by_account .get(&account) .cloned() @@ -1136,16 +1204,82 @@ where pub fn usage_for(&self, account: &ResourceAccount) -> Result { let account = account.clone(); let now = self.clock.now(); - self.store.update(move |snapshot| { - advance_period_if_rolled_over(&mut snapshot.state, &account, now); - Ok(snapshot - .state + self.update_active_state(move |state| { + advance_period_if_rolled_over(state, &account, now); + Ok(state .usage_by_account .get(&account) .cloned() .unwrap_or_default()) }) } + + fn unlimited_fast_path_seed(&self) -> Result, ResourceError> { + if !self.unlimited_fast_path { + return Ok(None); + } + self.store.inspect(|snapshot| { + if resource_state_has_finite_limits(&snapshot.state) + || resource_state_has_durable_activity(&snapshot.state) + { + Ok(None) + } else { + Ok(Some(snapshot.state.clone())) + } + }) + } + + fn update_active_state(&self, update: F) -> Result + where + T: Send + 'static, + F: FnOnce(&mut ResourceState) -> Result + Send + 'static, + { + if let Some(seed) = self.unlimited_fast_path_seed()? { + let mut local = self.lock_unlimited_state()?; + if !local.initialized { + local.state = seed; + local.initialized = true; + } + return update(&mut local.state); + } + self.store + .update(move |snapshot| update(&mut snapshot.state)) + } + + fn close_fast_path_or_durable( + &self, + local_close: L, + durable_close: D, + ) -> Result + where + T: Send + 'static, + L: FnOnce(&mut ResourceState) -> Result, + D: FnOnce(&mut ResourceGovernorSnapshot) -> Result + Send + 'static, + { + if let Some(seed) = self.unlimited_fast_path_seed()? { + let mut local = self.lock_unlimited_state()?; + if !local.initialized { + local.state = seed; + local.initialized = true; + } + match local_close(&mut local.state) { + Ok(value) => return Ok(value), + Err(ResourceError::UnknownReservation { .. }) => {} + Err(error) => return Err(error), + } + } + self.store.update(durable_close) + } + + fn lock_unlimited_state( + &self, + ) -> Result, ResourceError> { + self.unlimited_state + .lock() + .map_err(|_| ResourceError::Storage { + reason: "resource governor unlimited fast-path state lock poisoned".to_string(), + }) + } } impl ResourceGovernor for PersistentResourceGovernor @@ -1181,8 +1315,8 @@ where reservation_id: ResourceReservationId, ) -> Result { let now = self.clock.now(); - let result = self.store.update(move |snapshot| { - reserve_with_outcome_in_state(&mut snapshot.state, scope, estimate, reservation_id, now) + let result = self.update_active_state(move |state| { + reserve_with_outcome_in_state(state, scope, estimate, reservation_id, now) }); emit_reserve_events(self.event_sink.as_ref(), &result, now); result @@ -1194,9 +1328,11 @@ where actual: ResourceUsage, ) -> Result { let now = self.clock.now(); - let result = self.store.update(move |snapshot| { - reconcile_in_state(&mut snapshot.state, reservation_id, actual, now) - }); + let local_actual = actual.clone(); + let result = self.close_fast_path_or_durable( + move |state| reconcile_in_state(state, reservation_id, local_actual, now), + move |snapshot| reconcile_in_state(&mut snapshot.state, reservation_id, actual, now), + ); if let Ok(receipt) = &result { self.event_sink.emit(BudgetEvent::Reconciled { account: most_specific_account(&receipt.scope), @@ -1212,9 +1348,10 @@ where reservation_id: ResourceReservationId, ) -> Result { let now = self.clock.now(); - let result = self - .store - .update(move |snapshot| release_in_state(&mut snapshot.state, reservation_id, now)); + let result = self.close_fast_path_or_durable( + move |state| release_in_state(state, reservation_id, now), + move |snapshot| release_in_state(&mut snapshot.state, reservation_id, now), + ); if let Ok(receipt) = &result { self.event_sink.emit(BudgetEvent::Released { account: most_specific_account(&receipt.scope), @@ -1231,13 +1368,7 @@ where ) -> Result, ResourceError> { let account = account.clone(); let now = self.clock.now(); - self.store.update(move |snapshot| { - Ok(account_snapshot_in_state( - &mut snapshot.state, - &account, - now, - )) - }) + self.update_active_state(move |state| Ok(account_snapshot_in_state(state, &account, now))) } } @@ -1279,6 +1410,22 @@ struct ResourceState { period_anchors: HashMap>, } +#[derive(Debug, Clone, Default, PartialEq)] +struct UnlimitedFastPathState { + state: ResourceState, + initialized: bool, +} + +fn resource_state_has_finite_limits(state: &ResourceState) -> bool { + state.limits.values().any(|limits| !limits.is_unlimited()) +} + +fn resource_state_has_durable_activity(state: &ResourceState) -> bool { + !state.reserved_by_account.is_empty() + || !state.usage_by_account.is_empty() + || !state.reservations.is_empty() +} + /// Snapshot of accumulated period-scoped spend + reserved. /// /// Returned by [`ResourceGovernor`] query helpers so callers can render diff --git a/crates/ironclaw_resources/tests/resource_governor_contract.rs b/crates/ironclaw_resources/tests/resource_governor_contract.rs index e2550d9edbf..faef26809e2 100644 --- a/crates/ironclaw_resources/tests/resource_governor_contract.rs +++ b/crates/ironclaw_resources/tests/resource_governor_contract.rs @@ -23,6 +23,16 @@ impl ResourceGovernorStore for AlwaysFailingStore { reason: "forced durable write failure".to_string(), }) } + + fn inspect(&self, _inspect: F) -> Result + where + T: Send + 'static, + F: FnOnce(&ResourceGovernorSnapshot) -> Result + Send + 'static, + { + Err(ResourceError::Storage { + reason: "forced durable read failure".to_string(), + }) + } } #[test] @@ -722,6 +732,96 @@ fn persistent_governor_reloads_active_holds_and_usage_from_store() { )); } +#[test] +fn persistent_governor_unlimited_fast_path_avoids_durable_writes_until_finite_limit() { + let dir = tempdir().unwrap(); + let path = dir.path().join("resource-governor.json"); + let scope = sample_scope("tenant1", "user1", Some("project1")); + let account = ResourceAccount::tenant(scope.tenant_id.clone()); + + let governor = PersistentResourceGovernor::new(JsonFileResourceGovernorStore::new(&path)) + .with_unlimited_fast_path(); + let reservation = governor + .reserve( + scope.clone(), + ResourceEstimate { + usd: Some(dec!(0.20)), + concurrency_slots: Some(1), + ..ResourceEstimate::default() + }, + ) + .unwrap(); + assert!( + !path.exists(), + "unlimited fast path should not create the durable governor snapshot" + ); + + governor + .reconcile( + reservation.id, + ResourceUsage { + usd: dec!(0.20), + ..ResourceUsage::default() + }, + ) + .unwrap(); + assert!( + matches!( + governor.release(reservation.id), + Err(ResourceError::ReservationClosed { + status: ReservationStatus::Reconciled, + .. + }) + ), + "same-process lifecycle checks are still enforced" + ); + assert!( + !path.exists(), + "reconcile on the unlimited fast path should not create the durable snapshot" + ); + + let active = governor + .reserve( + scope.clone(), + ResourceEstimate { + concurrency_slots: Some(1), + ..ResourceEstimate::default() + }, + ) + .unwrap(); + assert!( + !path.exists(), + "active unlimited fast-path reservations should stay process-local before finite limits" + ); + + governor + .set_limit( + account.clone(), + ResourceLimits { + max_concurrency_slots: Some(1), + ..ResourceLimits::default() + }, + ) + .unwrap(); + let denied = governor + .reserve( + scope, + ResourceEstimate { + concurrency_slots: Some(1), + ..ResourceEstimate::default() + }, + ) + .unwrap_err(); + assert!(matches!( + denied, + ResourceError::LimitExceeded { + denial, + .. + } if denial.dimension == ResourceDimension::ConcurrencySlots + )); + governor.release(active.id).unwrap(); +} + #[test] fn persistent_governor_serializes_concurrent_reservations_across_handles() { let dir = tempdir().unwrap(); diff --git a/tools/ironclaw_stress/README.md b/tools/ironclaw_stress/README.md index 5b726d7dc86..7a99c29a926 100644 --- a/tools/ironclaw_stress/README.md +++ b/tools/ironclaw_stress/README.md @@ -184,7 +184,9 @@ fields: ### bottleneck-finder -The broad scan for libsql and general local bottleneck discovery: +The broad scan for libsql and general local bottleneck discovery. User-turn +cases use separate user threads by default; hot-thread serialization is kept as +an explicit preset so it does not distort capacity results. ```bash cargo run -p ironclaw_stress --release -- \ @@ -198,7 +200,6 @@ Cases include: - `resource-contention` - `chat-baseline` -- `hot-thread` - `large-context` - `tool-heavy` - `tool-wait` @@ -244,6 +245,7 @@ cargo run -p ironclaw_stress --release -- \ cargo run -p ironclaw_stress --release -- \ --backend libsql \ --preset chat-baseline \ + --users 128 \ --ramp-concurrency 64 \ --max-p95-ms 500 \ --max-failure-rate 0.01 \ @@ -251,6 +253,10 @@ cargo run -p ironclaw_stress --release -- \ --bottleneck-report ``` +For user-turn capacity runs, keep `--users >= --concurrency`. With the default +`--active-thread-count 0`, the harness partitions users/threads across workers +so concurrent workers do not target the same thread. + ### Test hot-thread behavior ```bash diff --git a/tools/ironclaw_stress/src/main.rs b/tools/ironclaw_stress/src/main.rs index 73a1935353c..0aa906da316 100644 --- a/tools/ironclaw_stress/src/main.rs +++ b/tools/ironclaw_stress/src/main.rs @@ -903,6 +903,15 @@ fn validate_args(args: &Args) -> Result<(), String> { if args.sweep_users.contains(&0) { return Err("--sweep-users values must be greater than 0".to_string()); } + let max_concurrency = args + .sweep_concurrency + .iter() + .copied() + .max() + .unwrap_or(args.concurrency); + let max_concurrency = args + .ramp_concurrency + .map_or(max_concurrency, |ramp_max| max_concurrency.max(ramp_max)); let min_user_count = args.sweep_users.iter().copied().min().unwrap_or(args.users); let max_active_thread_count = args .sweep_active_thread_count @@ -919,6 +928,20 @@ fn validate_args(args: &Args) -> Result<(), String> { .to_string(), ); } + if args.scenario.is_user_turn() + && args + .sweep_active_thread_count + .iter() + .copied() + .chain(std::iter::once(args.active_thread_count)) + .any(|active_thread_count| active_thread_count == 0) + && max_concurrency > min_user_count + { + return Err( + "user-turn scenarios with --active-thread-count 0 require --users to be greater than or equal to --concurrency" + .to_string(), + ); + } if args.sweep_context_max_messages.contains(&0) { return Err("--sweep-context-max-messages values must be greater than 0".to_string()); } @@ -1754,7 +1777,9 @@ where let view = resource_mount_view(run_id)?; let scoped = Arc::new(ScopedFilesystem::with_fixed_view(root, view)); let store = FilesystemResourceGovernorStore::new(scoped); - Ok(Arc::new(PersistentResourceGovernor::new(store))) + Ok(Arc::new( + PersistentResourceGovernor::new(store).with_unlimited_fast_path(), + )) } fn resource_mount_view(run_id: &str) -> Result { diff --git a/tools/ironclaw_stress/src/suite.rs b/tools/ironclaw_stress/src/suite.rs index 60f4768e35e..ac23f0f7767 100644 --- a/tools/ironclaw_stress/src/suite.rs +++ b/tools/ironclaw_stress/src/suite.rs @@ -187,12 +187,6 @@ pub(crate) fn build_cases(suite: StressSuite) -> Vec { preset: Some(StressPreset::ChatBaseline), scenario: Scenario::ChatTurn, }, - SuiteCase { - label: "hot-thread", - probe: "same-thread-serialization", - preset: Some(StressPreset::HotThread), - scenario: Scenario::ChatTurn, - }, SuiteCase { label: "large-context", probe: "context-read-amplification", @@ -243,12 +237,6 @@ pub(crate) fn build_cases(suite: StressSuite) -> Vec { preset: Some(StressPreset::ChatBaseline), scenario: Scenario::ChatTurn, }, - SuiteCase { - label: "postgres-hot-thread-pool", - probe: "postgres-hot-thread-pool", - preset: Some(StressPreset::HotThread), - scenario: Scenario::ChatTurn, - }, SuiteCase { label: "postgres-context-pool", probe: "postgres-context-read-pool", @@ -296,10 +284,6 @@ fn apply_case(base_args: &Args, case: &SuiteCase, case_args: &mut Args, run_id: } match case.label { - "hot-thread" => { - case_args.active_thread_count = 1; - case_args.span_log_failures = true; - } "large-context" => { case_args.prefill_threads = base_args.users.clamp(1, 100); case_args.prefill_turns_per_thread = base_args @@ -334,11 +318,6 @@ fn apply_case(base_args: &Args, case: &SuiteCase, case_args: &mut Args, run_id: "postgres-chat-pool" => { apply_postgres_pool_pressure_defaults(base_args, case_args); } - "postgres-hot-thread-pool" => { - apply_postgres_pool_pressure_defaults(base_args, case_args); - case_args.active_thread_count = 1; - case_args.span_log_failures = true; - } "postgres-context-pool" => { apply_postgres_pool_pressure_defaults(base_args, case_args); case_args.prefill_threads = base_args.users.clamp(1, 100); diff --git a/tools/ironclaw_stress/src/synthetic.rs b/tools/ironclaw_stress/src/synthetic.rs index e39a0b663e6..2a0b0734cfd 100644 --- a/tools/ironclaw_stress/src/synthetic.rs +++ b/tools/ironclaw_stress/src/synthetic.rs @@ -82,10 +82,15 @@ impl SyntheticIds { worker_index: usize, operation_index: usize, ) -> Result { - let (_, user_index, global_index) = - self.synthetic_indexes(args, worker_index, operation_index); - let thread_user_index = self.thread_user_index(args, user_index, global_index); - self.user_turn_context_for_indexes(user_index, thread_user_index) + let actor_user_index = partitioned_worker_index( + self.users.len(), + args.concurrency, + worker_index, + operation_index, + ); + let thread_user_index = + self.thread_user_index(args, actor_user_index, worker_index, operation_index); + self.user_turn_context_for_indexes(actor_user_index, thread_user_index) } pub(crate) fn user_turn_context_for_user_index( @@ -141,11 +146,22 @@ impl SyntheticIds { }) } - fn thread_user_index(&self, args: &Args, user_index: usize, global_index: usize) -> usize { + fn thread_user_index( + &self, + args: &Args, + user_index: usize, + worker_index: usize, + operation_index: usize, + ) -> usize { if args.active_thread_count == 0 { user_index } else { - global_index % args.active_thread_count + partitioned_worker_index( + args.active_thread_count, + args.concurrency, + worker_index, + operation_index, + ) } } @@ -163,3 +179,23 @@ impl SyntheticIds { (tenant_index, user_index, global_index) } } + +fn partitioned_worker_index( + pool_size: usize, + worker_count: usize, + worker_index: usize, + operation_index: usize, +) -> usize { + let worker_count = worker_count.max(1); + if pool_size < worker_count { + return worker_index % pool_size; + } + + let base_len = pool_size / worker_count; + let remainder = pool_size % worker_count; + let partition_len = base_len + usize::from(worker_index < remainder); + let partition_start = worker_index + .saturating_mul(base_len) + .saturating_add(worker_index.min(remainder)); + partition_start + operation_index % partition_len +} diff --git a/tools/ironclaw_stress/src/tests.rs b/tools/ironclaw_stress/src/tests.rs index 9bec14a6e05..9ea9c9f66ff 100644 --- a/tools/ironclaw_stress/src/tests.rs +++ b/tools/ironclaw_stress/src/tests.rs @@ -86,6 +86,56 @@ fn fixed_operation_mode_spreads_concurrent_workers_across_threads() { assert_eq!(thread_ids.len(), args.concurrency); } +#[test] +fn user_turn_workers_keep_disjoint_thread_partitions() { + let mut args = test_args(); + args.users = 10; + args.concurrency = 4; + args.operations = 100; + args.active_thread_count = 0; + let ids = SyntheticIds::new(&args).expect("synthetic ids build"); + + let mut seen_by_worker = Vec::new(); + for worker_index in 0..args.concurrency { + let thread_ids = (0..args.operations) + .map(|operation_index| { + ids.user_turn_context(&args, worker_index, operation_index) + .expect("context") + .thread_id + }) + .collect::>(); + seen_by_worker.push(thread_ids); + } + + for worker_index in 0..seen_by_worker.len() { + for other_worker_index in worker_index + 1..seen_by_worker.len() { + assert!( + seen_by_worker[worker_index].is_disjoint(&seen_by_worker[other_worker_index]), + "workers {worker_index} and {other_worker_index} should not share target threads" + ); + } + } + let all_threads = seen_by_worker + .iter() + .flat_map(|thread_ids| thread_ids.iter().cloned()) + .collect::>(); + assert_eq!(all_threads.len(), args.users); +} + +#[test] +fn user_turn_default_rejects_more_workers_than_users() { + let mut args = test_args(); + args.scenario = Scenario::ChatTurn; + args.users = 2; + args.concurrency = 3; + args.active_thread_count = 0; + + let error = + validate_args(&args).expect_err("default user-turn mode should require enough users"); + + assert!(error.contains("--users to be greater than or equal to --concurrency")); +} + #[test] fn chat_turn_rejects_multi_process_runs() { let mut args = test_args(); @@ -157,7 +207,7 @@ fn bottleneck_finder_suite_includes_core_pressure_cases() { assert!(labels.contains("resource-contention")); assert!(labels.contains("chat-baseline")); - assert!(labels.contains("hot-thread")); + assert!(!labels.contains("hot-thread")); assert!(labels.contains("large-context")); assert!(labels.contains("tool-heavy")); assert!(labels.contains("tool-wait")); @@ -176,7 +226,7 @@ fn postgres_pool_pressure_suite_includes_remote_pool_cases() { .collect::>(); assert!(labels.contains("postgres-chat-pool")); - assert!(labels.contains("postgres-hot-thread-pool")); + assert!(!labels.contains("postgres-hot-thread-pool")); assert!(labels.contains("postgres-context-pool")); assert!(labels.contains("postgres-tool-pool")); } @@ -381,6 +431,18 @@ fn sweep_concurrency_rejects_zero_values() { assert!(error.contains("--sweep-concurrency values must be greater than 0")); } +#[test] +fn sweep_concurrency_validation_uses_sweep_max_not_base_concurrency() { + let mut args = test_args(); + args.scenario = Scenario::MixedUserSession; + args.concurrency = 8; + args.sweep_concurrency = vec![2, 4]; + args.users = 4; + args.active_thread_count = 0; + + validate_args(&args).expect("sweep cases only require enough users for the sweep max"); +} + #[test] fn active_thread_count_rejects_values_above_users() { let mut args = test_args();