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