1use std::sync::Arc;
2use std::time::{Duration, Instant};
3
4use crate::chunk::{Chunk, ChunkRef, Op};
5use crate::value::{ModuleFunctionRegistry, VmError, VmValue};
6
7use super::callable_entry::TopLevelEntry;
8use super::state::ExecutionDeadlineState;
9use super::{CallFrame, LocalSlot, Vm};
10
11const CANCEL_GRACE_ASYNC_OP: Duration = Duration::from_millis(250);
12
13pub(super) fn new_execution_deadline_state(
14 deadline: Option<Instant>,
15) -> Arc<ExecutionDeadlineState> {
16 ExecutionDeadlineState::new(Instant::now(), deadline)
17}
18
19#[cfg(test)]
20thread_local! {
21 static SCOPE_INTERRUPT_ASYNC_DISPATCHES: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
22}
23
24#[cfg(test)]
25pub(super) fn reset_scope_interrupt_async_dispatches() {
26 SCOPE_INTERRUPT_ASYNC_DISPATCHES.set(0);
27}
28
29#[cfg(test)]
30pub(super) fn scope_interrupt_async_dispatches() -> u64 {
31 SCOPE_INTERRUPT_ASYNC_DISPATCHES.get()
32}
33
34#[derive(Clone, Copy)]
35enum DeadlineKind {
36 Execution,
37 Scope,
38 InterruptHandler,
39}
40
41impl Vm {
42 #[inline]
49 pub(crate) fn scope_interrupts_clean(&self) -> bool {
50 self.requested_process_exit().is_none()
51 && self.cancel_token.is_none()
52 && self.interrupt_signal_token.is_none()
53 && self.pending_interrupt_signal.is_none()
54 && self.interrupt_handler_deadline.is_none()
55 && !self.execution_deadline.is_active()
56 && self.deadlines.is_empty()
57 }
58
59 pub async fn execute(&mut self, chunk: &Chunk) -> Result<VmValue, VmError> {
68 self.execute_arc(Arc::new(chunk.clone())).await
69 }
70
71 pub async fn execute_with_timeout(
82 &mut self,
83 chunk: &Chunk,
84 timeout: Duration,
85 ) -> Result<VmValue, VmError> {
86 self.execute_top_level_with_timeout(TopLevelEntry::Chunk(Arc::new(chunk.clone())), timeout)
87 .await
88 }
89
90 pub(super) async fn execute_top_level_with_timeout(
91 &mut self,
92 entry: TopLevelEntry,
93 timeout: Duration,
94 ) -> Result<VmValue, VmError> {
95 let deadline = Instant::now().checked_add(timeout).ok_or_else(|| {
96 VmError::Runtime("execution timeout exceeds the platform clock range".to_string())
97 })?;
98 crate::orchestration::scope_ambient_transaction(async {
99 let pipeline_checkpoint = crate::orchestration::checkpoint_pipeline_lifecycle();
100 let deadline_guard = self.execution_deadline.install(deadline);
101 let result = crate::tracing::checkpoint_future(self.execute_top_level(entry)).await;
102 deadline_guard.complete();
103 pipeline_checkpoint.complete();
104 result
105 })
106 .await
107 }
108
109 pub async fn execute_arc(&mut self, chunk: ChunkRef) -> Result<VmValue, VmError> {
115 self.execute_top_level(TopLevelEntry::Chunk(chunk)).await
116 }
117
118 pub(super) async fn run_pipeline_finish_lifecycle(
125 &mut self,
126 value: VmValue,
127 ) -> Result<VmValue, VmError> {
128 use crate::orchestration::{
129 take_pipeline_on_finish, unsettled_state_snapshot_async, HookEvent,
130 };
131 let _tape_phase =
132 crate::testbench::tape::enter_phase(crate::testbench::tape::TapePhase::RuntimeFinalize);
133
134 let on_finish = take_pipeline_on_finish();
135 let unsettled = unsettled_state_snapshot_async().await;
136
137 let pre_payload = serde_json::json!({
138 "event": HookEvent::PreFinish.as_str(),
139 "return_value": crate::llm::vm_value_to_json(&value),
140 "unsettled": unsettled.to_json(),
141 "has_on_finish": on_finish.is_some(),
142 });
143 self.fire_finish_lifecycle_event(HookEvent::PreFinish, &pre_payload)
144 .await?;
145
146 if !unsettled.is_empty() {
147 let payload = serde_json::json!({
148 "event": HookEvent::OnUnsettledDetected.as_str(),
149 "unsettled": unsettled.to_json(),
150 });
151 self.fire_finish_lifecycle_event(HookEvent::OnUnsettledDetected, &payload)
152 .await?;
153 }
154
155 let final_value = if let Some(closure) = on_finish {
156 let harness_value = self.root_harness_value().ok_or_else(|| {
157 VmError::Runtime(
158 "pipeline finish callback requires Harness, but no root Harness is installed"
159 .to_string(),
160 )
161 })?;
162 self.call_closure_pub(&closure, &[harness_value, value])
163 .await?
164 } else {
165 value
166 };
167
168 let post_payload = serde_json::json!({
169 "event": HookEvent::PostFinish.as_str(),
170 "return_value": crate::llm::vm_value_to_json(&final_value),
171 "unsettled": unsettled.to_json(),
172 });
173 self.fire_finish_lifecycle_event(HookEvent::PostFinish, &post_payload)
174 .await?;
175
176 Ok(final_value)
177 }
178
179 async fn fire_finish_lifecycle_event(
197 &mut self,
198 event: crate::orchestration::HookEvent,
199 payload: &serde_json::Value,
200 ) -> Result<(), VmError> {
201 use crate::orchestration::{HookControl, HookEvent};
202 let invocations = crate::orchestration::matching_vm_lifecycle_hooks(event, payload);
203 if invocations.is_empty() {
204 return Ok(());
205 }
206 let harness = self.root_harness_value().ok_or_else(|| {
207 VmError::Runtime(
208 "pipeline lifecycle hook requires Harness, but no root Harness is installed"
209 .to_string(),
210 )
211 })?;
212 let mut current_payload = payload.clone();
213 for invocation in invocations {
214 let arg = crate::stdlib::json_to_vm_value(¤t_payload);
215 let closure = invocation.resolve(self).await?;
216 let raw = self
217 .call_closure_pub(&closure, &[harness.clone(), arg])
218 .await?;
219 let (action, effects) = crate::orchestration::collect_hook_effects_and_action(
220 event,
221 raw,
222 crate::value::VmValue::Nil,
223 )?;
224 crate::orchestration::inject_hook_effects_into_current_session(effects)?;
225 let control = crate::orchestration::parse_hook_control_for_finish(event, &action)?;
226 match control {
227 HookControl::Allow => {}
228 HookControl::Block { reason } => {
229 if matches!(event, HookEvent::PreFinish) {
230 return Err(VmError::Runtime(format!(
231 "PreFinish hook returned block, which is not a valid control: {reason}. \
232 To delay pipeline finish until unsettled work clears, use \
233 OnFinish.block_until_settled (std/lifecycle) or return Modify/Allow \
234 from PreFinish."
235 )));
236 }
237 if matches!(event, HookEvent::PostFinish) {
238 continue;
240 }
241 return Err(VmError::Runtime(format!(
243 "{} hook blocked pipeline finish: {reason}",
244 event.as_str()
245 )));
246 }
247 HookControl::Modify { payload: modified } => {
248 current_payload = modified;
249 }
250 HookControl::Decision { .. } => {}
251 }
252 }
253 Ok(())
254 }
255
256 pub(crate) fn handle_error(&mut self, error: VmError) -> Result<Option<VmValue>, VmError> {
258 if let Some(code) = error.process_exit_code() {
259 self.request_process_exit(code);
260 }
261 if error.is_uncatchable_control_flow() {
262 return Err(error);
263 }
264 let thrown_value = error.thrown_value();
265
266 if let Some(handler) = self.exception_handlers.pop() {
267 if let Some(error_type) = handler.error_type.as_deref() {
268 let matches = match &thrown_value {
270 VmValue::EnumVariant(enum_variant) => enum_variant.has_enum_name(error_type),
271 _ => false,
272 };
273 if !matches {
274 return self.handle_error(error);
275 }
276 }
277
278 self.release_sync_guards_after_unwind(handler.frame_depth, handler.env_scope_depth);
279
280 while self.frames.len() > handler.frame_depth {
281 if let Some(frame) = self.frames.pop() {
282 if let Some(ref dir) = frame.saved_source_dir {
283 crate::stdlib::set_thread_source_dir(dir);
284 }
285 self.iterators.truncate(frame.saved_iterator_depth);
286 self.env = frame.saved_env;
287 }
288 }
289 crate::step_runtime::prune_below_frame(self.frames.len());
290
291 while self
293 .deadlines
294 .last()
295 .is_some_and(|d| d.1 > handler.frame_depth)
296 {
297 self.deadlines.pop();
298 }
299
300 self.env.truncate_scopes(handler.env_scope_depth);
301
302 self.stack.truncate(handler.stack_depth);
303 self.stack.push(thrown_value);
304
305 if let Some(frame) = self.frames.last_mut() {
306 frame.ip = handler.catch_ip;
307 }
308
309 Ok(None)
310 } else {
311 Err(error)
312 }
313 }
314
315 pub(crate) async fn run_chunk(&mut self, chunk: ChunkRef) -> Result<VmValue, VmError> {
316 self.run_chunk_ref(chunk, 0, None, None, None, None).await
317 }
318
319 pub(crate) async fn run_chunk_ref(
320 &mut self,
321 chunk: ChunkRef,
322 argc: usize,
323 saved_source_dir: Option<std::path::PathBuf>,
324 module_functions: Option<ModuleFunctionRegistry>,
325 module_state: Option<crate::value::ModuleState>,
326 local_slots: Option<Vec<LocalSlot>>,
327 ) -> Result<VmValue, VmError> {
328 self.ensure_execution_available()?;
329 let debugger = self.debugger_attached();
330 let local_slots = local_slots.unwrap_or_else(|| Self::fresh_local_slots(&chunk));
331 let initial_env = if debugger {
332 Some(self.env.clone())
333 } else {
334 None
335 };
336 let initial_local_slots = if debugger {
337 Some(local_slots.clone())
338 } else {
339 None
340 };
341 let inline_cache_set = self.inline_cache_set_index_for_chunk(&chunk);
342 self.frames.push(CallFrame {
343 chunk,
344 inline_cache_set,
345 ip: 0,
346 stack_base: self.stack.len(),
347 saved_env: self.env.clone(),
348 initial_env,
349 initial_local_slots,
350 saved_iterator_depth: self.iterators.len(),
351 fn_name: crate::value::HarnStr::new(),
352 argc,
353 saved_source_dir,
354 module_functions,
355 module_state,
356 local_slots,
357 local_scope_base: self.env.scope_depth().saturating_sub(1),
358 local_scope_depth: 0,
359 });
360
361 self.drive_dispatch_loop(0, false).await
362 }
363
364 pub(crate) async fn drive_until_frame_depth(
371 &mut self,
372 target_depth: usize,
373 ) -> Result<VmValue, VmError> {
374 self.drive_dispatch_loop(target_depth, true).await
375 }
376
377 async fn drive_dispatch_loop(
391 &mut self,
392 target_depth: usize,
393 restore_on_final_pop: bool,
394 ) -> Result<VmValue, VmError> {
395 self.ensure_execution_available()?;
396 let _task_activity = self
397 .wait_for_graph
398 .register_task(self.runtime_context.task_id.clone());
399 loop {
400 if !self.scope_interrupts_clean() {
407 if let Some(err) = self.pending_scope_interrupt().await {
408 match self.handle_error(err) {
409 Ok(None) => continue,
410 Ok(Some(val)) => return Ok(val),
411 Err(e) => {
412 self.unwind_frames_to_depth(target_depth);
413 return Err(e);
414 }
415 }
416 }
417 }
418
419 let frame_depth = self.frames.len();
420 let frame = match self.frames.last_mut() {
421 Some(f) => f,
422 None => return Ok(self.stack.pop().unwrap_or(VmValue::Nil)),
423 };
424
425 if frame.ip >= frame.chunk.code.len() {
426 let val = self.stack.pop().unwrap_or(VmValue::Nil);
427 let val = self.run_step_post_hooks_for_current_frame(val).await?;
428 self.release_sync_guards_for_frame(self.frames.len());
429 let popped_frame = self.frames.pop().unwrap();
430 if let Some(ref dir) = popped_frame.saved_source_dir {
431 crate::stdlib::set_thread_source_dir(dir);
432 }
433 let current_depth = self.frames.len();
434 crate::step_runtime::prune_below_frame(current_depth);
435 while self.deadlines.last().is_some_and(|d| d.1 > current_depth) {
440 self.deadlines.pop();
441 }
442
443 let reached_target = current_depth <= target_depth;
444 if reached_target && !restore_on_final_pop {
445 return Ok(val);
448 }
449 self.iterators.truncate(popped_frame.saved_iterator_depth);
450 self.env = popped_frame.saved_env;
451 self.stack.truncate(popped_frame.stack_base);
452 if reached_target {
453 return Ok(val);
454 }
455 self.stack.push(val);
456 continue;
457 }
458
459 let op_byte = frame.chunk.code[frame.ip];
460 if let Some(coverage) = self.coverage.as_mut() {
466 coverage.record(&frame.chunk, frame.ip);
467 }
468 frame.ip += 1;
469
470 let op = match Op::from_byte(op_byte) {
475 Some(op) => op,
476 None => return Err(VmError::InvalidInstruction(op_byte)),
477 };
478 if let Some(recorder) = self.flight_recorder.as_ref() {
479 recorder.record_instruction(
480 &self.runtime_context.task_id,
481 frame_depth,
482 frame.fn_name.as_str(),
483 &frame.chunk,
484 self.source_file.as_deref(),
485 frame.ip - 1,
486 op,
487 );
488 }
489 let op_result: Result<(), VmError> = if let Some(result) = self.execute_op_sync(op) {
490 result
491 } else if self.scope_interrupts_clean() {
492 self.execute_op_async(op).await
493 } else {
494 match self.execute_op_with_scope_interrupts(op_byte).await {
495 Ok(Some(val)) => return Ok(val),
496 Ok(None) => Ok(()),
497 Err(e) => Err(e),
498 }
499 };
500
501 match op_result {
502 Ok(()) => continue,
503 Err(VmError::Return(val)) => {
504 let val = self.run_step_post_hooks_for_current_frame(val).await?;
505 if let Some(popped_frame) = self.frames.pop() {
506 self.release_sync_guards_for_frame(self.frames.len() + 1);
507 if let Some(ref dir) = popped_frame.saved_source_dir {
508 crate::stdlib::set_thread_source_dir(dir);
509 }
510 let current_depth = self.frames.len();
511 self.exception_handlers
512 .retain(|h| h.frame_depth <= current_depth);
513 crate::step_runtime::prune_below_frame(current_depth);
514 while self.deadlines.last().is_some_and(|d| d.1 > current_depth) {
515 self.deadlines.pop();
516 }
517
518 let reached_target = current_depth <= target_depth;
519 if reached_target && !restore_on_final_pop {
520 return Ok(val);
521 }
522 self.iterators.truncate(popped_frame.saved_iterator_depth);
523 self.env = popped_frame.saved_env;
524 self.stack.truncate(popped_frame.stack_base);
525 if reached_target {
526 return Ok(val);
527 }
528 self.stack.push(val);
529 } else {
530 return Ok(val);
531 }
532 }
533 Err(e) => {
534 if self.error_stack_trace.is_empty() {
536 self.error_stack_trace = self.capture_stack_trace();
537 }
538 let e = match self.apply_step_error_boundary(e) {
545 StepBoundaryOutcome::Returned(val) => {
546 self.error_stack_trace.clear();
547 if self.frames.len() <= target_depth {
548 return Ok(val);
549 }
550 self.stack.push(val);
551 continue;
552 }
553 StepBoundaryOutcome::Throw(err) => err,
554 };
555 match self.handle_error(e) {
556 Ok(None) => {
557 self.error_stack_trace.clear();
558 continue;
559 }
560 Ok(Some(val)) => return Ok(val),
561 Err(e) => {
562 self.unwind_frames_to_depth(target_depth);
563 return Err(self.enrich_error_with_line(e));
564 }
565 }
566 }
567 }
568 }
569 }
570
571 fn unwind_frames_to_depth(&mut self, target_depth: usize) {
578 while self.frames.len() > target_depth {
579 let frame_depth = self.frames.len();
580 if let Some(frame) = self.frames.pop() {
581 self.release_sync_guards_for_frame(frame_depth);
582 if let Some(ref dir) = frame.saved_source_dir {
583 crate::stdlib::set_thread_source_dir(dir);
584 }
585 self.iterators.truncate(frame.saved_iterator_depth);
586 self.env = frame.saved_env;
587 self.stack.truncate(frame.stack_base);
588 }
589 }
590 let current_depth = self.frames.len();
591 crate::step_runtime::prune_below_frame(current_depth);
592 while self.deadlines.last().is_some_and(|d| d.1 > current_depth) {
593 self.deadlines.pop();
594 }
595 }
596
597 pub(crate) fn apply_step_error_boundary(&mut self, error: VmError) -> StepBoundaryOutcome {
603 use crate::step_runtime;
604 if !step_runtime::is_step_budget_exhausted(&error) {
605 return StepBoundaryOutcome::Throw(error);
606 }
607 let Some(step_depth) = step_runtime::active_step_frame_depth() else {
608 return StepBoundaryOutcome::Throw(error);
609 };
610 if step_depth != self.frames.len() {
615 return StepBoundaryOutcome::Throw(error);
616 }
617 let boundary = step_runtime::with_active_step(|step| step.definition.boundary())
618 .unwrap_or(step_runtime::StepErrorBoundary::Fail);
619 match boundary {
620 step_runtime::StepErrorBoundary::Continue => {
621 if let Some(popped) = self.frames.pop() {
625 self.release_sync_guards_for_frame(self.frames.len() + 1);
626 if let Some(ref dir) = popped.saved_source_dir {
627 crate::stdlib::set_thread_source_dir(dir);
628 }
629 let current_depth = self.frames.len();
630 self.exception_handlers
631 .retain(|h| h.frame_depth <= current_depth);
632 step_runtime::pop_and_record(
633 current_depth + 1,
634 "skipped",
635 Some(step_runtime_error_message(&error)),
636 );
637 if self.frames.is_empty() {
638 return StepBoundaryOutcome::Returned(VmValue::Nil);
639 }
640 self.iterators.truncate(popped.saved_iterator_depth);
641 self.env = popped.saved_env;
642 self.stack.truncate(popped.stack_base);
643 }
644 StepBoundaryOutcome::Returned(VmValue::Nil)
645 }
646 step_runtime::StepErrorBoundary::Escalate => {
647 let identity = step_runtime::with_active_step(|step| {
648 (
649 step.definition.name.clone(),
650 step.definition.function.clone(),
651 )
652 });
653 step_runtime::pop_and_record(
654 step_depth,
655 "escalated",
656 Some(step_runtime_error_message(&error)),
657 );
658 let (step_name, function) = identity.unzip();
659 StepBoundaryOutcome::Throw(step_runtime::mark_escalated(
660 error,
661 step_name.as_deref(),
662 function.as_deref(),
663 ))
664 }
665 step_runtime::StepErrorBoundary::Fail => {
666 step_runtime::pop_and_record(
667 step_depth,
668 "failed",
669 Some(step_runtime_error_message(&error)),
670 );
671 StepBoundaryOutcome::Throw(error)
672 }
673 }
674 }
675}
676
677fn next_deadline(
678 execution_deadline: Option<Instant>,
679 scope_deadline: Option<Instant>,
680 interrupt_handler_deadline: Option<Instant>,
681) -> (Option<Instant>, Option<DeadlineKind>) {
682 [
683 (execution_deadline, DeadlineKind::Execution),
684 (scope_deadline, DeadlineKind::Scope),
685 (interrupt_handler_deadline, DeadlineKind::InterruptHandler),
686 ]
687 .into_iter()
688 .filter_map(|(deadline, kind)| deadline.map(|deadline| (deadline, kind)))
689 .min_by_key(|(deadline, _)| *deadline)
690 .map_or((None, None), |(deadline, kind)| {
691 (Some(deadline), Some(kind))
692 })
693}
694
695fn step_runtime_error_message(error: &VmError) -> String {
696 match error {
697 VmError::Thrown(VmValue::Dict(dict)) => dict
698 .get("message")
699 .map(|v| v.display())
700 .unwrap_or_else(|| error.to_string()),
701 _ => error.to_string(),
702 }
703}
704
705pub(crate) enum StepBoundaryOutcome {
706 Returned(VmValue),
707 Throw(VmError),
708}
709
710impl crate::vm::Vm {
711 pub(crate) async fn execute_one_cycle(&mut self) -> Result<Option<(VmValue, bool)>, VmError> {
712 if let Some(err) = self.pending_scope_interrupt().await {
713 match self.handle_error(err) {
714 Ok(None) => return Ok(None),
715 Ok(Some(val)) => return Ok(Some((val, false))),
716 Err(e) => return Err(e),
717 }
718 }
719
720 let frame_depth = self.frames.len();
721 let frame = match self.frames.last_mut() {
722 Some(f) => f,
723 None => {
724 let val = self.stack.pop().unwrap_or(VmValue::Nil);
725 return Ok(Some((val, false)));
726 }
727 };
728
729 if frame.ip >= frame.chunk.code.len() {
730 let val = self.stack.pop().unwrap_or(VmValue::Nil);
731 self.release_sync_guards_for_frame(self.frames.len());
732 let popped_frame = self.frames.pop().unwrap();
733 if self.frames.is_empty() {
734 return Ok(Some((val, false)));
735 }
736 self.iterators.truncate(popped_frame.saved_iterator_depth);
737 self.env = popped_frame.saved_env;
738 self.stack.truncate(popped_frame.stack_base);
739 self.stack.push(val);
740 return Ok(None);
741 }
742
743 let op_offset = frame.ip;
744 let op = frame.chunk.code[op_offset];
745 frame.ip += 1;
746
747 if let (Some(recorder), Some(decoded)) = (self.flight_recorder.as_ref(), Op::from_byte(op))
748 {
749 recorder.record_instruction(
750 &self.runtime_context.task_id,
751 frame_depth,
752 frame.fn_name.as_str(),
753 &frame.chunk,
754 self.source_file.as_deref(),
755 op_offset,
756 decoded,
757 );
758 }
759
760 match self.execute_op_with_scope_interrupts(op).await {
761 Ok(Some(val)) => Ok(Some((val, false))),
762 Ok(None) => Ok(None),
763 Err(VmError::Return(val)) => {
764 if let Some(popped_frame) = self.frames.pop() {
765 self.release_sync_guards_for_frame(self.frames.len() + 1);
766 if let Some(ref dir) = popped_frame.saved_source_dir {
767 crate::stdlib::set_thread_source_dir(dir);
768 }
769 let current_depth = self.frames.len();
770 self.exception_handlers
771 .retain(|h| h.frame_depth <= current_depth);
772 if self.frames.is_empty() {
773 return Ok(Some((val, false)));
774 }
775 self.iterators.truncate(popped_frame.saved_iterator_depth);
776 self.env = popped_frame.saved_env;
777 self.stack.truncate(popped_frame.stack_base);
778 self.stack.push(val);
779 Ok(None)
780 } else {
781 Ok(Some((val, false)))
782 }
783 }
784 Err(e) => {
785 if self.error_stack_trace.is_empty() {
786 self.error_stack_trace = self.capture_stack_trace();
787 }
788 match self.handle_error(e) {
789 Ok(None) => {
790 self.error_stack_trace.clear();
791 Ok(None)
792 }
793 Ok(Some(val)) => Ok(Some((val, false))),
794 Err(e) => Err(self.enrich_error_with_line(e)),
795 }
796 }
797 }
798 }
799
800 async fn execute_op_with_scope_interrupts(
801 &mut self,
802 op: u8,
803 ) -> Result<Option<VmValue>, VmError> {
804 #[cfg(test)]
805 SCOPE_INTERRUPT_ASYNC_DISPATCHES
806 .set(SCOPE_INTERRUPT_ASYNC_DISPATCHES.get().saturating_add(1));
807
808 enum ScopeInterruptResult {
809 Op(Result<Option<VmValue>, VmError>),
810 Deadline(DeadlineKind),
811 CancelTimedOut,
812 }
813
814 if matches!(
821 Op::from_byte(op),
822 Some(
823 Op::Import | Op::SelectiveImport | Op::NamespaceImport | Op::NamespaceImportMembers
824 )
825 ) {
826 return self.execute_op(op).await;
827 }
828
829 let execution_deadline = Arc::clone(&self.execution_deadline);
830 let scope_deadline = self.deadlines.last().map(|(deadline, _)| *deadline);
831 let interrupt_handler_deadline = self.interrupt_handler_deadline;
832 let cancel_token = self.cancel_token.clone();
833
834 let has_deadline = execution_deadline.is_active()
835 || scope_deadline.is_some()
836 || interrupt_handler_deadline.is_some();
837 if !has_deadline && cancel_token.is_none() {
838 return self.execute_op(op).await;
839 }
840
841 let cancel_requested_at_start = cancel_token
842 .as_ref()
843 .is_some_and(|token| token.load(std::sync::atomic::Ordering::SeqCst));
844 let has_cancel = cancel_token.is_some() && !cancel_requested_at_start;
845 let deadline_sleep = async move {
846 loop {
847 let changed = execution_deadline.changed();
850 let (deadline, kind) = next_deadline(
851 execution_deadline.current(),
852 scope_deadline,
853 interrupt_handler_deadline,
854 );
855 if let Some(deadline) = deadline {
856 tokio::select! {
857 _ = tokio::time::sleep_until(tokio::time::Instant::from_std(deadline)) => {
858 return kind.unwrap_or(DeadlineKind::Scope);
859 }
860 _ = changed => {}
861 }
862 } else {
863 changed.await;
864 }
865 }
866 };
867 let cancel_sleep = async move {
868 if let Some(token) = cancel_token {
869 while !token.load(std::sync::atomic::Ordering::SeqCst) {
870 tokio::time::sleep(Duration::from_millis(10)).await;
871 }
872 } else {
873 std::future::pending::<()>().await;
874 }
875 };
876
877 let result = {
878 let op_future = self.execute_op(op);
879 tokio::pin!(op_future);
880 tokio::select! {
881 result = &mut op_future => ScopeInterruptResult::Op(result),
882 kind = deadline_sleep, if has_deadline => {
883 ScopeInterruptResult::Deadline(kind)
884 },
885 _ = cancel_sleep, if has_cancel => {
886 let grace = tokio::time::sleep(CANCEL_GRACE_ASYNC_OP);
887 tokio::pin!(grace);
888 tokio::select! {
889 result = &mut op_future => ScopeInterruptResult::Op(result),
890 _ = &mut grace => ScopeInterruptResult::CancelTimedOut,
891 }
892 }
893 }
894 };
895
896 match result {
897 ScopeInterruptResult::Op(result) => result,
898 ScopeInterruptResult::Deadline(DeadlineKind::Execution) => {
899 self.cancel_spawned_tasks();
900 Err(VmError::ExecutionDeadlineExceeded)
901 }
902 ScopeInterruptResult::Deadline(DeadlineKind::Scope) => {
903 self.deadlines.pop();
904 self.cancel_spawned_tasks();
905 Err(Self::deadline_exceeded_error())
906 }
907 ScopeInterruptResult::Deadline(DeadlineKind::InterruptHandler) => {
908 Err(Self::interrupt_handler_timeout_error())
909 }
910 ScopeInterruptResult::CancelTimedOut => {
911 self.cancel_spawned_tasks();
912 let signal = self
913 .take_host_interrupt_signal()
914 .unwrap_or_else(|| "SIGINT".to_string());
915 if self.has_interrupt_handler_for(&signal) {
916 self.dispatch_interrupt_handlers(&signal).await?;
917 }
918 Err(Self::cancelled_error())
919 }
920 }
921 }
922
923 pub(crate) fn deadline_exceeded_error() -> VmError {
924 VmError::Thrown(VmValue::String(arcstr::ArcStr::from("Deadline exceeded")))
925 }
926
927 pub(crate) fn cancelled_error() -> VmError {
928 VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
929 "kind:cancelled:VM cancelled by host",
930 )))
931 }
932
933 pub(crate) fn capture_stack_trace(&self) -> Vec<(String, usize, usize, Option<String>)> {
935 self.frames
936 .iter()
937 .map(|f| {
938 let idx = if f.ip > 0 { f.ip - 1 } else { 0 };
939 let line = f.chunk.lines.get(idx).copied().unwrap_or(0) as usize;
940 let col = f.chunk.columns.get(idx).copied().unwrap_or(0) as usize;
941 (
942 f.fn_name.to_string(),
943 line,
944 col,
945 f.chunk.source_file.clone(),
946 )
947 })
948 .collect()
949 }
950
951 pub(crate) fn enrich_error_with_line(&self, error: VmError) -> VmError {
955 let (line, file) = self
960 .error_stack_trace
961 .last()
962 .map(|(_, l, _, f)| (*l, f.clone()))
963 .unwrap_or_else(|| (self.current_line(), None));
964 if line == 0 {
965 return error;
966 }
967 let suffix = match file.as_deref() {
968 Some(path) => {
969 let name = std::path::Path::new(path)
970 .file_name()
971 .and_then(|n| n.to_str())
972 .unwrap_or(path);
973 format!(" ({name}:{line})")
974 }
975 None => format!(" (line {line})"),
976 };
977 match error {
978 VmError::Runtime(msg) => VmError::Runtime(format!("{msg}{suffix}")),
979 VmError::TypeError(msg) => VmError::TypeError(format!("{msg}{suffix}")),
980 VmError::DivisionByZero => VmError::Runtime(format!("Division by zero{suffix}")),
981 VmError::UndefinedVariable(name) => {
982 VmError::Runtime(format!("Undefined variable: {name}{suffix}"))
983 }
984 VmError::UndefinedBuiltin(name) => {
985 VmError::Runtime(format!("Undefined builtin: {name}{suffix}"))
986 }
987 VmError::ImmutableAssignment(name) => VmError::Runtime(format!(
988 "Cannot assign to immutable binding: {name}{suffix}"
989 )),
990 VmError::StackOverflow => {
991 VmError::Runtime(format!("Stack overflow: too many nested calls{suffix}"))
992 }
993 other => other,
999 }
1000 }
1001}
1002
1003#[cfg(test)]
1004mod tests {
1005 use super::*;
1006 use crate::compiler::Compiler;
1007 use crate::stdlib::register_vm_stdlib;
1008 use harn_lexer::Lexer;
1009 use harn_parser::Parser;
1010 use std::sync::atomic::{AtomicBool, Ordering};
1011
1012 fn compile_harn(source: &str) -> Chunk {
1013 let mut lexer = Lexer::new(source);
1014 let tokens = lexer.tokenize().unwrap();
1015 let mut parser = Parser::new(tokens);
1016 let program = parser.parse().unwrap();
1017 Compiler::new().compile(&program).unwrap()
1018 }
1019
1020 #[tokio::test(flavor = "current_thread")]
1021 async fn dropping_timed_execution_restores_ambient_state_and_poison_vm_reuse() {
1022 let local = tokio::task::LocalSet::new();
1023 local
1024 .run_until(async {
1025 crate::reset_thread_local_state();
1026 let baseline_dir = tempfile::tempdir().unwrap();
1027 let imported_dir = tempfile::tempdir().unwrap();
1028 let poisoned_dir = tempfile::tempdir().unwrap();
1029 let quick = compile_harn("pipeline default(harness: Harness) { return 42 }");
1030 let mut vm = Vm::new();
1031 register_vm_stdlib(&mut vm);
1032 crate::tracing::set_tracing_enabled(true);
1033
1034 let child_started = Arc::new(AtomicBool::new(false));
1035 let child_effect = Arc::new(AtomicBool::new(false));
1036 let child_release = Arc::new(tokio::sync::Notify::new());
1037 let (ambient_started_tx, ambient_started_rx) = tokio::sync::oneshot::channel();
1038 let ambient_started_tx = Arc::new(std::sync::Mutex::new(Some(ambient_started_tx)));
1039 let started_for_builtin = Arc::clone(&child_started);
1040 vm.register_builtin("child_started", move |_args, _output| {
1041 started_for_builtin.store(true, Ordering::Release);
1042 Ok(VmValue::Nil)
1043 });
1044 let effect_for_builtin = Arc::clone(&child_effect);
1045 vm.register_builtin("child_effect", move |_args, _output| {
1046 effect_for_builtin.store(true, Ordering::Release);
1047 Ok(VmValue::Nil)
1048 });
1049 let release_for_builtin = Arc::clone(&child_release);
1050 vm.register_async_builtin("wait_for_child_release", move |_ctx, _args| {
1051 let release = Arc::clone(&release_for_builtin);
1052 async move {
1053 release.notified().await;
1054 Ok(VmValue::Nil)
1055 }
1056 });
1057 vm.register_async_builtin("wait_forever", |_ctx, _args| async move {
1058 std::future::pending::<()>().await;
1059 Ok(VmValue::Nil)
1060 });
1061
1062 let imported_path = imported_dir.path().join("cancelled.harn");
1063 let imported_source_dir = imported_dir.path().to_path_buf();
1064 let poisoned_source_dir = poisoned_dir.path().to_path_buf();
1065 let ambient_started_for_builtin = Arc::clone(&ambient_started_tx);
1066 vm.register_builtin("strand_ambient_state", move |_args, _output| {
1067 crate::step_runtime::register_persona(
1068 "cancellation_entry",
1069 crate::step_runtime::PersonaDefinition {
1070 name: "cancel_persona".into(),
1071 stages: vec![crate::personas::StageDecl {
1072 name: "cancel_step".into(),
1073 allowed_tools: Some(vec!["cancel_tool".into()]),
1074 ..Default::default()
1075 }],
1076 ..Default::default()
1077 },
1078 );
1079 crate::step_runtime::register_step(
1080 "cancelled_step",
1081 crate::step_runtime::StepDefinition {
1082 name: "cancel_step".into(),
1083 function: "cancelled_step".into(),
1084 model: Some("cancel-model".into()),
1085 ..Default::default()
1086 },
1087 );
1088 assert!(crate::step_runtime::maybe_push_active_persona(
1089 "cancellation_entry",
1090 1
1091 ));
1092 assert!(crate::step_runtime::maybe_push_active_step(
1093 "cancelled_step",
1094 2,
1095 &[]
1096 ));
1097 assert_eq!(
1098 crate::stdlib::process::source_root_path(),
1099 imported_source_dir
1100 );
1101 assert_eq!(
1102 crate::step_runtime::current_persona_name().as_deref(),
1103 Some("cancel_persona")
1104 );
1105 assert_eq!(
1106 crate::step_runtime::active_step_model_default().as_deref(),
1107 Some("cancel-model")
1108 );
1109 assert_eq!(
1110 crate::orchestration::current_execution_policy()
1111 .unwrap()
1112 .tools,
1113 vec!["cancel_tool"]
1114 );
1115 crate::stdlib::process::set_thread_source_dir(&poisoned_source_dir);
1116 crate::stdlib::process::set_thread_execution_context(Some(
1117 crate::orchestration::RunExecutionRecord {
1118 adapter: Some("cancelled".into()),
1119 ..Default::default()
1120 },
1121 ));
1122 crate::orchestration::push_approval_policy(
1123 crate::orchestration::ToolApprovalPolicy {
1124 auto_deny: vec!["cancel_tool".into()],
1125 ..Default::default()
1126 },
1127 );
1128 if let Some(sender) = ambient_started_for_builtin.lock().unwrap().take() {
1129 let _ = sender.send(());
1130 }
1131 Ok(VmValue::Nil)
1132 });
1133
1134 std::fs::write(
1135 &imported_path,
1136 r#"
1137@persona(name: "cancel_persona", stages: [{name: "cancel_step", allowed_tools: ["cancel_tool"]}])
1138pub fn cancellation_entry(agent: HarnessAgent) {
1139 return cancelled_step(agent)
1140}
1141
1142@step(name: "cancel_step", model: "cancel-model")
1143fn cancelled_step(agent: HarnessAgent) {
1144 agent.pipeline_on_finish({ _h, value -> value })
1145 const child = spawn {
1146 child_started()
1147 wait_for_child_release()
1148 child_effect()
1149 }
1150 strand_ambient_state()
1151 wait_forever()
1152}
1153"#,
1154 )
1155 .unwrap();
1156 let mut helper_exports = vm
1157 .load_module_exports_from_source(
1158 "<cancellation-helper>",
1159 "pub fn outer_callback(_h, value) { return value }\n\
1160 pub fn answer() { return 42 }",
1161 )
1162 .await
1163 .unwrap();
1164 let callable = helper_exports.remove("answer").unwrap();
1165 let outer_callback = helper_exports.remove("outer_callback").unwrap();
1166 let slow = compile_harn(&format!(
1167 "import {{ cancellation_entry }} from \"{}\"\n\
1168 pipeline default(harness: Harness) {{ return cancellation_entry(harness.agent) }}",
1169 imported_path.display()
1170 ));
1171
1172 let baseline_execution = crate::orchestration::RunExecutionRecord {
1173 cwd: Some(baseline_dir.path().display().to_string()),
1174 source_dir: Some(baseline_dir.path().display().to_string()),
1175 adapter: Some("baseline".into()),
1176 ..Default::default()
1177 };
1178 let baseline_policy = crate::orchestration::CapabilityPolicy {
1179 tools: vec!["baseline_tool".into(), "cancel_tool".into()],
1180 ..Default::default()
1181 };
1182 let baseline_approval = crate::orchestration::ToolApprovalPolicy {
1183 auto_approve: vec!["baseline_tool".into()],
1184 ..Default::default()
1185 };
1186 crate::stdlib::process::set_thread_source_dir(baseline_dir.path());
1187 crate::stdlib::process::set_thread_execution_context(Some(
1188 baseline_execution.clone(),
1189 ));
1190 crate::orchestration::push_execution_policy(baseline_policy.clone());
1191 crate::orchestration::push_approval_policy(baseline_approval.clone());
1192 crate::orchestration::set_pipeline_on_finish(Arc::clone(&outer_callback));
1193 let outer_span =
1194 crate::tracing::span_start(crate::tracing::SpanKind::Pipeline, "outer".into());
1195
1196 let mut execution =
1197 Box::pin(vm.execute_with_timeout(&slow, Duration::from_secs(30)));
1198 tokio::select! {
1199 biased;
1200 result = &mut execution => panic!("slow execution unexpectedly finished: {result:?}"),
1201 started = async {
1202 ambient_started_rx.await.expect("step did not reach cancellation point");
1203 while !child_started.load(Ordering::Acquire) {
1204 tokio::task::yield_now().await;
1205 }
1206 } => started,
1207 }
1208 drop(execution);
1209
1210 assert!(!vm.execution_deadline.is_active());
1211 assert!(!vm.frames.is_empty(), "fixture must abandon a live frame");
1212 assert_eq!(
1213 crate::stdlib::process::source_root_path(),
1214 baseline_dir.path()
1215 );
1216 assert_eq!(
1217 crate::stdlib::process::current_execution_context(),
1218 Some(baseline_execution.clone())
1219 );
1220 assert_eq!(
1221 crate::orchestration::current_execution_policy(),
1222 Some(baseline_policy.clone())
1223 );
1224 assert_eq!(
1225 crate::orchestration::current_approval_policy(),
1226 Some(baseline_approval.clone())
1227 );
1228 assert!(crate::step_runtime::current_persona_name().is_none());
1229 assert!(crate::step_runtime::active_step_model_default().is_none());
1230 assert_eq!(crate::tracing::current_span_id(), Some(outer_span));
1231 let restored_callback =
1232 crate::orchestration::take_pipeline_on_finish().unwrap();
1233 assert!(Arc::ptr_eq(&restored_callback, &outer_callback));
1234 crate::orchestration::set_pipeline_on_finish(restored_callback);
1235 let abandoned_spans = crate::tracing::peek_spans();
1236 assert!(abandoned_spans.iter().any(|span| {
1237 span.name == "cancel_step"
1238 && span.metadata.get("status") == Some(&serde_json::json!("abandoned"))
1239 }));
1240 assert!(abandoned_spans.iter().any(|span| {
1241 span.name == "main"
1242 && span.metadata.get("status") == Some(&serde_json::json!("abandoned"))
1243 }));
1244
1245 let frame_depth = vm.frames.len();
1246 let output = vm.output().to_string();
1247 let error = vm.execute(&quick).await.unwrap_err();
1248 assert!(matches!(error, VmError::AbandonedExecution));
1249 let closure_error = vm.call_closure_pub(&callable, &[]).await.unwrap_err();
1250 assert!(matches!(closure_error, VmError::AbandonedExecution));
1251 let source_cache_len = vm.source_cache.len();
1252 let module_cache_len = vm.module_cache.len();
1253 let module_error = vm
1254 .load_module_exports_from_source(
1255 "<poisoned-module-load>",
1256 "pub fn poisoned() { return 0 }",
1257 )
1258 .await
1259 .unwrap_err();
1260 assert!(matches!(module_error, VmError::AbandonedExecution));
1261 assert_eq!(vm.source_cache.len(), source_cache_len);
1262 assert_eq!(vm.module_cache.len(), module_cache_len);
1263 let start_error = vm.start(&quick).unwrap_err();
1264 assert!(matches!(start_error, VmError::AbandonedExecution));
1265 let restart_error = vm.restart_frame(0).unwrap_err();
1266 assert!(matches!(restart_error, VmError::AbandonedExecution));
1267 assert_eq!(vm.frames.len(), frame_depth);
1268 assert_eq!(vm.output(), output);
1269
1270 let mut vm_b = Vm::new();
1271 register_vm_stdlib(&mut vm_b);
1272 let observed_callback = Arc::clone(&outer_callback);
1273 let observed_dir = baseline_dir.path().to_path_buf();
1274 vm_b.register_builtin("observe_baseline", move |_args, _output| {
1275 assert_eq!(crate::stdlib::process::source_root_path(), observed_dir);
1276 assert_eq!(
1277 crate::stdlib::process::current_execution_context(),
1278 Some(baseline_execution.clone())
1279 );
1280 assert_eq!(
1281 crate::orchestration::current_execution_policy(),
1282 Some(baseline_policy.clone())
1283 );
1284 assert_eq!(
1285 crate::orchestration::current_approval_policy(),
1286 Some(baseline_approval.clone())
1287 );
1288 assert!(crate::step_runtime::current_persona_name().is_none());
1289 assert!(crate::step_runtime::active_step_model_default().is_none());
1290 let callback = crate::orchestration::take_pipeline_on_finish().unwrap();
1291 assert!(Arc::ptr_eq(&callback, &observed_callback));
1292 crate::orchestration::set_pipeline_on_finish(callback);
1293 Ok(VmValue::Nil)
1294 });
1295 let vm_b_chunk =
1296 compile_harn("pipeline default(harness: Harness) { observe_baseline(); return 42 }");
1297 assert!(matches!(
1298 vm_b.execute(&vm_b_chunk).await.unwrap(),
1299 VmValue::Int(42)
1300 ));
1301 assert_eq!(crate::tracing::current_span_id(), Some(outer_span));
1302
1303 drop(vm);
1304 child_release.notify_one();
1305 for _ in 0..10 {
1306 tokio::task::yield_now().await;
1307 }
1308 assert!(
1309 !child_effect.load(Ordering::Acquire),
1310 "dropping an abandoned VM must abort spawned side effects"
1311 );
1312 crate::reset_thread_local_state();
1313 })
1314 .await;
1315 }
1316
1317 #[tokio::test(flavor = "current_thread")]
1318 async fn natural_host_deadline_is_terminal_not_abandoned() {
1319 let local = tokio::task::LocalSet::new();
1320 local
1321 .run_until(async {
1322 let infinite = compile_harn("pipeline default(harness: Harness) { while true {} }");
1323 let quick = compile_harn("pipeline default(harness: Harness) { return 42 }");
1324 let mut vm = Vm::new();
1325 register_vm_stdlib(&mut vm);
1326
1327 let error = vm
1328 .execute_with_timeout(&infinite, Duration::ZERO)
1329 .await
1330 .unwrap_err();
1331 assert!(matches!(error, VmError::ExecutionDeadlineExceeded));
1332 assert!(!vm.execution_deadline.is_abandoned());
1333 assert!(matches!(
1334 vm.execute(&quick).await.unwrap(),
1335 VmValue::Int(42)
1336 ));
1337 })
1338 .await;
1339 }
1340
1341 #[tokio::test(flavor = "current_thread")]
1342 async fn host_admission_extends_execution_deadline_by_injected_clock_delta() {
1343 let chunk = compile_harn(
1344 r"
1345pipeline default(harness: Harness) {
1346 wait_for_admission()
1347 return 42
1348}
1349",
1350 );
1351 let mut vm = Vm::new();
1352 register_vm_stdlib(&mut vm);
1353 let clock = harn_clock::PausedClock::new(time::OffsetDateTime::UNIX_EPOCH);
1354 let admission_clock: Arc<dyn harn_clock::Clock> = clock.clone();
1355 vm.register_async_builtin("wait_for_admission", move |ctx, _args| {
1356 let clock = Arc::clone(&clock);
1357 let admission_clock = Arc::clone(&admission_clock);
1358 async move {
1359 let before = ctx.execution_deadline_offset_for_test();
1360 let pause = ctx
1361 .pause_execution_deadline(admission_clock)
1362 .expect("timed execution exposes its outer deadline to inline host work");
1363 clock.advance(Duration::from_millis(25));
1364 drop(pause);
1365 let after = ctx.execution_deadline_offset_for_test();
1366 assert_eq!(after.saturating_sub(before), 25_000_000);
1367 Ok(VmValue::Nil)
1368 }
1369 });
1370
1371 let value = vm
1372 .execute_with_timeout(&chunk, Duration::from_secs(5))
1373 .await
1374 .expect("host admission extends rather than spends the execution budget");
1375 assert!(matches!(value, VmValue::Int(42)));
1376 }
1377
1378 #[tokio::test(flavor = "current_thread")]
1379 async fn timed_finite_loop_keeps_sync_opcodes_on_direct_dispatch() {
1380 let local = tokio::task::LocalSet::new();
1381 local
1382 .run_until(async {
1383 let chunk = compile_harn(
1384 r"
1385pipeline default(harness: Harness) {
1386 let total = 0
1387 for i in 0 to 10000 {
1388 total = total + i
1389 }
1390 return total
1391}
1392",
1393 );
1394 let mut vm = Vm::new();
1395 register_vm_stdlib(&mut vm);
1396 reset_scope_interrupt_async_dispatches();
1397
1398 let value = vm
1399 .execute_with_timeout(&chunk, Duration::from_secs(1))
1400 .await
1401 .unwrap();
1402
1403 assert!(matches!(value, VmValue::Int(_)));
1404 assert!(
1405 scope_interrupt_async_dispatches() <= 4,
1406 "finite sync loop fell back to per-op async dispatch"
1407 );
1408 })
1409 .await;
1410 }
1411}