diff --git a/scripts/core-boundaries/explicit-test-topology.mjs b/scripts/core-boundaries/explicit-test-topology.mjs index e3c93aa5ec..307fe8fc84 100644 --- a/scripts/core-boundaries/explicit-test-topology.mjs +++ b/scripts/core-boundaries/explicit-test-topology.mjs @@ -41,6 +41,7 @@ export const servicesCoreIntegrationTestTargets = [ { name: 'permission_store_contracts', path: 'tests/permission_store_contracts.rs' }, { name: 'workspace_instruction_contracts', path: 'tests/workspace_instruction_contracts.rs' }, { name: 'session_write_lock_contracts', path: 'tests/session_write_lock_contracts.rs' }, + { name: 'exclusive_file_lease_contracts', path: 'tests/exclusive_file_lease_contracts.rs' }, { name: 'process_runtime_contracts', path: 'tests/process_runtime_contracts.rs' }, { name: 'service_contracts', path: 'tests/service_contracts.rs' }, { name: 'storage_owner_contracts', path: 'tests/storage_owner_contracts.rs' }, diff --git a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs index f8da52705b..36a49cdf09 100644 --- a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs +++ b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs @@ -71,6 +71,15 @@ pub use openbitfun_runtime_ports::{ DialogSubmitOutcome, }; +/// Rejection prefix for a submission that reuses a dialog turn ID the session +/// already owns. +/// +/// The scheduled-job service classifies enqueue failures from the port message +/// text, so producers and classifiers share this constant instead of repeating +/// the wording. +pub(crate) const DIALOG_TURN_ID_ALREADY_SETTLED_MESSAGE: &str = + "Dialog turn ID is already active or completed"; + /// A message waiting to be dispatched to the coordinator #[derive(Debug, Clone)] pub struct QueuedTurn { @@ -2772,7 +2781,7 @@ impl DialogScheduler { PortError::new( PortErrorKind::InvalidRequest, format!( - "Dialog turn ID is already active or completed: session_id={}, turn_id={resolved_turn_id}", + "{DIALOG_TURN_ID_ALREADY_SETTLED_MESSAGE}: session_id={}, turn_id={resolved_turn_id}", request.session_id ), ) diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/cron_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/cron_tool.rs index 8de27be2bb..e77ecb4c14 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/cron_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/cron_tool.rs @@ -1127,6 +1127,14 @@ Patch schema for "update": session_id )); } + if !cron_service.is_scheduling_owner() { + // Two instances can share one user data directory. Saying so + // here stops the agent from presenting a standby list as the + // state that will actually run. + result_for_assistant.push_str( + "\n\nScheduling note: another running OpenBitFun instance owns scheduled job execution right now, so these jobs run there and edits or manual runs are refused in this instance.", + ); + } Ok(vec![ToolResult::Result { data: json!({ diff --git a/src/crates/assembly/core/src/infrastructure/app_paths/path_manager.rs b/src/crates/assembly/core/src/infrastructure/app_paths/path_manager.rs index ea05e40d82..5dd023d391 100644 --- a/src/crates/assembly/core/src/infrastructure/app_paths/path_manager.rs +++ b/src/crates/assembly/core/src/infrastructure/app_paths/path_manager.rs @@ -320,6 +320,15 @@ impl PathManager { self.user_cron_dir().join("jobs.json") } + /// Lease file identifying the process that owns scheduled job scheduling. + /// + /// It sits beside `jobs.json` because it guards exactly that store: the + /// holder is the only process allowed to schedule jobs and write the file. + /// The file itself is inert — ownership is the OS lock on it. + pub fn cron_scheduler_lease_file(&self) -> PathBuf { + self.user_cron_dir().join("scheduler.lock") + } + /// Get miniapps root directory: ~/.config/openbitfun/data/miniapps/ pub fn miniapps_dir(&self) -> PathBuf { self.user_data_dir().join("miniapps") diff --git a/src/crates/assembly/core/src/service/cron/service.rs b/src/crates/assembly/core/src/service/cron/service.rs index 6f43af415c..4168316f46 100644 --- a/src/crates/assembly/core/src/service/cron/service.rs +++ b/src/crates/assembly/core/src/service/cron/service.rs @@ -5,10 +5,11 @@ use super::schedule::{ }; use super::store::CronJobStore; use super::types::{ - CreateCronJobRequest, CronJob, CronJobPayload, CronJobTarget, CronJobTargetKind, + CreateCronJobRequest, CronJob, CronJobPayload, CronJobTarget, CronJobTargetKind, CronJobsFile, CronLaunchSpec, CronSchedule, CronWorkspaceRef, UpdateCronJobRequest, DEFAULT_RETRY_DELAY_MS, }; use super::{CronJobsChangedEvent, CronJobsChangedReason, CRON_JOBS_CHANGED_EVENT}; +use crate::agentic::coordination::scheduler::DIALOG_TURN_ID_ALREADY_SETTLED_MESSAGE; use crate::agentic::coordination::{ ConversationCoordinator, DialogQueuePriority, DialogScheduler, DialogSubmissionPolicy, DialogTriggerSource, @@ -20,11 +21,13 @@ use crate::infrastructure::PathManager; use crate::service_agent_runtime::CoreServiceAgentRuntime; use crate::util::errors::{OpenBitFunError, OpenBitFunResult}; use chrono::{Local, SecondsFormat, TimeZone, Utc}; -use log::{debug, info, warn}; +use log::{debug, error, info, warn}; use openbitfun_agent_runtime::scheduled_job::ScheduledJobEnqueueFailureAction; use openbitfun_agent_runtime::sdk::AgentRuntime; use openbitfun_runtime_ports::{AgentDialogPrependedReminder, AgentDialogTurnRequest}; +use openbitfun_services_core::exclusive_file_lease::{ExclusiveFileLease, ExclusiveFileLeaseError}; use std::collections::HashMap; +use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, OnceLock}; use tokio::sync::{Mutex, Notify, RwLock}; @@ -33,6 +36,11 @@ use uuid::Uuid; static GLOBAL_CRON_SERVICE: OnceLock> = OnceLock::new(); +/// How long a standby instance waits before re-checking whether the scheduling +/// owner released the lease. Only a process exit releases it, so this is a +/// takeover latency bound, not a lock timeout. +const STANDBY_TAKEOVER_RETRY: Duration = Duration::from_secs(30); + pub struct CronService { coordinator: Arc, runtime: AgentRuntime, @@ -41,6 +49,12 @@ pub struct CronService { mutation_lock: Arc>, wakeup: Arc, runner_started: AtomicBool, + /// Held by the process that schedules jobs and writes `jobs.json`. + /// + /// Two desktop instances share one user data directory, so without this + /// lease both schedulers fire the same trigger, and the loser retries a + /// turn id the winner already consumed. + scheduling_lease: SchedulingLease, } impl CronService { @@ -49,40 +63,16 @@ impl CronService { coordinator: Arc, scheduler: Arc, ) -> OpenBitFunResult> { - let store = Arc::new(CronJobStore::new(path_manager).await?); - let loaded = store.load().await?; - let current_ms = now_ms(); - - let mut jobs = HashMap::new(); - let mut needs_save = false; - - for mut job in loaded.jobs { - if jobs.contains_key(&job.id) { - return Err(OpenBitFunError::service(format!( - "Duplicate scheduled job id found in jobs.json: {}", - job.id - ))); - } - - // Upgrade persisted pre-ID targets once. Unavailable/ambiguous records - // stay on disk and fail explicitly when executed, never select a folder. - let old_target = job.target.clone(); - match crate::service::workspace::legacy_compat::upgrade_legacy_cron_target( - job.target.clone(), - ) - .await - { - Ok(target) => job.target = target, - Err(error) => warn!( - "Unable to upgrade scheduled job workspace: job_id={}, error={}", - job.id, error - ), - } - needs_save |= job.target != old_target; - needs_save |= reconcile_loaded_job(&mut job, current_ms)?; - jobs.insert(job.id.clone(), job); + // Take the scheduling lease before touching the store: it is what makes + // this instance the only writer of jobs.json. + let scheduling_lease = SchedulingLease::new(path_manager.cron_scheduler_lease_file()); + if scheduling_lease.claim() { + scheduling_lease.publish_ownership(); } + let store = Arc::new(CronJobStore::new(path_manager).await?); + let (jobs, needs_save) = materialize_loaded_jobs(store.load().await?).await?; + let runtime = CoreServiceAgentRuntime::agent_runtime_with_dialog_turns( coordinator.clone(), scheduler, @@ -97,15 +87,39 @@ impl CronService { mutation_lock: Arc::new(Mutex::new(())), wakeup: Arc::new(Notify::new()), runner_started: AtomicBool::new(false), + scheduling_lease, }); + if !service.is_scheduling_owner() { + warn!( + "Another OpenBitFun instance owns scheduled jobs; this instance runs as standby and only reads them: lease={}", + service.scheduling_lease.path().display() + ); + } + if needs_save { - service.persist_snapshot().await?; + if service.is_scheduling_owner() { + service.persist_snapshot().await?; + } else { + // Writing here would clobber state the owner is advancing. + warn!( + "Deferring scheduled job store upgrade to the owning instance: lease={}", + service.scheduling_lease.path().display() + ); + } } Ok(service) } + /// Whether this instance schedules jobs and owns `jobs.json`. + /// + /// Standby instances reject job changes and read the store from disk so + /// they never report state the owner has already replaced. + pub fn is_scheduling_owner(&self) -> bool { + self.scheduling_lease.is_owned() + } + pub fn start(self: &Arc) { if self .runner_started @@ -116,12 +130,91 @@ impl CronService { } let service = Arc::clone(self); + if service.is_scheduling_owner() { + tokio::spawn(async move { + service.run_loop().await; + }); + return; + } tokio::spawn(async move { - service.run_loop().await; + service.run_standby_loop().await; }); } + fn require_scheduling_owner(&self, operation: &str) -> OpenBitFunResult<()> { + self.scheduling_lease.require_owned(operation) + } + + /// Waits for the owner's lease so this instance can take over scheduling. + async fn run_standby_loop(self: Arc) { + loop { + tokio::time::sleep(STANDBY_TAKEOVER_RETRY).await; + match self.try_take_over_scheduling().await { + Ok(true) => { + self.run_loop().await; + return; + } + Ok(false) => {} + Err(error) => { + warn!("Failed to take over scheduled job scheduling: {}", error); + } + } + } + } + + async fn try_take_over_scheduling(&self) -> OpenBitFunResult { + if !self.scheduling_lease.claim() { + return Ok(false); + } + if self.is_scheduling_owner() { + return Ok(true); + } + + // The previous owner advanced job state while this instance was standby, + // so adopt the store before scheduling from it. Ownership stays + // unpublished until the adoption is in memory, and standby callers are + // refused during that window instead of writing a stale map. + let (jobs, needs_save) = self.load_jobs_from_store().await?; + *self.jobs.write().await = jobs; + self.scheduling_lease.publish_ownership(); + if needs_save { + self.persist_snapshot().await?; + } + self.notify_jobs_changed(CronJobsChangedReason::StateChanged, None) + .await; + info!( + "Scheduled job scheduler took ownership from a released lease: lease={}", + self.scheduling_lease.path().display() + ); + Ok(true) + } + + async fn load_jobs_from_store(&self) -> OpenBitFunResult<(HashMap, bool)> { + materialize_loaded_jobs(self.store.load().await?).await + } + + /// Standby readers see the owner's current state instead of their own + /// snapshot from startup. Owner reads stay on the in-memory map, which is + /// authoritative while it holds the lease. + async fn refresh_from_store_when_standby(&self) { + if self.is_scheduling_owner() { + return; + } + match self.load_jobs_from_store().await { + Ok((jobs, _)) => { + *self.jobs.write().await = jobs; + } + Err(error) => { + warn!( + "Failed to read scheduled jobs from the owning instance's store: {}", + error + ); + } + } + } + pub async fn list_jobs(&self) -> Vec { + self.refresh_from_store_when_standby().await; let jobs = self.jobs.read().await; jobs.values().cloned().collect::>() } @@ -155,6 +248,7 @@ impl CronService { session_id: Option<&str>, target_kind: Option, ) -> Vec { + self.refresh_from_store_when_standby().await; let jobs = self.jobs.read().await; jobs.values() .filter(|job| { @@ -172,10 +266,12 @@ impl CronService { } pub async fn get_job(&self, job_id: &str) -> Option { + self.refresh_from_store_when_standby().await; self.jobs.read().await.get(job_id).cloned() } pub async fn create_job(&self, request: CreateCronJobRequest) -> OpenBitFunResult { + self.require_scheduling_owner("creating a scheduled job")?; let target = self.canonicalize_target(request.target).await?; let _guard = self.mutation_lock.lock().await; let mut jobs = self.jobs.write().await; @@ -218,6 +314,7 @@ impl CronService { job_id: &str, request: UpdateCronJobRequest, ) -> OpenBitFunResult { + self.require_scheduling_owner("updating a scheduled job")?; let canonicalized_target = match request.target { Some(target) => Some(self.canonicalize_target(target).await?), None => None, @@ -286,6 +383,7 @@ impl CronService { } pub async fn delete_job(&self, job_id: &str) -> OpenBitFunResult { + self.require_scheduling_owner("deleting a scheduled job")?; let _guard = self.mutation_lock.lock().await; let mut jobs = self.jobs.write().await; let existed = jobs.remove(job_id).is_some(); @@ -301,6 +399,7 @@ impl CronService { /// Remove all scheduled jobs bound to the given session (e.g. after session delete). pub async fn delete_jobs_for_session(&self, session_id: &str) -> OpenBitFunResult { + self.require_scheduling_owner("deleting scheduled jobs of a session")?; let session_id = session_id.trim(); if session_id.is_empty() { return Ok(0); @@ -321,6 +420,7 @@ impl CronService { } pub async fn run_job_now(&self, job_id: &str) -> OpenBitFunResult { + self.require_scheduling_owner("running a scheduled job now")?; { let _guard = self.mutation_lock.lock().await; let mut jobs = self.jobs.write().await; @@ -388,6 +488,15 @@ impl CronService { where F: FnOnce(&mut CronJob, i64), { + if !self.is_scheduling_owner() { + // Turn lifecycle events reach every instance, but only the owner + // tracks job state; writing here would race the owner's store. + debug!( + "Ignoring scheduled job turn state change while standby: turn_id={}", + turn_id + ); + return Ok(()); + } let _guard = self.mutation_lock.lock().await; let mut jobs = self.jobs.write().await; let Some(job_id) = jobs @@ -569,6 +678,22 @@ impl CronService { scheduled_at_ms ); } + Err(error) if cron_enqueue_error_is_turn_already_settled(&error) => { + job.state.mark_trigger_delivered_elsewhere(now_after_submit); + job.updated_at_ms = now_after_submit; + + if job.is_one_shot() { + job.enabled = false; + } + + info!( + "Scheduled job trigger was already delivered by another owner: job_id={}, target_kind={:?}, target_session_id={}, scheduled_at_ms={}", + job.id, + job.target_kind(), + submit_target_session_id(&enqueue_input), + scheduled_at_ms + ); + } Err(error) => { let missing_session = matches!(job.target_kind(), CronJobTargetKind::Session) && cron_enqueue_error_is_missing_session(&error); @@ -769,6 +894,123 @@ pub fn set_global_cron_service(service: Arc) { let _ = GLOBAL_CRON_SERVICE.set(service); } +/// Process-local view of the cross-process lease that decides which instance +/// schedules jobs and writes `jobs.json`. +/// +/// The lease is held until the process exits; dropping it releases the OS lock +/// so a standby instance can take over. +struct SchedulingLease { + path: PathBuf, + lease: std::sync::Mutex>, + owned: AtomicBool, +} + +impl SchedulingLease { + fn new(path: PathBuf) -> Self { + Self { + path, + lease: std::sync::Mutex::new(None), + owned: AtomicBool::new(false), + } + } + + fn path(&self) -> &Path { + &self.path + } + + fn is_owned(&self) -> bool { + self.owned.load(Ordering::SeqCst) + } + + /// Takes the cross-process lease without publishing ownership, so the + /// caller can adopt the store it protects first. `false` means another + /// process holds it and this instance stays standby. + /// + /// A lease that cannot be created at all fails open: scheduled jobs keep + /// working exactly as they did before the lease existed, instead of a + /// product feature going silently dead because one lock file is + /// unavailable. That case stays loud in the log and leaves the + /// multi-instance hazard in place. + fn claim(&self) -> bool { + if self.is_owned() { + return true; + } + match ExclusiveFileLease::try_acquire(&self.path) { + Ok(lease) => { + *self + .lease + .lock() + .expect("Scheduled job lease slot poisoned") = Some(lease); + true + } + Err(ExclusiveFileLeaseError::InUse) => false, + Err(error) => { + error!( + "Failed to acquire the scheduled job lease {}; scheduling without cross-instance protection: {}", + self.path.display(), + error + ); + true + } + } + } + + fn publish_ownership(&self) { + self.owned.store(true, Ordering::SeqCst); + } + + fn require_owned(&self, operation: &str) -> OpenBitFunResult<()> { + if self.is_owned() { + return Ok(()); + } + Err(OpenBitFunError::service(format!( + "Scheduled jobs are managed by the other running OpenBitFun instance (lease: {}); {} was refused here instead of overwriting its state", + self.path.display(), + operation + ))) + } +} + +/// Turns the persisted file into the in-memory job map, applying the one-time +/// target upgrade and the startup reconciliation. The returned flag reports +/// whether the result differs from what is on disk. +async fn materialize_loaded_jobs( + loaded: CronJobsFile, +) -> OpenBitFunResult<(HashMap, bool)> { + let current_ms = now_ms(); + let mut jobs = HashMap::new(); + let mut needs_save = false; + + for mut job in loaded.jobs { + if jobs.contains_key(&job.id) { + return Err(OpenBitFunError::service(format!( + "Duplicate scheduled job id found in jobs.json: {}", + job.id + ))); + } + + // Upgrade persisted pre-ID targets once. Unavailable/ambiguous records + // stay on disk and fail explicitly when executed, never select a folder. + let old_target = job.target.clone(); + match crate::service::workspace::legacy_compat::upgrade_legacy_cron_target( + job.target.clone(), + ) + .await + { + Ok(target) => job.target = target, + Err(error) => warn!( + "Unable to upgrade scheduled job workspace: job_id={}, error={}", + job.id, error + ), + } + needs_save |= job.target != old_target; + needs_save |= reconcile_loaded_job(&mut job, current_ms)?; + jobs.insert(job.id.clone(), job); + } + + Ok((jobs, needs_save)) +} + fn reconcile_loaded_job(job: &mut CronJob, now_ms: i64) -> OpenBitFunResult { let original = job.clone(); @@ -1084,6 +1326,15 @@ fn submit_target_session_id(enqueue_input: &EnqueueInput) -> &str { } } +/// The trigger's dialog turn already exists, so this trigger was delivered by +/// another owner instead of failing. +/// +/// The trigger's turn ID is derived from its scheduled timestamp, so every retry +/// for the same trigger is rejected the same way; retrying can never succeed. +fn cron_enqueue_error_is_turn_already_settled(error: &str) -> bool { + error.contains(DIALOG_TURN_ID_ALREADY_SETTLED_MESSAGE) +} + /// Permanent failure: coordinator cannot load session metadata (session deleted from disk). fn cron_enqueue_error_is_missing_session(error: &str) -> bool { error.contains("Session metadata not found") @@ -1092,7 +1343,7 @@ fn cron_enqueue_error_is_missing_session(error: &str) -> bool { #[cfg(test)] mod tests { use super::*; - use crate::service::cron::CronJobState; + use crate::service::cron::{CronJobState, CRON_JOBS_VERSION}; fn sample_job(schedule: CronSchedule) -> CronJob { CronJob { @@ -1121,6 +1372,128 @@ mod tests { } } + #[test] + fn enqueue_failure_classifies_already_delivered_trigger() { + assert!(cron_enqueue_error_is_turn_already_settled(&format!( + "{}: session_id=session_1, turn_id=cronjob_cron_test_1000", + DIALOG_TURN_ID_ALREADY_SETTLED_MESSAGE + ))); + assert!(!cron_enqueue_error_is_turn_already_settled( + "Session is already open for writing: session_1" + )); + assert!(!cron_enqueue_error_is_turn_already_settled( + "Session metadata not found" + )); + } + + #[tokio::test] + async fn loaded_jobs_are_rewritten_when_they_gain_their_first_anchor() { + let mut job = sample_job(CronSchedule::Every { + every_ms: 60_000, + anchor_ms: None, + }); + // An id-carrying workspace needs no legacy upgrade, so the pending save + // can only come from the anchor reconciliation. + if let CronJobTarget::Session { workspace, .. } = &mut job.target { + workspace.workspace_id = Some("workspace_1".to_string()); + } + + let (jobs, needs_save) = materialize_loaded_jobs(CronJobsFile { + version: CRON_JOBS_VERSION, + jobs: vec![job], + }) + .await + .expect("load"); + + assert_eq!(jobs.len(), 1); + assert!( + needs_save, + "a job that gains its first schedule anchor must be written back" + ); + } + + #[tokio::test] + async fn loaded_jobs_reject_duplicate_ids_before_scheduling() { + let mut job = sample_job(CronSchedule::Every { + every_ms: 60_000, + anchor_ms: Some(0), + }); + if let CronJobTarget::Session { workspace, .. } = &mut job.target { + workspace.workspace_id = Some("workspace_1".to_string()); + } + + let error = materialize_loaded_jobs(CronJobsFile { + version: CRON_JOBS_VERSION, + jobs: vec![job.clone(), job], + }) + .await + .expect_err("duplicate ids must be rejected"); + + assert!(error.to_string().contains("Duplicate scheduled job id")); + } + + #[test] + fn a_second_lease_on_the_same_path_stays_standby_and_refuses_writes() { + let project_root = tempfile::tempdir().expect("project root"); + let lease_path = project_root.path().join("cron").join("scheduler.lock"); + + let owner = SchedulingLease::new(lease_path.clone()); + assert!(owner.claim(), "first claim"); + owner.publish_ownership(); + assert!(owner.is_owned()); + + let standby = SchedulingLease::new(lease_path.clone()); + assert!(!standby.claim(), "the lease is held elsewhere"); + assert!(!standby.is_owned()); + + let error = standby + .require_owned("updating a scheduled job") + .expect_err("standby must refuse job changes"); + assert!(error.to_string().contains("managed by the other running")); + assert!(error.to_string().contains("scheduler.lock")); + + // The leaseless owner stays able to act on its own store. + owner + .require_owned("updating a scheduled job") + .expect("owner may write"); + } + + #[test] + fn dropping_the_owner_releases_the_lease_for_takeover() { + let project_root = tempfile::tempdir().expect("project root"); + let lease_path = project_root.path().join("cron").join("scheduler.lock"); + + let owner = SchedulingLease::new(lease_path.clone()); + assert!(owner.claim(), "first claim"); + owner.publish_ownership(); + drop(owner); + + let standby = SchedulingLease::new(lease_path); + assert!( + standby.claim(), + "a released lease must be claimable for takeover" + ); + assert!(!standby.is_owned(), "ownership needs an explicit publish"); + standby.publish_ownership(); + assert!(standby.is_owned()); + } + + #[test] + fn an_unusable_lease_path_fails_open_instead_of_disabling_scheduling() { + let project_root = tempfile::tempdir().expect("project root"); + // A regular file where the lease directory should be makes the lease + // impossible to create, which is not the same as "someone else holds it". + let blocking_file = project_root.path().join("cron"); + std::fs::write(&blocking_file, b"not a directory").expect("blocking file"); + + let lease = SchedulingLease::new(blocking_file.join("scheduler.lock")); + + assert!( + lease.claim(), + "an unusable lease must keep the pre-lease scheduling behavior" + ); + } + #[test] fn generate_cron_job_id_uses_short_hex_suffix() { let jobs = HashMap::new(); diff --git a/src/crates/execution/agent-runtime/src/scheduled_job.rs b/src/crates/execution/agent-runtime/src/scheduled_job.rs index b9f8c90e36..2b8c62424c 100644 --- a/src/crates/execution/agent-runtime/src/scheduled_job.rs +++ b/src/crates/execution/agent-runtime/src/scheduled_job.rs @@ -150,6 +150,23 @@ impl ScheduledJobRuntimeState { } } + /// The trigger's dialog turn already exists, so the trigger was delivered by + /// another owner: a second application instance racing the same trigger, or + /// this instance before a restart that lost the enqueue result. + /// + /// A trigger's turn ID is derived from its scheduled timestamp, so every + /// later delivery attempt for that trigger is rejected the same way. + /// Retrying can never succeed, so record the trigger as delivered and stop + /// counting failures instead of retrying until the process exits. + pub fn mark_trigger_delivered_elsewhere(&mut self, delivered_at_ms: i64) { + self.pending_trigger_at_ms = None; + self.retry_at_ms = None; + self.last_run_status = Some(ScheduledJobRunStatus::Ok); + self.last_error = None; + self.last_run_finished_at_ms = Some(delivered_at_ms); + self.consecutive_failures = 0; + } + pub fn mark_turn_started(&mut self, started_at_ms: i64) { self.last_run_status = Some(ScheduledJobRunStatus::Running); self.last_run_started_at_ms = Some(started_at_ms); diff --git a/src/crates/execution/agent-runtime/tests/agent_session_contracts/scheduled_job_contracts.rs b/src/crates/execution/agent-runtime/tests/agent_session_contracts/scheduled_job_contracts.rs index 8c796365ff..d32cb2fa98 100644 --- a/src/crates/execution/agent-runtime/tests/agent_session_contracts/scheduled_job_contracts.rs +++ b/src/crates/execution/agent-runtime/tests/agent_session_contracts/scheduled_job_contracts.rs @@ -195,6 +195,34 @@ fn enqueue_failure_preserves_retry_and_missing_session_disable_semantics() { assert_eq!(missing_session_state.consecutive_failures, 3); } +#[test] +fn trigger_delivered_elsewhere_clears_pending_and_stops_retrying() { + let mut state = ScheduledJobRuntimeState { + next_run_at_ms: Some(300), + pending_trigger_at_ms: Some(100), + retry_at_ms: Some(200), + last_run_status: Some(ScheduledJobRunStatus::Error), + last_error: Some( + "Dialog turn ID is already active or completed: session_id=session_1, turn_id=turn-1" + .to_string(), + ), + consecutive_failures: 416, + ..Default::default() + }; + + state.mark_trigger_delivered_elsewhere(250); + + assert_eq!(state.pending_trigger_at_ms, None); + assert_eq!(state.retry_at_ms, None); + assert_eq!(state.next_run_at_ms, Some(300)); + assert_eq!(state.last_run_status, Some(ScheduledJobRunStatus::Ok)); + assert_eq!(state.last_error, None); + assert_eq!(state.last_run_finished_at_ms, Some(250)); + assert_eq!(state.consecutive_failures, 0); + assert!(!state.pending_is_due(1_000)); + assert_eq!(state.next_wakeup_at_ms(), Some(300)); +} + #[test] fn turn_completion_failure_and_cancel_preserve_legacy_status_fields() { let mut completed = ScheduledJobRuntimeState { diff --git a/src/crates/services/services-core/AGENTS.md b/src/crates/services/services-core/AGENTS.md index 57cf93286a..13ebb03fdf 100644 --- a/src/crates/services/services-core/AGENTS.md +++ b/src/crates/services/services-core/AGENTS.md @@ -96,6 +96,7 @@ cargo test -p openbitfun-services-core --no-default-features --features workspac cargo test -p openbitfun-services-core --no-default-features --features local-storage --test session_contracts session_metadata_contracts:: cargo test -p openbitfun-services-core --no-default-features --features local-storage --lib session::metadata cargo test -p openbitfun-services-core --no-default-features --features local-storage --test session_write_lock_contracts +cargo test -p openbitfun-services-core --no-default-features --features local-storage --test exclusive_file_lease_contracts cargo test -p openbitfun-services-core --no-default-features --features memory-store --lib memory_store::tests:: cargo test -p openbitfun-services-core --no-default-features --features token-usage-statistics --lib token_usage:: cargo test -p openbitfun-services-core --no-default-features --features process-runtime --test process_runtime_contracts diff --git a/src/crates/services/services-core/Cargo.toml b/src/crates/services/services-core/Cargo.toml index a30c1ae867..7dbf14168c 100644 --- a/src/crates/services/services-core/Cargo.toml +++ b/src/crates/services/services-core/Cargo.toml @@ -227,6 +227,11 @@ name = "session_write_lock_contracts" path = "tests/session_write_lock_contracts.rs" required-features = ["local-storage"] +[[test]] +name = "exclusive_file_lease_contracts" +path = "tests/exclusive_file_lease_contracts.rs" +required-features = ["local-storage"] + [[test]] name = "process_runtime_contracts" path = "tests/process_runtime_contracts.rs" diff --git a/src/crates/services/services-core/src/exclusive_file_lease.rs b/src/crates/services/services-core/src/exclusive_file_lease.rs new file mode 100644 index 0000000000..087da4924a --- /dev/null +++ b/src/crates/services/services-core/src/exclusive_file_lease.rs @@ -0,0 +1,138 @@ +use crate::file_lock::{is_lock_contention, FileLock, FileLockError, FileLockMode}; +use std::collections::HashMap; +use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex, OnceLock, Weak}; + +/// Holds one shared resource for a single process at a time. +/// +/// Where `session::write_lock` guards one Session, this guards a resource that +/// an owning process keeps for its whole lifetime — a scheduler loop, a store +/// writer. Ownership is represented by the OS lock, not by file existence: +/// the lock file may remain on disk after the lease is released, and a stale +/// file never implies ownership. +pub struct ExclusiveFileLease { + inner: Arc, +} + +struct ExclusiveFileLeaseInner { + lock: Option, + lock_path: PathBuf, +} + +impl ExclusiveFileLease { + /// Acquires the lease without waiting so callers can degrade to a standby + /// state instead of blocking startup on another process. + /// + /// The parent directory is created when missing, so the caller does not have + /// to pre-create the lease location. + pub fn try_acquire(lock_path: &Path) -> Result { + if let Some(parent) = lock_path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + { + std::fs::create_dir_all(parent).map_err(|source| { + ExclusiveFileLeaseError::CreateLeaseDirectory { + path: parent.to_path_buf(), + source, + } + })?; + } + let mut process_leases = process_leases() + .lock() + .expect("Exclusive file lease registry poisoned"); + if process_leases + .get(lock_path) + .and_then(Weak::upgrade) + .is_some() + { + return Err(ExclusiveFileLeaseError::InUse); + } + let lock = + FileLock::try_acquire(lock_path, FileLockMode::Exclusive).map_err( + |error| match error { + FileLockError::Open(source) => ExclusiveFileLeaseError::OpenLeaseFile { + path: lock_path.to_path_buf(), + source, + }, + FileLockError::Unavailable(source) if is_lock_contention(&source) => { + ExclusiveFileLeaseError::InUse + } + FileLockError::Unavailable(source) => { + ExclusiveFileLeaseError::LockFailed { source } + } + }, + )?; + let inner = Arc::new(ExclusiveFileLeaseInner { + lock: Some(lock), + lock_path: lock_path.to_path_buf(), + }); + process_leases.insert(lock_path.to_path_buf(), Arc::downgrade(&inner)); + Ok(Self { inner }) + } + + /// The path whose OS lock represents ownership. Useful for diagnostics. + pub fn lock_path(&self) -> &Path { + &self.inner.lock_path + } +} + +impl Drop for ExclusiveFileLeaseInner { + fn drop(&mut self) { + if let Ok(mut process_leases) = process_leases().lock() { + // Keep the in-process registry authoritative until the OS lock is + // released so an immediate reacquire cannot observe a false gap. + drop(self.lock.take()); + process_leases.remove(&self.lock_path); + } + } +} + +impl std::fmt::Debug for ExclusiveFileLease { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("ExclusiveFileLease") + .field("lock_path", &self.inner.lock_path) + .field("strong_count", &Arc::strong_count(&self.inner)) + .finish() + } +} + +fn process_leases() -> &'static Mutex>> { + static PROCESS_LEASES: OnceLock>>> = + OnceLock::new(); + PROCESS_LEASES.get_or_init(|| Mutex::new(HashMap::new())) +} + +#[derive(Debug, thiserror::Error)] +pub enum ExclusiveFileLeaseError { + #[error("failed to create lease directory {path}")] + CreateLeaseDirectory { + path: PathBuf, + #[source] + source: std::io::Error, + }, + #[error("failed to open lease file {path}")] + OpenLeaseFile { + path: PathBuf, + #[source] + source: std::io::Error, + }, + #[error("resource is already leased by another owner")] + InUse, + #[error("failed to acquire lease")] + LockFailed { + #[source] + source: std::io::Error, + }, +} + +impl ExclusiveFileLeaseError { + pub fn code(&self) -> &'static str { + match self { + Self::CreateLeaseDirectory { .. } => "lease_directory_create_failed", + Self::OpenLeaseFile { .. } => "lease_open_failed", + Self::InUse => "lease_in_use", + Self::LockFailed { .. } => "lease_lock_failed", + } + } +} diff --git a/src/crates/services/services-core/src/file_lock.rs b/src/crates/services/services-core/src/file_lock.rs index 551dc7fc5f..914c3a0e0e 100644 --- a/src/crates/services/services-core/src/file_lock.rs +++ b/src/crates/services/services-core/src/file_lock.rs @@ -43,6 +43,27 @@ impl FileLock { } } +/// Whether a failed lock attempt means "someone else holds it" rather than a +/// real filesystem failure. Callers map this to their own `InUse` error. +/// +/// A contended `flock` reports `EAGAIN`/`EWOULDBLOCK`, which `std` already +/// surfaces as [`std::io::ErrorKind::WouldBlock`] on every supported platform, +/// so this stays free of platform crates. +pub(crate) fn is_lock_contention(error: &std::io::Error) -> bool { + if error.kind() == std::io::ErrorKind::WouldBlock { + return true; + } + #[cfg(windows)] + { + // ERROR_LOCK_VIOLATION: an exclusive lock is held by someone else. + error.raw_os_error() == Some(33) + } + #[cfg(not(windows))] + { + false + } +} + fn open_lock_file(path: &Path) -> Result { OpenOptions::new() .create(true) diff --git a/src/crates/services/services-core/src/lib.rs b/src/crates/services/services-core/src/lib.rs index 193aab3856..972e0b7032 100644 --- a/src/crates/services/services-core/src/lib.rs +++ b/src/crates/services/services-core/src/lib.rs @@ -16,6 +16,8 @@ pub mod dispatch_contract; #[cfg(feature = "dispatch-workspace")] pub mod dispatch_workspace; #[cfg(any(feature = "local-storage", feature = "runtime-ownership"))] +pub mod exclusive_file_lease; +#[cfg(any(feature = "local-storage", feature = "runtime-ownership"))] mod file_lock; #[cfg(any(feature = "filesystem", feature = "workspace-transfer"))] pub mod file_write_lock; diff --git a/src/crates/services/services-core/src/session/write_lock.rs b/src/crates/services/services-core/src/session/write_lock.rs index 4bfa04e024..022c173eb5 100644 --- a/src/crates/services/services-core/src/session/write_lock.rs +++ b/src/crates/services/services-core/src/session/write_lock.rs @@ -1,4 +1,4 @@ -use crate::file_lock::{FileLock, FileLockError, FileLockMode}; +use crate::file_lock::{is_lock_contention, FileLock, FileLockError, FileLockMode}; use sha2::{Digest, Sha256}; use std::collections::HashMap; use std::fmt::Write as _; @@ -89,7 +89,7 @@ impl SessionWriteLock { path: lock_path.clone(), source, }, - FileLockError::Unavailable(source) if is_contention(&source) => { + FileLockError::Unavailable(source) if is_lock_contention(&source) => { SessionWriteLockError::InUse } FileLockError::Unavailable(source) => { @@ -232,20 +232,6 @@ fn hash_path(hasher: &mut Sha256, path: &Path) { } } -fn is_contention(error: &std::io::Error) -> bool { - if error.kind() == std::io::ErrorKind::WouldBlock { - return true; - } - #[cfg(windows)] - { - error.raw_os_error() == Some(33) - } - #[cfg(unix)] - { - matches!(error.raw_os_error(), Some(libc::EAGAIN)) - } -} - #[cfg(test)] mod tests { use super::{session_lock_root, SessionWriteLockError}; diff --git a/src/crates/services/services-core/tests/exclusive_file_lease_contracts.rs b/src/crates/services/services-core/tests/exclusive_file_lease_contracts.rs new file mode 100644 index 0000000000..ea6dcf1cb9 --- /dev/null +++ b/src/crates/services/services-core/tests/exclusive_file_lease_contracts.rs @@ -0,0 +1,140 @@ +#![cfg(feature = "local-storage")] + +use openbitfun_services_core::exclusive_file_lease::{ExclusiveFileLease, ExclusiveFileLeaseError}; +use std::path::Path; +use std::process::{Command, Stdio}; +use std::time::{Duration, Instant}; +use tempfile::tempdir; + +#[test] +fn one_resource_can_have_only_one_lease_holder() { + let project_root = tempdir().expect("project runtime root"); + let lock_path = project_root.path().join("lease.lock"); + + let holder = ExclusiveFileLease::try_acquire(&lock_path).expect("first holder"); + let error = ExclusiveFileLease::try_acquire(&lock_path) + .expect_err("second holder must fail immediately"); + + assert!(matches!(error, ExclusiveFileLeaseError::InUse)); + assert_eq!(error.code(), "lease_in_use"); + assert_eq!(holder.lock_path(), lock_path.as_path()); + + drop(holder); + ExclusiveFileLease::try_acquire(&lock_path).expect("holder after release"); +} + +#[test] +fn different_lease_paths_are_independent() { + let project_root = tempdir().expect("project runtime root"); + + let _scheduler = ExclusiveFileLease::try_acquire(&project_root.path().join("scheduler.lock")) + .expect("scheduler lease"); + let _store = ExclusiveFileLease::try_acquire(&project_root.path().join("store.lock")) + .expect("store lease"); +} + +#[test] +fn a_stale_lock_file_does_not_block_a_new_holder() { + let project_root = tempdir().expect("project runtime root"); + let lock_path = project_root.path().join("lease.lock"); + + let holder = ExclusiveFileLease::try_acquire(&lock_path).expect("first holder"); + assert!( + lock_path.exists(), + "acquiring the lease creates its lock file" + ); + drop(holder); + + assert!( + lock_path.exists(), + "the lock file may remain after the OS lock is released" + ); + ExclusiveFileLease::try_acquire(&lock_path).expect("stale file must not imply ownership"); +} + +#[test] +fn acquiring_creates_a_missing_parent_directory() { + let project_root = tempdir().expect("project runtime root"); + let lock_path = project_root.path().join("nested").join("lease.lock"); + + let _holder = ExclusiveFileLease::try_acquire(&lock_path).expect("holder in a new directory"); + + assert!(lock_path.exists()); +} + +#[test] +fn a_path_alias_resolves_to_the_same_lease() { + let project_root = tempdir().expect("project runtime root"); + let lock_path = project_root.path().join("lease.lock"); + let alias = project_root.path().join(".").join("lease.lock"); + + let _holder = ExclusiveFileLease::try_acquire(&lock_path).expect("first holder"); + let error = ExclusiveFileLease::try_acquire(&alias) + .expect_err("a path alias must identify the same lease"); + + assert_eq!(error.code(), "lease_in_use"); +} + +#[test] +fn abnormal_process_exit_releases_the_lease() { + let project_root = tempdir().expect("project runtime root"); + let lock_path = project_root.path().join("lease.lock"); + let ready_path = project_root.path().join("child-ready"); + let mut command = Command::new(std::env::current_exe().expect("current test executable")); + #[cfg(windows)] + { + use std::os::windows::process::CommandExt; + command.creation_flags(0x0800_0000); + } + let mut child = command + .arg("--exact") + .arg("abnormal_exit_child_holds_lease") + .arg("--nocapture") + .env("OPENBITFUN_EXCLUSIVE_FILE_LEASE_CHILD", "1") + .env("OPENBITFUN_EXCLUSIVE_FILE_LEASE_LOCK_PATH", &lock_path) + .env("OPENBITFUN_EXCLUSIVE_FILE_LEASE_READY_PATH", &ready_path) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("spawn lease child"); + + let deadline = Instant::now() + Duration::from_secs(5); + while !ready_path.exists() && Instant::now() < deadline { + if let Some(status) = child.try_wait().expect("poll lease child") { + panic!("lease child exited before acquiring the lease: {status}"); + } + std::thread::sleep(Duration::from_millis(10)); + } + if !ready_path.exists() { + let _ = child.kill(); + let _ = child.wait(); + panic!("lease child did not become ready"); + } + let was_blocked = matches!( + ExclusiveFileLease::try_acquire(&lock_path), + Err(ExclusiveFileLeaseError::InUse) + ); + + child.kill().expect("terminate lease child"); + child.wait().expect("reap lease child"); + assert!(was_blocked, "child process must own the lease"); + ExclusiveFileLease::try_acquire(&lock_path).expect("lease after abnormal process exit"); +} + +#[test] +fn abnormal_exit_child_holds_lease() { + if std::env::var_os("OPENBITFUN_EXCLUSIVE_FILE_LEASE_CHILD").is_none() { + return; + } + let lock_path = std::path::PathBuf::from( + std::env::var_os("OPENBITFUN_EXCLUSIVE_FILE_LEASE_LOCK_PATH").expect("child lock path"), + ); + let ready_path = std::path::PathBuf::from( + std::env::var_os("OPENBITFUN_EXCLUSIVE_FILE_LEASE_READY_PATH").expect("child ready path"), + ); + let _lease = ExclusiveFileLease::try_acquire(Path::new(&lock_path)).expect("child lease"); + std::fs::write(ready_path, b"ready").expect("publish child readiness"); + loop { + std::thread::park(); + } +}