Skip to main content

roder_core/
execution_lease.rs

1//! Local execution authority for a durable hosted-runtime owner. The supervisor
2//! extends this deadline only after a confirmed renewal of the same generation.
3//! Durable writes and remote/browser actions still require their own fencing.
4use crate::Runtime;
5use std::{
6    sync::{
7        Arc, Mutex,
8        atomic::{AtomicUsize, Ordering},
9    },
10    time::Instant,
11};
12
13#[derive(Debug)]
14pub struct RuntimeExecutionLease {
15    deadline: Mutex<Option<Instant>>,
16    executions: AtomicUsize,
17}
18
19impl RuntimeExecutionLease {
20    /// Use request-start + TTL, never a database wall-clock timestamp. If the
21    /// request took too long, this deadline can already be expired and fails closed.
22    pub fn new(deadline: Instant) -> Self {
23        Self {
24            deadline: Mutex::new(Some(deadline)),
25            executions: AtomicUsize::new(0),
26        }
27    }
28
29    /// A failed or late renewal cannot resurrect a runtime that lost authority.
30    pub fn renew(&self, deadline: Instant) -> anyhow::Result<()> {
31        let mut current = self
32            .deadline
33            .lock()
34            .map_err(|_| anyhow::anyhow!("runtime execution lease unavailable"))?;
35        let now = Instant::now();
36        match *current {
37            Some(existing) if existing > now && deadline > now => {
38                *current = Some(existing.max(deadline));
39                Ok(())
40            }
41            _ => {
42                *current = None;
43                anyhow::bail!("runtime execution lease expired or revoked")
44            }
45        }
46    }
47
48    pub(crate) fn enter(self: &Arc<Self>) -> anyhow::Result<ExecutionPermit> {
49        let mut deadline = self
50            .deadline
51            .lock()
52            .map_err(|_| anyhow::anyhow!("runtime execution lease unavailable"))?;
53        if !deadline.is_some_and(|deadline| deadline > Instant::now()) {
54            *deadline = None;
55            anyhow::bail!("runtime execution lease expired or revoked");
56        }
57        self.executions.fetch_add(1, Ordering::AcqRel);
58        Ok(ExecutionPermit(self.clone()))
59    }
60
61    pub(crate) fn seal_if_idle(&self) -> anyhow::Result<bool> {
62        let mut deadline = self
63            .deadline
64            .lock()
65            .map_err(|_| anyhow::anyhow!("runtime execution lease unavailable"))?;
66        if self.executions.load(Ordering::Acquire) != 0 {
67            return Ok(false);
68        }
69        *deadline = None;
70        Ok(true)
71    }
72
73    pub fn revoke(&self) {
74        if let Ok(mut deadline) = self.deadline.lock() {
75            *deadline = None;
76        }
77    }
78
79    pub fn require_live(&self) -> anyhow::Result<()> {
80        let mut deadline = self
81            .deadline
82            .lock()
83            .map_err(|_| anyhow::anyhow!("runtime execution lease unavailable"))?;
84        if deadline.is_some_and(|deadline| deadline > Instant::now()) {
85            return Ok(());
86        }
87        *deadline = None;
88        anyhow::bail!("runtime execution lease expired or revoked")
89    }
90}
91
92/// Counts local tool futures, including approval and cleanup waits. A successful
93/// planned handoff cannot race a new dispatch or abandon an admitted future.
94pub(crate) struct ExecutionPermit(Arc<RuntimeExecutionLease>);
95impl Drop for ExecutionPermit {
96    fn drop(&mut self) {
97        self.0.executions.fetch_sub(1, Ordering::AcqRel);
98    }
99}
100
101impl Runtime {
102    /// Bind before publishing the runtime or accepting work. Replacing a lost
103    /// generation requires constructing a new runtime, not resetting this guard.
104    pub fn with_execution_lease(mut self, lease: Arc<RuntimeExecutionLease>) -> Self {
105        self.execution_lease = Some(lease);
106        self
107    }
108
109    pub fn uses_execution_lease(&self, lease: &Arc<RuntimeExecutionLease>) -> bool {
110        self.execution_lease
111            .as_ref()
112            .is_some_and(|bound| Arc::ptr_eq(bound, lease))
113    }
114
115    pub fn ensure_execution_authority(&self) -> anyhow::Result<()> {
116        if let Some(lease) = &self.execution_lease {
117            lease.require_live()?;
118        }
119        Ok(())
120    }
121}
122
123#[cfg(test)]
124mod tests {
125    use super::*;
126    use crate::{RuntimeConfig, StartTurnRequest, fake_provider::FakeInferenceEngine};
127    use roder_api::{extension::ExtensionRegistryBuilder, tools::*};
128    use std::{
129        sync::atomic::{AtomicUsize, Ordering},
130        time::Duration,
131    };
132
133    #[test]
134    fn late_renewal_and_revocation_cannot_restore_authority() {
135        let expired = RuntimeExecutionLease::new(Instant::now());
136        assert!(
137            expired
138                .renew(Instant::now() + Duration::from_secs(30))
139                .is_err()
140        );
141        let live = RuntimeExecutionLease::new(Instant::now() + Duration::from_secs(30));
142        live.renew(Instant::now() + Duration::from_secs(60))
143            .unwrap();
144        live.require_live().unwrap();
145        live.revoke();
146        assert!(live.require_live().is_err());
147        assert!(
148            live.renew(Instant::now() + Duration::from_secs(60))
149                .is_err()
150        );
151    }
152
153    #[test]
154    fn revoked_owner_sealing_waits_for_admitted_execution() {
155        let lease = Arc::new(RuntimeExecutionLease::new(
156            Instant::now() + Duration::from_secs(30),
157        ));
158        let permit = lease.enter().unwrap();
159        lease.revoke();
160        assert!(!lease.seal_if_idle().unwrap());
161        drop(permit);
162        assert!(lease.seal_if_idle().unwrap());
163        assert!(lease.seal_if_idle().unwrap());
164        assert!(lease.require_live().is_err());
165        assert!(lease.enter().is_err());
166        assert!(
167            lease
168                .renew(Instant::now() + Duration::from_secs(30))
169                .is_err()
170        );
171    }
172
173    struct Probe(Arc<AtomicUsize>);
174    impl ToolContributor for Probe {
175        fn id(&self) -> String {
176            "lease-probe".into()
177        }
178        fn contribute(&self, registry: &mut ToolRegistry) -> anyhow::Result<()> {
179            registry.register(Arc::new(Probe(self.0.clone())))
180        }
181    }
182    #[async_trait::async_trait]
183    impl ToolExecutor for Probe {
184        fn spec(&self) -> ToolSpec {
185            ToolSpec {
186                name: "write_file".into(),
187                description: "probe".into(),
188                parameters: serde_json::json!({"type":"object","properties":{}}),
189            }
190        }
191        async fn execute(
192            &self,
193            _ctx: ToolExecutionContext,
194            call: ToolCall,
195        ) -> anyhow::Result<ToolResult> {
196            self.0.fetch_add(1, Ordering::SeqCst);
197            Ok(ToolResult {
198                id: call.id,
199                name: call.name,
200                text: "ok".into(),
201                data: serde_json::json!({}),
202                is_error: false,
203            })
204        }
205    }
206
207    #[tokio::test]
208    async fn revoked_owner_cannot_start_turns_or_execute_more_tools() {
209        let count = Arc::new(AtomicUsize::new(0));
210        let lease = Arc::new(RuntimeExecutionLease::new(
211            Instant::now() + Duration::from_secs(30),
212        ));
213        let mut registry = ExtensionRegistryBuilder::new();
214        registry.inference_engine(Arc::new(FakeInferenceEngine));
215        registry.tool_contributor(Arc::new(Probe(count.clone())));
216        let runtime = Arc::new(
217            Runtime::new(
218                registry.build().unwrap(),
219                RuntimeConfig {
220                    policy_mode: roder_api::policy_mode::PolicyMode::Bypass,
221                    ..Default::default()
222                },
223            )
224            .unwrap()
225            .with_execution_lease(lease.clone()),
226        );
227        let thread = runtime.create_thread(None).await.unwrap().thread_id;
228        let call = || roder_api::inference::ToolCallCompleted {
229            id: uuid::Uuid::new_v4().to_string(),
230            name: "write_file".into(),
231            arguments: "{}".into(),
232        };
233        runtime
234            .route_tool_call(&thread, &"turn".into(), call(), None, None)
235            .await
236            .unwrap();
237        assert_eq!(count.load(Ordering::SeqCst), 1);
238        lease.revoke();
239        assert!(
240            runtime
241                .route_tool_call(&thread, &"turn".into(), call(), None, None)
242                .await
243                .unwrap_err()
244                .to_string()
245                .contains("execution lease")
246        );
247        assert_eq!(count.load(Ordering::SeqCst), 1);
248        let error = runtime
249            .start_turn(StartTurnRequest {
250                thread_id: thread,
251                message: "do work".into(),
252                images: Vec::new(),
253                provider_override: None,
254                model_override: None,
255                reasoning_override: None,
256                workspace: std::env::current_dir().unwrap().display().to_string(),
257                instructions: crate::default_instructions(),
258                developer_context: None,
259                task_ledger_required: false,
260                service_tier_override: None,
261            })
262            .await
263            .unwrap_err();
264        assert!(error.to_string().contains("execution lease"));
265    }
266    #[tokio::test]
267    async fn ownership_is_rechecked_after_a_user_approval_wait() {
268        let count = Arc::new(AtomicUsize::new(0));
269        let lease = Arc::new(RuntimeExecutionLease::new(
270            Instant::now() + Duration::from_secs(30),
271        ));
272        let mut registry = ExtensionRegistryBuilder::new();
273        registry.inference_engine(Arc::new(FakeInferenceEngine));
274        registry.tool_contributor(Arc::new(Probe(count.clone())));
275        let runtime = Arc::new(
276            Runtime::new(registry.build().unwrap(), RuntimeConfig::default())
277                .unwrap()
278                .with_execution_lease(lease.clone()),
279        );
280        let thread = runtime.create_thread(None).await.unwrap().thread_id;
281        let mut events = runtime.subscribe_events();
282        let task_runtime = runtime.clone();
283        let task = tokio::spawn(async move {
284            task_runtime
285                .route_tool_call(
286                    &thread,
287                    &"turn".into(),
288                    roder_api::inference::ToolCallCompleted {
289                        id: "approval-lease-call".into(),
290                        name: "write_file".into(),
291                        arguments: "{}".into(),
292                    },
293                    None,
294                    None,
295                )
296                .await
297        });
298        tokio::time::timeout(Duration::from_secs(3), async {
299            loop {
300                if matches!(
301                    events.recv().await.unwrap().event,
302                    roder_api::events::RoderEvent::ApprovalRequested(_)
303                ) {
304                    break;
305                }
306            }
307        })
308        .await
309        .unwrap();
310        lease.revoke();
311        assert!(
312            runtime
313                .resolve_tool_approval("approval-lease-call", true)
314                .await
315                .unwrap()
316        );
317        let result = tokio::time::timeout(Duration::from_secs(3), task)
318            .await
319            .unwrap()
320            .unwrap()
321            .unwrap();
322        assert!(result.is_error);
323        assert_eq!(count.load(Ordering::SeqCst), 0);
324    }
325    #[tokio::test]
326    async fn handoff_waits_for_approved_tool_to_finish_without_interrupting_it() {
327        let count = Arc::new(AtomicUsize::new(0));
328        let lease = Arc::new(RuntimeExecutionLease::new(
329            Instant::now() + Duration::from_secs(30),
330        ));
331        let mut registry = ExtensionRegistryBuilder::new();
332        registry.inference_engine(Arc::new(FakeInferenceEngine));
333        registry.tool_contributor(Arc::new(Probe(count.clone())));
334        let runtime = Arc::new(
335            Runtime::new(registry.build().unwrap(), RuntimeConfig::default())
336                .unwrap()
337                .with_execution_lease(lease.clone()),
338        );
339        let thread = runtime.create_thread(None).await.unwrap().thread_id;
340        let mut events = runtime.subscribe_events();
341        let task_runtime = runtime.clone();
342        let task = tokio::spawn(async move {
343            task_runtime
344                .route_tool_call(
345                    &thread,
346                    &"turn".into(),
347                    roder_api::inference::ToolCallCompleted {
348                        id: "approval-lease-call".into(),
349                        name: "write_file".into(),
350                        arguments: "{}".into(),
351                    },
352                    None,
353                    None,
354                )
355                .await
356        });
357        tokio::time::timeout(Duration::from_secs(3), async {
358            loop {
359                if matches!(
360                    events.recv().await.unwrap().event,
361                    roder_api::events::RoderEvent::ApprovalRequested(_)
362                ) {
363                    break;
364                }
365            }
366        })
367        .await
368        .unwrap();
369        assert!(!runtime.seal_idle_owner().await.unwrap());
370        lease.require_live().unwrap();
371        assert!(
372            runtime
373                .resolve_tool_approval("approval-lease-call", true)
374                .await
375                .unwrap()
376        );
377        let result = tokio::time::timeout(Duration::from_secs(3), task)
378            .await
379            .unwrap()
380            .unwrap()
381            .unwrap();
382        assert!(!result.is_error);
383        assert_eq!(count.load(Ordering::SeqCst), 1);
384        assert!(runtime.seal_idle_owner().await.unwrap());
385        assert!(lease.require_live().is_err());
386    }
387    #[tokio::test]
388    async fn handoff_leaves_an_active_provider_turn_running_until_completion() {
389        use roder_api::inference::*;
390        struct WaitingEngine {
391            entered: tokio::sync::Notify,
392            release: tokio::sync::Notify,
393        }
394        #[async_trait::async_trait]
395        impl InferenceEngine for WaitingEngine {
396            fn id(&self) -> String {
397                "mock".into()
398            }
399            fn capabilities(&self) -> InferenceCapabilities {
400                InferenceCapabilities::text_only()
401            }
402            async fn list_models(
403                &self,
404                _: InferenceProviderContext<'_>,
405            ) -> anyhow::Result<Vec<ModelDescriptor>> {
406                Ok(vec![])
407            }
408            async fn stream_turn(
409                &self,
410                _: InferenceTurnContext<'_>,
411                _: AgentInferenceRequest,
412            ) -> anyhow::Result<InferenceEventStream> {
413                self.entered.notify_one();
414                self.release.notified().await;
415                Ok(Box::pin(futures::stream::iter([Ok(
416                    InferenceEvent::Completed(CompletionMetadata {
417                        stop_reason: Some("completed".into()),
418                        provider_response_id: None,
419                    }),
420                )])))
421            }
422        }
423        let engine = Arc::new(WaitingEngine {
424            entered: Default::default(),
425            release: Default::default(),
426        });
427        let lease = Arc::new(RuntimeExecutionLease::new(
428            Instant::now() + Duration::from_secs(30),
429        ));
430        let mut registry = ExtensionRegistryBuilder::new();
431        registry.inference_engine(engine.clone());
432        let runtime = Arc::new(
433            Runtime::new(registry.build().unwrap(), RuntimeConfig::default())
434                .unwrap()
435                .with_execution_lease(lease.clone()),
436        );
437        let thread = runtime.create_thread(None).await.unwrap().thread_id;
438        runtime
439            .start_turn(StartTurnRequest {
440                thread_id: thread,
441                message: "finish normally".into(),
442                images: vec![],
443                provider_override: None,
444                model_override: None,
445                reasoning_override: None,
446                workspace: std::env::current_dir().unwrap().display().to_string(),
447                instructions: crate::default_instructions(),
448                developer_context: None,
449                task_ledger_required: false,
450                service_tier_override: None,
451            })
452            .await
453            .unwrap();
454        tokio::time::timeout(Duration::from_secs(3), engine.entered.notified())
455            .await
456            .unwrap();
457        assert!(!runtime.seal_idle_owner().await.unwrap());
458        assert_eq!(runtime.active_turn_count().await, 1);
459        lease.require_live().unwrap();
460        engine.release.notify_one();
461        tokio::time::timeout(Duration::from_secs(3), async {
462            while !runtime.seal_idle_owner().await.unwrap() {
463                tokio::task::yield_now().await;
464            }
465        })
466        .await
467        .unwrap();
468        assert_eq!(runtime.active_turn_count().await, 0);
469        assert!(lease.require_live().is_err());
470    }
471}