1use 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 pub fn new(deadline: Instant) -> Self {
23 Self {
24 deadline: Mutex::new(Some(deadline)),
25 executions: AtomicUsize::new(0),
26 }
27 }
28
29 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
92pub(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 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}