From 8dc317470ae25540dc6da8d4fc98d7a6e1b7abf1 Mon Sep 17 00:00:00 2001 From: nonoqing Date: Mon, 21 Sep 2026 17:59:37 +0800 Subject: [PATCH] fix(cron): keep scheduled jobs single-owner across desktop instances Two desktop instances share one user data directory: the user root is not keyed by bundle identifier, so the dev and release apps each ran their own scheduler over the same jobs.json. Both fired the same trigger, and the loser retried a turn id the winner had already consumed every five seconds until the job was permanently stuck as overdue. Scheduling now has one owner per user data directory: - services-core gains ExclusiveFileLease, a typed cross-process lease over one resource path that reuses the existing FileLock. A second holder gets InUse, a leftover lock file never implies ownership, and an abnormal process exit releases the lease. - CronService takes that lease before touching the store. The holder schedules jobs and is the only writer of jobs.json. A leapless instance becomes standby: it refuses job changes with an explicit error, skips turn lifecycle bookkeeping, and reads jobs from disk so it reports the owner's state instead of a stale startup snapshot. - A standby re-checks the lease every 30 seconds and, on takeover, adopts the store from disk before scheduling from it. Ownership is published only after that adoption, so a takeover cannot schedule from a stale map. - A lease that cannot be created at all fails open with an error log, which keeps the pre-lease behavior instead of silently disabling scheduled jobs. - The cron tool reports the standby state in its list result, so an agent does not present another instance's job list as the state that will run. Delivery also converges when a trigger was consumed already: that outcome now clears the pending trigger and stops retrying instead of failing forever. Verified on base e60f92a01: - cargo test -p openbitfun-core --no-default-features --features agent-runtime,scheduled-jobs,git --lib service::cron (22 passed) - cargo test -p openbitfun-core --no-default-features --features agent-runtime,scheduled-jobs,git --lib agentic::tools::implementations::cron_tool (11 passed) - cargo test -p openbitfun-agent-runtime --features agent-runtime --test agent_session_contracts scheduled_job (11 passed) - cargo test -p openbitfun-services-core --no-default-features --features local-storage --test exclusive_file_lease_contracts (7 passed) - cargo test -p openbitfun-services-core --no-default-features --features local-storage --test session_write_lock_contracts (10 passed) - cargo check -p openbitfun-core --no-default-features --features agent-runtime,scheduled-jobs,git --lib - pnpm run check:core-boundaries Not covered: other files under the user data root (config/app.json, token usage) still have no cross-process write protection, and an older binary does not know about this lease, so mixed-version instances stay unprotected. No test builds a full CronService, so the service wiring is covered by lease unit tests and contract tests, not end to end. Co-authored-by: bitfun-ai --- .../explicit-test-topology.mjs | 1 + .../src/agentic/coordination/scheduler.rs | 11 +- .../tools/implementations/cron_tool.rs | 8 + .../infrastructure/app_paths/path_manager.rs | 9 + .../assembly/core/src/service/cron/service.rs | 447 ++++++++++++++++-- .../agent-runtime/src/scheduled_job.rs | 17 + .../scheduled_job_contracts.rs | 28 ++ src/crates/services/services-core/AGENTS.md | 1 + src/crates/services/services-core/Cargo.toml | 5 + .../services-core/src/exclusive_file_lease.rs | 138 ++++++ .../services/services-core/src/file_lock.rs | 21 + src/crates/services/services-core/src/lib.rs | 2 + .../services-core/src/session/write_lock.rs | 18 +- .../tests/exclusive_file_lease_contracts.rs | 140 ++++++ 14 files changed, 792 insertions(+), 54 deletions(-) create mode 100644 src/crates/services/services-core/src/exclusive_file_lease.rs create mode 100644 src/crates/services/services-core/tests/exclusive_file_lease_contracts.rs 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(); + } +}