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 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
96pub(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 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}