diff --git a/frontend/src-tauri/src/extension_container.rs b/frontend/src-tauri/src/extension_container.rs index ee02d27..b737f75 100644 --- a/frontend/src-tauri/src/extension_container.rs +++ b/frontend/src-tauri/src/extension_container.rs @@ -707,12 +707,51 @@ mod tests { .code, "EXTENSION_PROCESS_WAIT_INVALID" ); - running.terminate().unwrap(); + let completed = running.start_tool_call().unwrap(); + completed.finish().unwrap(); + assert_eq!(running.wait(std::time::Duration::ZERO).unwrap(), None); + let deadline = running + .start_test_tool_call(std::time::Duration::from_millis(100)) + .unwrap(); + // No check() or finish() drives expiration; wait only on the OS process. + assert!(running + .wait(std::time::Duration::from_secs(5)) + .unwrap() + .is_some()); + assert_eq!( + deadline.finish().unwrap_err().code, + "EXTENSION_TOOL_DEADLINE_EXCEEDED" + ); assert!(running .wait(std::time::Duration::from_secs(5)) .unwrap() .is_some()); drop(running); + for explicit_cancel in [false, true] { + let data = crate::extension_launch_data::LaunchData::new( + &executable, + &["wait".to_owned()], + &system, + &folder, + &folder.join("Temp"), + &std::collections::BTreeMap::new(), + ) + .unwrap(); + let suspended = + crate::extension_process::Suspended::create(&profile, &executable, data).unwrap(); + let running = unsafe { suspended.resume().unwrap() }; + let deadline = running.start_tool_call().unwrap(); + assert_eq!(running.wait(std::time::Duration::ZERO).unwrap(), None); + if explicit_cancel { + deadline.cancel().unwrap(); + } else { + drop(deadline); + } + assert!(running + .wait(std::time::Duration::from_secs(5)) + .unwrap() + .is_some()); + } for ip in ["127.0.0.1:0", "[::1]:0"] { let tcp = TcpListener::bind(ip).unwrap(); let udp = UdpSocket::bind(ip).unwrap(); diff --git a/frontend/src-tauri/src/extension_deadline.rs b/frontend/src-tauri/src/extension_deadline.rs new file mode 100644 index 0000000..300a1f2 --- /dev/null +++ b/frontend/src-tauri/src/extension_deadline.rs @@ -0,0 +1,158 @@ +//! Per-tool-call deadline. A persistent server does not have a 60-second lifetime. +use crate::{ + extension_job::Job, + workspace::{HostError, Result}, +}; +use std::{ + sync::{ + atomic::{AtomicU8, Ordering}, + Arc, Condvar, Mutex, + }, + thread::JoinHandle, + time::{Duration, Instant}, +}; + +pub const TOOL_BUDGET: Duration = Duration::from_secs(60); +struct State { + cancelled: Mutex, + wake: Condvar, + outcome: AtomicU8, +} +/// Host-owned guard, created before dispatching a tool call. Finish it only when +/// the call completes. Expiration kills the entire instance group even if the +/// caller never polls. Dropping without finish aborts the instance; only explicit +/// successful completion cancels the timer while preserving the server. +pub struct ToolDeadline { + state: Arc, + job: Arc, + expires: Instant, + worker: Option>, +} +impl ToolDeadline { + pub(crate) fn arm(job: &Job) -> Result { + Self::with_budget(job, TOOL_BUDGET) + } + fn with_budget(job: &Job, budget: Duration) -> Result { + if budget.is_zero() || budget > TOOL_BUDGET { + return Err(HostError::new("EXTENSION_TOOL_BUDGET_INVALID")); + } + let expires = Instant::now() + budget; + let job = Arc::new(job.clone_for_deadline()?); + let state = Arc::new(State { + cancelled: Mutex::new(false), + wake: Condvar::new(), + outcome: AtomicU8::new(0), + }); + let thread_job = Arc::clone(&job); + let thread_state = Arc::clone(&state); + let worker = std::thread::Builder::new() + .name("extension-tool-deadline".into()) + .spawn(move || { + let mut cancelled = thread_state + .cancelled + .lock() + .unwrap_or_else(|e| e.into_inner()); + loop { + if *cancelled { + return; + } + let remaining = expires.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + break; + } + cancelled = thread_state + .wake + .wait_timeout(cancelled, remaining) + .unwrap_or_else(|e| e.into_inner()) + .0; + } + thread_state.outcome.store(1, Ordering::Release); + drop(cancelled); + if thread_job.terminate().is_err() { + thread_state.outcome.store(2, Ordering::Release); + } + }) + .map_err(|_| HostError::new("EXTENSION_TOOL_WATCHDOG_UNAVAILABLE"))?; + Ok(Self { + state, + job, + expires, + worker: Some(worker), + }) + } + pub fn check(&self) -> Result<()> { + match self.state.outcome.load(Ordering::Acquire) { + 2 => Err(HostError::new("EXTENSION_RESOURCE_TERMINATE_FAILED")), + 3 => Err(HostError::new("EXTENSION_TOOL_WATCHDOG_FAILED")), + 1 => Err(HostError::new("EXTENSION_TOOL_DEADLINE_EXCEEDED")), + _ if Instant::now() >= self.expires => { + Err(HostError::new("EXTENSION_TOOL_DEADLINE_EXCEEDED")) + } + _ => Ok(()), + } + } + /// Completion cannot cancel an already elapsed budget, even if the timer + /// thread has not yet been scheduled to observe expiration. + pub fn finish(mut self) -> Result<()> { + self.stop(); + match self.state.outcome.load(Ordering::Acquire) { + 0 => Ok(()), + 2 => Err(HostError::new("EXTENSION_RESOURCE_TERMINATE_FAILED")), + 3 => Err(HostError::new("EXTENSION_TOOL_WATCHDOG_FAILED")), + _ => Err(HostError::new("EXTENSION_TOOL_DEADLINE_EXCEEDED")), + } + } + /// Explicit abandonment terminates the instance and reports kill failures. + pub fn cancel(mut self) -> Result<()> { + let result = self.job.terminate(); + self.stop(); + result + } + fn stop(&mut self) { + if self.worker.is_none() { + return; + } + let mut cancelled = self + .state + .cancelled + .lock() + .unwrap_or_else(|e| e.into_inner()); + if Instant::now() < self.expires { + *cancelled = true; + } + self.state.wake.notify_all(); + drop(cancelled); + if self.worker.take().unwrap().join().is_err() { + let _ = self.job.terminate(); + self.state.outcome.store(3, Ordering::Release); + } + } + #[cfg(test)] + pub(crate) fn arm_test(job: &Job, budget: Duration) -> Result { + Self::with_budget(job, budget) + } +} +impl Drop for ToolDeadline { + fn drop(&mut self) { + if self.worker.is_some() { + let _ = self.job.terminate(); + self.stop(); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + #[test] + fn rejects_invalid_budgets_and_finishes_without_waiting_for_full_budget() { + let job = Job::new().unwrap(); + assert!(ToolDeadline::with_budget(&job, Duration::ZERO).is_err()); + assert!(ToolDeadline::with_budget(&job, Duration::from_secs(61)).is_err()); + let deadline = ToolDeadline::arm(&job).unwrap(); + assert!(deadline.check().is_ok()); + let start = Instant::now(); + deadline.finish().unwrap(); + assert!(start.elapsed() < Duration::from_secs(5)); + } +} diff --git a/frontend/src-tauri/src/extension_job.rs b/frontend/src-tauri/src/extension_job.rs index 7bcb6c9..989e4ec 100644 --- a/frontend/src-tauri/src/extension_job.rs +++ b/frontend/src-tauri/src/extension_job.rs @@ -13,6 +13,14 @@ pub struct Job { handle: OwnedHandle, } impl Job { + pub(crate) fn clone_for_deadline(&self) -> Result { + Ok(Self { + handle: self + .handle + .try_clone() + .map_err(|_| HostError::new("EXTENSION_RESOURCE_UNAVAILABLE"))?, + }) + } pub fn new() -> Result { Self::with_process_limit(16) } diff --git a/frontend/src-tauri/src/extension_process.rs b/frontend/src-tauri/src/extension_process.rs index 6f8a6c0..59c34a6 100644 --- a/frontend/src-tauri/src/extension_process.rs +++ b/frontend/src-tauri/src/extension_process.rs @@ -175,6 +175,18 @@ impl<'a> Suspended<'a> { } } impl Running<'_> { + /// Arm before dispatching a tool request; finish after receiving its result. + /// Failure to arm must prevent dispatch. This does not time server lifetime. + pub fn start_tool_call(&self) -> Result { + crate::extension_deadline::ToolDeadline::arm(&self.0.job) + } + #[cfg(test)] + pub(crate) fn start_test_tool_call( + &self, + budget: Duration, + ) -> Result { + crate::extension_deadline::ToolDeadline::arm_test(&self.0.job, budget) + } /// A bounded observation only. The runtime must enforce the tool deadline. pub fn wait(&self, timeout: Duration) -> Result> { let milliseconds = u32::try_from(timeout.as_millis()) diff --git a/frontend/src-tauri/src/lib.rs b/frontend/src-tauri/src/lib.rs index 4dedfea..c0f8ef9 100644 --- a/frontend/src-tauri/src/lib.rs +++ b/frontend/src-tauri/src/lib.rs @@ -63,3 +63,6 @@ pub mod extension_launch_data; #[cfg(windows)] pub mod extension_process; + +#[cfg(windows)] +pub mod extension_deadline;