1use crate::value::VmDictExt;
2use std::cell::RefCell;
3use std::collections::{BTreeMap, BTreeSet};
4use std::path::Path;
5use std::sync::Arc;
6use std::time::Duration as StdDuration;
7
8use serde::{Deserialize, Serialize};
9use serde_json::{json, Value as JsonValue};
10use sha2::Digest;
11use time::OffsetDateTime;
12use uuid::Uuid;
13
14use crate::event_log::{
15 active_event_log, install_default_for_base_dir, install_memory_for_current_thread, AnyEventLog,
16 EventLog, LogEvent, Topic,
17};
18use crate::runtime_limits::RuntimeLimits;
19use crate::schema::schema_expect_value;
20use crate::stdlib::args::{ArgError, Args, ErrorKind};
21use crate::stdlib::macros::{harn_builtin, BuiltinSignature, Param, VmBuiltinDef, TY_ANY, TY_DICT};
22use crate::stdlib::waitpoint::{
23 cancel_waitpoint_on, complete_waitpoint_on, create_waitpoint_on, inspect_waitpoint_on,
24 wait_on_waitpoints, WaitpointRecord, WaitpointStatus, WaitpointWaitFailure,
25 WaitpointWaitOptions,
26};
27use crate::triggers::dispatcher::current_dispatch_context;
28use crate::value::{categorized_error, ErrorCategory, VmError, VmValue};
29use crate::vm::{AsyncBuiltinCtx, Vm};
30
31mod fixtures;
32use fixtures::{maybe_apply_mock_response, maybe_apply_mock_response_with_harness};
33
34const HITL_EVENT_LOG_QUEUE_DEPTH: usize = RuntimeLimits::DEFAULT.default_event_log_queue_depth;
35const HITL_APPROVAL_TIMEOUT_MS: u64 = 24 * 60 * 60 * 1000;
36const HITL_QUESTION_TIMEOUT_MS: u64 = 24 * 60 * 60 * 1000;
37
38pub const HITL_QUESTIONS_TOPIC: &str = "hitl.questions";
39pub const HITL_APPROVALS_TOPIC: &str = "hitl.approvals";
40pub const HITL_DUAL_CONTROL_TOPIC: &str = "hitl.dual_control";
41pub const HITL_ESCALATIONS_TOPIC: &str = "hitl.escalations";
42
43thread_local! {
44 static REQUEST_SEQUENCE: RefCell<RequestSequenceState> = RefCell::new(RequestSequenceState::default());
45}
46
47#[derive(Default)]
48pub(crate) struct RequestSequenceState {
49 pub(crate) instance_key: String,
50 pub(crate) next_seq: u64,
51}
52
53#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
54#[serde(rename_all = "snake_case")]
55pub enum HitlRequestKind {
56 Question,
57 Approval,
58 DualControl,
59 Escalation,
60}
61
62impl HitlRequestKind {
63 pub(crate) fn as_str(self) -> &'static str {
64 match self {
65 Self::Question => "question",
66 Self::Approval => "approval",
67 Self::DualControl => "dual_control",
68 Self::Escalation => "escalation",
69 }
70 }
71
72 fn topic(self) -> &'static str {
73 match self {
74 Self::Question => HITL_QUESTIONS_TOPIC,
75 Self::Approval => HITL_APPROVALS_TOPIC,
76 Self::DualControl => HITL_DUAL_CONTROL_TOPIC,
77 Self::Escalation => HITL_ESCALATIONS_TOPIC,
78 }
79 }
80
81 fn request_event_kind(self) -> &'static str {
82 match self {
83 Self::Question => "hitl.question_asked",
84 Self::Approval => "hitl.approval_requested",
85 Self::DualControl => "hitl.dual_control_requested",
86 Self::Escalation => "hitl.escalation_issued",
87 }
88 }
89
90 pub(crate) fn from_request_id(request_id: &str) -> Option<Self> {
91 if request_id.starts_with("hitl_question_") {
92 Some(Self::Question)
93 } else if request_id.starts_with("hitl_approval_") {
94 Some(Self::Approval)
95 } else if request_id.starts_with("hitl_dual_control_") {
96 Some(Self::DualControl)
97 } else if request_id.starts_with("hitl_escalation_") {
98 Some(Self::Escalation)
99 } else {
100 None
101 }
102 }
103}
104
105#[derive(Clone, Debug, Serialize, Deserialize)]
106pub struct HitlHostResponse {
107 pub request_id: String,
108 #[serde(skip_serializing_if = "Option::is_none")]
109 pub answer: Option<JsonValue>,
110 #[serde(skip_serializing_if = "Option::is_none")]
111 pub approved: Option<bool>,
112 #[serde(skip_serializing_if = "Option::is_none")]
113 pub accepted: Option<bool>,
114 #[serde(skip_serializing_if = "Option::is_none")]
115 pub reviewer: Option<String>,
116 #[serde(skip_serializing_if = "Option::is_none")]
117 pub reason: Option<String>,
118 #[serde(skip_serializing_if = "Option::is_none")]
119 pub metadata: Option<JsonValue>,
120 #[serde(skip_serializing_if = "Option::is_none")]
121 pub responded_at: Option<String>,
122 #[serde(skip_serializing_if = "Option::is_none")]
123 pub signature: Option<String>,
124}
125
126#[derive(Clone, Debug, Serialize, Deserialize)]
127struct HitlRequestEnvelope {
128 request_id: String,
129 kind: HitlRequestKind,
130 #[serde(default)]
131 agent: String,
132 trace_id: String,
133 #[serde(skip_serializing_if = "Option::is_none")]
134 run_id: Option<String>,
135 requested_at: String,
136 payload: JsonValue,
137}
138
139#[derive(Clone, Debug, Serialize, Deserialize)]
140struct HitlTimeoutRecord {
141 request_id: String,
142 kind: HitlRequestKind,
143 trace_id: String,
144 timed_out_at: String,
145}
146
147#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
148pub struct ApprovalRequest {
149 pub id: String,
150 pub action: String,
151 #[serde(default)]
152 pub args: JsonValue,
153 pub principal: String,
154 pub requested_at: String,
155 #[serde(skip_serializing_if = "Option::is_none")]
156 pub deadline: Option<String>,
157 pub approvers_required: u32,
158 #[serde(default)]
159 pub evidence_refs: Vec<JsonValue>,
160 #[serde(default)]
161 pub undo_metadata: JsonValue,
162 #[serde(default)]
163 pub capabilities_requested: Vec<String>,
164}
165
166impl ApprovalRequest {
167 pub fn new(
168 id: impl Into<String>,
169 action: impl Into<String>,
170 args: JsonValue,
171 principal: impl Into<String>,
172 requested_at: impl Into<String>,
173 ) -> Self {
174 Self {
175 id: id.into(),
176 action: action.into(),
177 args,
178 principal: principal.into(),
179 requested_at: requested_at.into(),
180 deadline: None,
181 approvers_required: 1,
182 evidence_refs: Vec::new(),
183 undo_metadata: JsonValue::Null,
184 capabilities_requested: Vec::new(),
185 }
186 }
187}
188
189pub(crate) fn approval_request_for_host_permission(
190 id: impl Into<String>,
191 action: impl Into<String>,
192 args: JsonValue,
193 principal: impl Into<String>,
194 evidence_refs: Vec<JsonValue>,
195 undo_metadata: JsonValue,
196 capabilities_requested: Vec<String>,
197) -> ApprovalRequest {
198 let mut request = ApprovalRequest::new(id, action, args, principal, now_rfc3339());
199 request.evidence_refs = evidence_refs;
200 request.undo_metadata = undo_metadata;
201 request.capabilities_requested = capabilities_requested;
202 request
203}
204
205#[derive(Clone, Debug)]
206struct DispatchKeys {
207 instance_key: String,
208 stable_base: String,
209 agent: String,
210 trace_id: String,
211}
212
213#[derive(Clone, Debug)]
214struct AskUserOptions {
215 schema: Option<VmValue>,
216 timeout: Option<StdDuration>,
217 default: Option<VmValue>,
218}
219
220#[derive(Clone, Debug)]
221struct ApprovalOptions {
222 detail: Option<VmValue>,
223 args: Option<VmValue>,
224 quorum: u32,
225 reviewers: Vec<String>,
226 deadline: StdDuration,
227 principal: Option<String>,
228 evidence_refs: Vec<JsonValue>,
229 undo_metadata: Option<JsonValue>,
230 capabilities_requested: Vec<String>,
231}
232
233#[derive(Clone, Debug)]
234struct ApprovalProgress {
235 request_id: String,
236 reviewers: BTreeSet<String>,
237 signatures: Vec<ApprovalSignature>,
238 reason: Option<String>,
239 approved_at: Option<String>,
240}
241
242#[derive(Clone, Debug, Serialize)]
243struct ApprovalSignature {
244 reviewer: String,
245 signed_at: String,
246 signature: String,
247}
248
249#[derive(Clone, Debug)]
250enum ApprovalResolution {
251 Pending,
252 Approved(ApprovalProgress),
253 Denied(HitlHostResponse),
254}
255
256#[allow(clippy::large_enum_variant)]
264#[derive(Clone, Debug)]
265enum WaitpointOutcome {
266 Completed(WaitpointRecord),
267 Timeout,
268 Cancelled {
269 wait_id: String,
270 waitpoint_ids: Vec<String>,
271 reason: Option<String>,
272 },
273}
274
275pub(crate) fn register_hitl_builtins(vm: &mut Vm) {
276 for def in MODULE_BUILTINS {
277 vm.register_builtin_def(def);
278 }
279}
280
281pub(crate) const MODULE_BUILTINS: &[&VmBuiltinDef] = &[
282 &ASK_USER_BUILTIN_DEF,
283 &REQUEST_APPROVAL_BUILTIN_DEF,
284 &DUAL_CONTROL_BUILTIN_DEF,
285 &ESCALATE_TO_BUILTIN_DEF,
286];
287
288#[harn_builtin(
289 exposure = "runtime_internal",
290 effects = [],
291 sig = "ask_user(prompt: string, options?: dict) -> any",
292 kind = "async",
293 category = "hitl"
294)]
295async fn ask_user_builtin(
296 ctx: crate::vm::AsyncBuiltinCtx,
297 args: Vec<VmValue>,
298) -> Result<VmValue, VmError> {
299 ask_user_impl(Some(&ctx), &args).await
300}
301
302#[harn_builtin(
303 exposure = "runtime_internal",
304 effects = [],
305 sig_expr = BuiltinSignature::variadic("request_approval", &[Param::new("args", TY_ANY
306)], TY_DICT),
307 kind = "async",
308 category = "hitl"
309)]
310async fn request_approval_builtin(
311 ctx: crate::vm::AsyncBuiltinCtx,
312 args: Vec<VmValue>,
313) -> Result<VmValue, VmError> {
314 request_approval_impl(Some(&ctx), None, &args).await
315}
316
317#[harn_builtin(
318 exposure = "runtime_internal",
319 effects = [],
320 sig = "dual_control(n: int, m: int, action: closure, approvers?: list) -> dict",
321 kind = "async",
322 category = "hitl"
323)]
324async fn dual_control_builtin(
325 ctx: crate::vm::AsyncBuiltinCtx,
326 args: Vec<VmValue>,
327) -> Result<VmValue, VmError> {
328 dual_control_impl(&ctx, &args).await
329}
330
331#[harn_builtin(
332 exposure = "runtime_internal",
333 effects = [],
334 sig = "escalate_to(role: string, reason: string) -> dict",
335 kind = "async",
336 category = "hitl"
337)]
338async fn escalate_to_builtin(
339 ctx: crate::vm::AsyncBuiltinCtx,
340 args: Vec<VmValue>,
341) -> Result<VmValue, VmError> {
342 escalate_to_impl(Some(&ctx), &args).await
343}
344
345pub(crate) fn reset_hitl_state() {
346 REQUEST_SEQUENCE.with(|slot| {
347 *slot.borrow_mut() = RequestSequenceState::default();
348 });
349}
350
351pub(crate) fn take_hitl_state() -> RequestSequenceState {
352 REQUEST_SEQUENCE.with(|slot| std::mem::take(&mut *slot.borrow_mut()))
353}
354
355pub(crate) fn restore_hitl_state(state: RequestSequenceState) {
356 REQUEST_SEQUENCE.with(|slot| {
357 *slot.borrow_mut() = state;
358 });
359}
360
361pub async fn append_hitl_response(
362 base_dir: Option<&Path>,
363 mut response: HitlHostResponse,
364) -> Result<u64, String> {
365 let kind = HitlRequestKind::from_request_id(&response.request_id)
366 .ok_or_else(|| format!("unknown HITL request id '{}'", response.request_id))?;
367 if response.responded_at.is_none() {
368 response.responded_at = Some(now_rfc3339());
369 }
370 let log = ensure_hitl_event_log_for(base_dir)?;
371 let headers = response_headers(&response.request_id);
372 let topic = Topic::new(kind.topic()).map_err(|error| error.to_string())?;
373 let event_id = log
374 .append(
375 &topic,
376 LogEvent::new(
377 match kind {
378 HitlRequestKind::Escalation => "hitl.escalation_accepted",
379 _ => "hitl.response_received",
380 },
381 serde_json::to_value(&response).map_err(|error| error.to_string())?,
382 )
383 .with_headers(headers),
384 )
385 .await
386 .map_err(|error| error.to_string())?;
387 finalize_hitl_response(&log, kind, &response).await?;
388 Ok(event_id)
389}
390
391pub async fn append_approval_request_on(
392 log: &Arc<AnyEventLog>,
393 agent: impl Into<String>,
394 trace_id: impl Into<String>,
395 action: impl Into<String>,
396 detail: JsonValue,
397 reviewers: Vec<String>,
398) -> Result<String, VmError> {
399 let request_id = next_request_id(HitlRequestKind::Approval, current_dispatch_keys().as_ref());
400 let trace_id = trace_id.into();
401 let agent = agent.into();
402 let requested_at_time = OffsetDateTime::now_utc();
403 let requested_at = format_rfc3339(requested_at_time);
404 let mut approval_request = ApprovalRequest::new(
405 request_id.clone(),
406 action.into(),
407 detail.clone(),
408 agent.clone(),
409 requested_at.clone(),
410 );
411 approval_request.deadline = deadline_after(
412 requested_at_time,
413 StdDuration::from_millis(HITL_APPROVAL_TIMEOUT_MS),
414 );
415 approval_request.approvers_required = 1;
416 let approval_request_json = serde_json::to_value(&approval_request)
417 .map_err(|error| VmError::Runtime(error.to_string()))?;
418 let request = HitlRequestEnvelope {
419 request_id: request_id.clone(),
420 kind: HitlRequestKind::Approval,
421 agent,
422 trace_id: trace_id.clone(),
423 run_id: None,
424 requested_at: requested_at.clone(),
425 payload: json!({
426 "approval_request": approval_request_json,
427 "id": approval_request.id,
428 "action": approval_request.action,
429 "args": approval_request.args,
430 "principal": approval_request.principal,
431 "requested_at": requested_at,
432 "deadline": approval_request.deadline,
433 "approvers_required": approval_request.approvers_required,
434 "evidence_refs": approval_request.evidence_refs,
435 "undo_metadata": approval_request.undo_metadata,
436 "capabilities_requested": approval_request.capabilities_requested,
437 "detail": detail,
438 "quorum": 1,
439 "reviewers": reviewers,
440 "deadline_ms": HITL_APPROVAL_TIMEOUT_MS,
441 }),
442 };
443 create_request_waitpoint(log, &request).await?;
444 append_request(log, &request).await?;
445 maybe_notify_host(None, &request);
446 Ok(request_id)
447}
448
449async fn ask_user_impl(
450 ctx: Option<&AsyncBuiltinCtx>,
451 args: &[VmValue],
452) -> Result<VmValue, VmError> {
453 let prompt = Args::runtime("ask_user", args)
454 .non_empty_string(0, "prompt")?
455 .to_string();
456 let options = parse_ask_user_options(args.get(1))?;
457 let keys = current_dispatch_keys();
458 let request_id = next_request_id(HitlRequestKind::Question, keys.as_ref());
459 let trace_id = keys
460 .as_ref()
461 .map(|keys| keys.trace_id.clone())
462 .unwrap_or_else(new_trace_id);
463 let log = ensure_hitl_event_log();
464 let request = HitlRequestEnvelope {
465 request_id: request_id.clone(),
466 kind: HitlRequestKind::Question,
467 agent: keys
468 .as_ref()
469 .map(|keys| keys.agent.clone())
470 .unwrap_or_default(),
471 trace_id: trace_id.clone(),
472 run_id: crate::orchestration::current_mutation_session().and_then(|session| session.run_id),
473 requested_at: now_rfc3339(),
474 payload: json!({
475 "prompt": prompt,
476 "schema": options.schema.as_ref().map(crate::llm::vm_value_to_json),
477 "default": options.default.as_ref().map(crate::llm::vm_value_to_json),
478 "timeout_ms": options.timeout.map(|timeout| timeout.as_millis() as u64),
479 }),
480 };
481 create_request_waitpoint(&log, &request).await?;
482 append_request(&log, &request).await?;
483 maybe_notify_host(ctx, &request);
484 emit_hitl_requested(&request);
485 maybe_apply_mock_response(
486 ctx,
487 HitlRequestKind::Question,
488 &request_id,
489 &request.payload,
490 )
491 .await?;
492
493 match wait_for_request_waitpoint_with_events(
494 &request_id,
495 HitlRequestKind::Question,
496 options.timeout,
497 )
498 .await?
499 {
500 WaitpointOutcome::Completed(record) => {
501 let answer = record
502 .value
503 .as_ref()
504 .map(crate::stdlib::json_to_vm_value)
505 .unwrap_or(VmValue::Nil);
506 if let Some(schema) = options.schema.as_ref() {
507 return schema_expect_value(&answer, schema, true);
508 }
509 if let Some(default) = options.default.as_ref() {
510 return Ok(coerce_like_default(&answer, default));
511 }
512 Ok(answer)
513 }
514 WaitpointOutcome::Timeout => {
515 append_timeout_once(&log, HitlRequestKind::Question, &request_id, &trace_id).await?;
516 if let Some(default) = options.default {
517 return Ok(default);
518 }
519 Err(timeout_error(&request_id, HitlRequestKind::Question))
520 }
521 WaitpointOutcome::Cancelled {
522 wait_id,
523 waitpoint_ids,
524 reason,
525 } => Err(hitl_cancelled_error(
526 &request_id,
527 HitlRequestKind::Question,
528 &wait_id,
529 &waitpoint_ids,
530 reason,
531 )),
532 }
533}
534
535async fn request_approval_impl(
536 ctx: Option<&AsyncBuiltinCtx>,
537 harness: Option<&crate::harness::VmHarness>,
538 args: &[VmValue],
539) -> Result<VmValue, VmError> {
540 let action = Args::runtime("request_approval", args)
541 .non_empty_string(0, "action")?
542 .to_string();
543 let options = parse_approval_options(args.get(1), "request_approval")?;
544 let keys = current_dispatch_keys();
545 let request_id = next_request_id(HitlRequestKind::Approval, keys.as_ref());
546 let trace_id = keys
547 .as_ref()
548 .map(|keys| keys.trace_id.clone())
549 .unwrap_or_else(new_trace_id);
550 let agent = keys
551 .as_ref()
552 .map(|keys| keys.agent.clone())
553 .unwrap_or_default();
554 let requested_at_time = OffsetDateTime::now_utc();
555 let requested_at = format_rfc3339(requested_at_time);
556 let principal = options.principal.clone().unwrap_or_else(|| agent.clone());
557 let approval_args = options
558 .args
559 .as_ref()
560 .or(options.detail.as_ref())
561 .map(crate::llm::vm_value_to_json)
562 .unwrap_or(JsonValue::Null);
563 let mut approval_request = ApprovalRequest::new(
564 request_id.clone(),
565 action.clone(),
566 approval_args,
567 principal,
568 requested_at.clone(),
569 );
570 approval_request.deadline = deadline_after(requested_at_time, options.deadline);
571 approval_request.approvers_required = options.quorum;
572 approval_request.evidence_refs = options.evidence_refs.clone();
573 approval_request.undo_metadata = options
574 .undo_metadata
575 .clone()
576 .or_else(|| {
577 crate::orchestration::current_mutation_session()
578 .and_then(|session| serde_json::to_value(session).ok())
579 })
580 .unwrap_or(JsonValue::Null);
581 approval_request.capabilities_requested = options.capabilities_requested.clone();
582 let approval_request_json = serde_json::to_value(&approval_request)
583 .map_err(|error| VmError::Runtime(error.to_string()))?;
584 let log = ensure_hitl_event_log();
585 let request = HitlRequestEnvelope {
586 request_id: request_id.clone(),
587 kind: HitlRequestKind::Approval,
588 agent,
589 trace_id: trace_id.clone(),
590 run_id: crate::orchestration::current_mutation_session().and_then(|session| session.run_id),
591 requested_at: requested_at.clone(),
592 payload: json!({
593 "approval_request": approval_request_json,
594 "id": approval_request.id,
595 "action": action,
596 "args": approval_request.args,
597 "principal": approval_request.principal,
598 "requested_at": requested_at,
599 "deadline": approval_request.deadline,
600 "approvers_required": approval_request.approvers_required,
601 "evidence_refs": approval_request.evidence_refs,
602 "undo_metadata": approval_request.undo_metadata,
603 "capabilities_requested": approval_request.capabilities_requested,
604 "detail": options.detail.as_ref().map(crate::llm::vm_value_to_json),
605 "quorum": options.quorum,
606 "reviewers": options.reviewers,
607 "deadline_ms": options.deadline.as_millis() as u64,
608 }),
609 };
610 create_request_waitpoint(&log, &request).await?;
611 append_request(&log, &request).await?;
612 maybe_notify_host(ctx, &request);
613 emit_hitl_requested(&request);
614 maybe_apply_mock_response_with_harness(
615 ctx,
616 harness,
617 HitlRequestKind::Approval,
618 &request_id,
619 &request.payload,
620 )
621 .await?;
622
623 match wait_for_request_waitpoint_with_events(
624 &request_id,
625 HitlRequestKind::Approval,
626 Some(options.deadline),
627 )
628 .await?
629 {
630 WaitpointOutcome::Completed(record) => {
631 approval_record_from_waitpoint(&record, "request_approval")
632 }
633 WaitpointOutcome::Timeout => {
634 append_timeout_once(&log, HitlRequestKind::Approval, &request_id, &trace_id).await?;
635 Err(timeout_error(&request_id, HitlRequestKind::Approval))
636 }
637 WaitpointOutcome::Cancelled { .. } => {
638 Err(approval_wait_error(&log, HitlRequestKind::Approval, &request_id).await)
639 }
640 }
641}
642
643pub(crate) async fn request_approval_for_side_effect(
644 harness: Option<&crate::harness::VmHarness>,
645 action: &str,
646 detail: JsonValue,
647 principal: String,
648 reviewers: Vec<String>,
649 capabilities_requested: Vec<String>,
650) -> Result<VmValue, VmError> {
651 let mut options = crate::value::DictMap::new();
652 options.insert(
653 crate::value::intern_key("args"),
654 crate::stdlib::json_to_vm_value(&detail),
655 );
656 options.insert(
657 crate::value::intern_key("detail"),
658 crate::stdlib::json_to_vm_value(&detail),
659 );
660 options.put_str("principal", principal);
661 options.insert(
662 crate::value::intern_key("reviewers"),
663 VmValue::List(std::sync::Arc::new(
664 reviewers
665 .into_iter()
666 .map(|reviewer| VmValue::String(arcstr::ArcStr::from(reviewer)))
667 .collect(),
668 )),
669 );
670 options.insert(
671 crate::value::intern_key("capabilities_requested"),
672 VmValue::List(std::sync::Arc::new(
673 capabilities_requested
674 .into_iter()
675 .map(|capability| VmValue::String(arcstr::ArcStr::from(capability)))
676 .collect(),
677 )),
678 );
679 let args = vec![
680 VmValue::String(arcstr::ArcStr::from(action.to_string())),
681 VmValue::dict(options),
682 ];
683 request_approval_impl(None, harness, &args).await
684}
685
686async fn dual_control_impl(ctx: &AsyncBuiltinCtx, args: &[VmValue]) -> Result<VmValue, VmError> {
687 let control = Args::runtime("dual_control", args);
688 let n = quorum_arg(&control, 0, "n")?;
689 let m = quorum_arg(&control, 1, "m")?;
690 if n > m {
691 return Err(VmError::Runtime(
692 "dual_control: n must be less than or equal to m".to_string(),
693 ));
694 }
695 let action = args
696 .get(2)
697 .and_then(|value| match value {
698 VmValue::Closure(closure) => Some(closure.clone()),
699 _ => None,
700 })
701 .ok_or_else(|| VmError::Runtime("dual_control: action must be a closure".to_string()))?;
702 let approvers = optional_string_list(args.get(3), "dual_control")?;
703 if !approvers.is_empty() && approvers.len() < m as usize {
704 return Err(VmError::Runtime(format!(
705 "dual_control: expected at least {m} approvers, got {}",
706 approvers.len()
707 )));
708 }
709
710 let keys = current_dispatch_keys();
711 let request_id = next_request_id(HitlRequestKind::DualControl, keys.as_ref());
712 let trace_id = keys
713 .as_ref()
714 .map(|keys| keys.trace_id.clone())
715 .unwrap_or_else(new_trace_id);
716 let action_name = if action.func.name.is_empty() {
717 "anonymous".to_string()
718 } else {
719 action.func.name.to_string()
720 };
721 let agent = keys
722 .as_ref()
723 .map(|keys| keys.agent.clone())
724 .unwrap_or_default();
725 let requested_at_time = OffsetDateTime::now_utc();
726 let requested_at = format_rfc3339(requested_at_time);
727 let mut approval_request = ApprovalRequest::new(
728 request_id.clone(),
729 action_name.clone(),
730 json!({
731 "n": n,
732 "m": m,
733 "approvers": approvers.clone(),
734 }),
735 agent.clone(),
736 requested_at.clone(),
737 );
738 approval_request.deadline = deadline_after(
739 requested_at_time,
740 StdDuration::from_millis(HITL_APPROVAL_TIMEOUT_MS),
741 );
742 approval_request.approvers_required = n as u32;
743 approval_request.undo_metadata = crate::orchestration::current_mutation_session()
744 .and_then(|session| serde_json::to_value(session).ok())
745 .unwrap_or(JsonValue::Null);
746 let approval_request_json = serde_json::to_value(&approval_request)
747 .map_err(|error| VmError::Runtime(error.to_string()))?;
748 let log = ensure_hitl_event_log();
749 let request = HitlRequestEnvelope {
750 request_id: request_id.clone(),
751 kind: HitlRequestKind::DualControl,
752 agent,
753 trace_id: trace_id.clone(),
754 run_id: crate::orchestration::current_mutation_session().and_then(|session| session.run_id),
755 requested_at: requested_at.clone(),
756 payload: json!({
757 "approval_request": approval_request_json,
758 "id": approval_request.id,
759 "args": approval_request.args,
760 "principal": approval_request.principal,
761 "requested_at": requested_at,
762 "deadline": approval_request.deadline,
763 "approvers_required": approval_request.approvers_required,
764 "evidence_refs": approval_request.evidence_refs,
765 "undo_metadata": approval_request.undo_metadata,
766 "capabilities_requested": approval_request.capabilities_requested,
767 "n": n,
768 "m": m,
769 "action": action_name,
770 "approvers": approvers,
771 "deadline_ms": HITL_APPROVAL_TIMEOUT_MS,
772 }),
773 };
774 create_request_waitpoint(&log, &request).await?;
775 append_request(&log, &request).await?;
776 maybe_notify_host(Some(ctx), &request);
777 emit_hitl_requested(&request);
778 maybe_apply_mock_response(
779 Some(ctx),
780 HitlRequestKind::DualControl,
781 &request_id,
782 &request.payload,
783 )
784 .await?;
785
786 match wait_for_request_waitpoint_with_events(
787 &request_id,
788 HitlRequestKind::DualControl,
789 Some(StdDuration::from_millis(HITL_APPROVAL_TIMEOUT_MS)),
790 )
791 .await?
792 {
793 WaitpointOutcome::Completed(record) => {
794 let _ = approval_record_from_waitpoint(&record, "dual_control")?;
795 let mut vm = ctx.child_vm();
796 let result = vm.call_closure_pub(&action, &[]).await?;
797 ctx.forward_output(&vm.take_output());
798
799 append_named_event(
800 &log,
801 HitlRequestKind::DualControl,
802 "hitl.dual_control_executed",
803 &request_id,
804 &trace_id,
805 json!({
806 "request_id": request_id,
807 "result": crate::llm::vm_value_to_json(&result),
808 }),
809 )
810 .await?;
811
812 Ok(result)
813 }
814 WaitpointOutcome::Timeout => {
815 append_timeout_once(&log, HitlRequestKind::DualControl, &request_id, &trace_id).await?;
816 Err(timeout_error(&request_id, HitlRequestKind::DualControl))
817 }
818 WaitpointOutcome::Cancelled { .. } => {
819 Err(approval_wait_error(&log, HitlRequestKind::DualControl, &request_id).await)
820 }
821 }
822}
823
824async fn escalate_to_impl(
825 ctx: Option<&AsyncBuiltinCtx>,
826 args: &[VmValue],
827) -> Result<VmValue, VmError> {
828 let escalate = Args::runtime("escalate_to", args);
829 let role = escalate.non_empty_string(0, "role")?.to_string();
830 let reason = escalate.non_empty_string(1, "reason")?.to_string();
831 let keys = current_dispatch_keys();
832 let request_id = next_request_id(HitlRequestKind::Escalation, keys.as_ref());
833 let trace_id = keys
834 .as_ref()
835 .map(|keys| keys.trace_id.clone())
836 .unwrap_or_else(new_trace_id);
837 let log = ensure_hitl_event_log();
838 let request = HitlRequestEnvelope {
839 request_id: request_id.clone(),
840 kind: HitlRequestKind::Escalation,
841 agent: keys
842 .as_ref()
843 .map(|keys| keys.agent.clone())
844 .unwrap_or_default(),
845 trace_id: trace_id.clone(),
846 run_id: crate::orchestration::current_mutation_session().and_then(|session| session.run_id),
847 requested_at: now_rfc3339(),
848 payload: json!({
849 "role": role,
850 "reason": reason,
851 "capability_policy": escalation_capability_policy(),
852 }),
853 };
854 create_request_waitpoint(&log, &request).await?;
855 append_request(&log, &request).await?;
856 maybe_notify_host(ctx, &request);
857 emit_hitl_requested(&request);
858 maybe_apply_mock_response(
859 ctx,
860 HitlRequestKind::Escalation,
861 &request_id,
862 &request.payload,
863 )
864 .await?;
865
866 match wait_for_request_waitpoint_with_events(&request_id, HitlRequestKind::Escalation, None)
867 .await?
868 {
869 WaitpointOutcome::Completed(record) => {
870 let accepted_at = record.completed_at.clone();
871 let reviewer = record.completed_by.clone();
872 let accepted = record
873 .value
874 .as_ref()
875 .and_then(|value| value.get("accepted"))
876 .and_then(JsonValue::as_bool)
877 .unwrap_or(true);
878 Ok(crate::stdlib::json_to_vm_value(&json!({
879 "request_id": request_id,
880 "role": role,
881 "reason": reason,
882 "trace_id": trace_id,
883 "status": if accepted { "accepted" } else { "pending" },
884 "accepted_at": accepted_at,
885 "reviewer": reviewer,
886 })))
887 }
888 WaitpointOutcome::Timeout => Err(timeout_error(&request_id, HitlRequestKind::Escalation)),
889 WaitpointOutcome::Cancelled {
890 wait_id,
891 waitpoint_ids,
892 reason,
893 } => Err(hitl_cancelled_error(
894 &request_id,
895 HitlRequestKind::Escalation,
896 &wait_id,
897 &waitpoint_ids,
898 reason,
899 )),
900 }
901}
902
903async fn create_request_waitpoint(
904 log: &Arc<AnyEventLog>,
905 request: &HitlRequestEnvelope,
906) -> Result<(), VmError> {
907 create_waitpoint_on(
908 log,
909 Some(request.request_id.clone()),
910 Some(json!({
911 "kind": request.kind.as_str(),
912 "agent": request.agent.clone(),
913 "trace_id": request.trace_id.clone(),
914 "requested_at": request.requested_at.clone(),
915 "payload": request.payload.clone(),
916 })),
917 )
918 .await?;
919 Ok(())
920}
921
922async fn wait_for_request_waitpoint(
923 request_id: &str,
924 timeout: Option<StdDuration>,
925) -> Result<WaitpointOutcome, VmError> {
926 match wait_on_waitpoints(
927 vec![request_id.to_string()],
928 WaitpointWaitOptions { timeout },
929 )
930 .await
931 {
932 Ok(records) => Ok(WaitpointOutcome::Completed(
933 records
934 .into_iter()
935 .next()
936 .expect("single waitpoint wait result"),
937 )),
938 Err(WaitpointWaitFailure::Timeout { .. }) => Ok(WaitpointOutcome::Timeout),
939 Err(WaitpointWaitFailure::Cancelled {
940 wait_id,
941 waitpoint_ids,
942 reason,
943 }) => Ok(WaitpointOutcome::Cancelled {
944 wait_id,
945 waitpoint_ids,
946 reason,
947 }),
948 Err(WaitpointWaitFailure::Vm(error)) => {
949 if let Some(outcome) = waitpoint_outcome_from_vm_error(&error) {
950 return Ok(outcome);
951 }
952 Err(error)
953 }
954 }
955}
956
957fn waitpoint_outcome_from_vm_error(error: &VmError) -> Option<WaitpointOutcome> {
958 let VmError::Thrown(VmValue::Dict(dict)) = error else {
959 return None;
960 };
961 let name = dict.get("name").and_then(vm_string)?;
962 match name {
963 "WaitpointTimeoutError" => Some(WaitpointOutcome::Timeout),
964 "WaitpointCancelledError" => Some(WaitpointOutcome::Cancelled {
965 wait_id: dict
966 .get("wait_id")
967 .and_then(vm_string)
968 .unwrap_or_default()
969 .to_string(),
970 waitpoint_ids: dict
971 .get("waitpoint_ids")
972 .and_then(vm_string_list)
973 .unwrap_or_default(),
974 reason: dict
975 .get("reason")
976 .and_then(vm_string)
977 .map(ToString::to_string),
978 }),
979 _ => None,
980 }
981}
982
983async fn finalize_hitl_response(
984 log: &Arc<AnyEventLog>,
985 kind: HitlRequestKind,
986 response: &HitlHostResponse,
987) -> Result<(), String> {
988 match kind {
989 HitlRequestKind::Question => {
990 if waitpoint_is_terminal(log, &response.request_id).await? {
991 return Ok(());
992 }
993 complete_waitpoint_on(
994 log,
995 &response.request_id,
996 response.answer.clone(),
997 response.reviewer.clone(),
998 response.reason.clone(),
999 response.metadata.clone(),
1000 )
1001 .await
1002 .map(|_| ())
1003 .map_err(|error| error.to_string())
1004 }
1005 HitlRequestKind::Escalation => {
1006 if !response.accepted.unwrap_or(false)
1007 || waitpoint_is_terminal(log, &response.request_id).await?
1008 {
1009 return Ok(());
1010 }
1011 complete_waitpoint_on(
1012 log,
1013 &response.request_id,
1014 Some(json!({
1015 "accepted": true,
1016 "reviewer": response.reviewer,
1017 "reason": response.reason,
1018 "responded_at": response.responded_at,
1019 })),
1020 response.reviewer.clone(),
1021 response.reason.clone(),
1022 response.metadata.clone(),
1023 )
1024 .await
1025 .map(|_| ())
1026 .map_err(|error| error.to_string())
1027 }
1028 HitlRequestKind::Approval | HitlRequestKind::DualControl => {
1029 if waitpoint_is_terminal(log, &response.request_id).await? {
1030 return Ok(());
1031 }
1032 let request = load_request_envelope(log, kind, &response.request_id)
1033 .await
1034 .map_err(|error| error.to_string())?;
1035 match resolve_approval_state(log, kind, &request)
1036 .await
1037 .map_err(|error| error.to_string())?
1038 {
1039 ApprovalResolution::Pending => Ok(()),
1040 ApprovalResolution::Approved(progress) => {
1041 let record = approval_record_json(&progress);
1042 append_named_event(
1043 log,
1044 kind,
1045 approved_event_kind(kind),
1046 &request.request_id,
1047 &request.trace_id,
1048 json!({
1049 "request_id": request.request_id.clone(),
1050 "record": record.clone(),
1051 }),
1052 )
1053 .await
1054 .map_err(|error| error.to_string())?;
1055 complete_waitpoint_on(
1056 log,
1057 &request.request_id,
1058 Some(record),
1059 response.reviewer.clone(),
1060 progress.reason.clone(),
1061 response.metadata.clone(),
1062 )
1063 .await
1064 .map(|_| ())
1065 .map_err(|error| error.to_string())
1066 }
1067 ApprovalResolution::Denied(denied) => {
1068 append_named_event(
1069 log,
1070 kind,
1071 denied_event_kind(kind),
1072 &request.request_id,
1073 &request.trace_id,
1074 json!({
1075 "request_id": request.request_id.clone(),
1076 "reviewer": denied.reviewer.clone(),
1077 "reason": denied.reason.clone(),
1078 }),
1079 )
1080 .await
1081 .map_err(|error| error.to_string())?;
1082 cancel_waitpoint_on(
1083 log,
1084 &request.request_id,
1085 denied.reviewer.clone(),
1086 denied.reason.clone(),
1087 denied.metadata.clone(),
1088 )
1089 .await
1090 .map(|_| ())
1091 .map_err(|error| error.to_string())
1092 }
1093 }
1094 }
1095 }
1096}
1097
1098async fn waitpoint_is_terminal(log: &Arc<AnyEventLog>, request_id: &str) -> Result<bool, String> {
1099 Ok(inspect_waitpoint_on(log, request_id)
1100 .await
1101 .map_err(|error| error.to_string())?
1102 .is_some_and(|record| record.status != WaitpointStatus::Open))
1103}
1104
1105async fn load_request_envelope(
1106 log: &Arc<AnyEventLog>,
1107 kind: HitlRequestKind,
1108 request_id: &str,
1109) -> Result<HitlRequestEnvelope, VmError> {
1110 let topic = topic(kind)?;
1111 let events = log
1112 .read_range(&topic, None, usize::MAX)
1113 .await
1114 .map_err(log_error)?;
1115 events
1116 .into_iter()
1117 .filter(|(_, event)| event.kind == kind.request_event_kind())
1118 .find_map(|(_, event)| {
1119 if !event_matches_request(&event, request_id) {
1120 return None;
1121 }
1122 serde_json::from_value::<HitlRequestEnvelope>(event.payload).ok()
1123 })
1124 .ok_or_else(|| {
1125 VmError::Runtime(format!("missing HITL request envelope for '{request_id}'"))
1126 })
1127}
1128
1129async fn resolve_approval_state(
1130 log: &Arc<AnyEventLog>,
1131 kind: HitlRequestKind,
1132 request: &HitlRequestEnvelope,
1133) -> Result<ApprovalResolution, VmError> {
1134 let quorum = approval_quorum_from_request(kind, request)?;
1135 let allowed_reviewers = approval_reviewers_from_request(kind, request)
1136 .into_iter()
1137 .collect::<BTreeSet<_>>();
1138 let mut progress = ApprovalProgress {
1139 request_id: request.request_id.clone(),
1140 reviewers: BTreeSet::new(),
1141 signatures: Vec::new(),
1142 reason: None,
1143 approved_at: None,
1144 };
1145 let topic = topic(kind)?;
1146 let events = log
1147 .read_range(&topic, None, usize::MAX)
1148 .await
1149 .map_err(log_error)?;
1150 for (_, event) in events {
1151 if !event_matches_request(&event, &request.request_id)
1152 || event.kind != "hitl.response_received"
1153 {
1154 continue;
1155 }
1156 let response: HitlHostResponse = serde_json::from_value(event.payload)
1157 .map_err(|error| VmError::Runtime(error.to_string()))?;
1158 if let Some(reviewer) = response.reviewer.as_deref() {
1159 if !allowed_reviewers.is_empty() && !allowed_reviewers.contains(reviewer) {
1160 continue;
1161 }
1162 if progress.reviewers.contains(reviewer) {
1163 continue;
1164 }
1165 }
1166 if response.approved.unwrap_or(false) {
1167 if let Some(reviewer) = response.reviewer.clone() {
1168 let signed_at = response.responded_at.clone().unwrap_or_else(now_rfc3339);
1169 progress.reviewers.insert(reviewer.clone());
1170 progress.signatures.push(ApprovalSignature {
1171 reviewer: reviewer.clone(),
1172 signed_at: signed_at.clone(),
1173 signature: response.signature.clone().unwrap_or_else(|| {
1174 approval_receipt_signature(
1175 &request.request_id,
1176 &reviewer,
1177 &signed_at,
1178 true,
1179 response.reason.as_deref(),
1180 )
1181 }),
1182 });
1183 }
1184 progress.reason = response.reason.clone();
1185 progress.approved_at = response.responded_at.clone();
1186 if progress.reviewers.len() as u32 >= quorum {
1187 return Ok(ApprovalResolution::Approved(progress));
1188 }
1189 continue;
1190 }
1191 return Ok(ApprovalResolution::Denied(response));
1192 }
1193 Ok(ApprovalResolution::Pending)
1194}
1195
1196fn approval_quorum_from_request(
1197 kind: HitlRequestKind,
1198 request: &HitlRequestEnvelope,
1199) -> Result<u32, VmError> {
1200 let key = match kind {
1201 HitlRequestKind::DualControl => "n",
1202 _ => "quorum",
1203 };
1204 let quorum = request
1205 .payload
1206 .get(key)
1207 .or_else(|| request.payload.get("approvers_required"))
1208 .or_else(|| {
1209 request
1210 .payload
1211 .get("approval_request")
1212 .and_then(|approval| approval.get("approvers_required"))
1213 })
1214 .and_then(JsonValue::as_u64)
1215 .unwrap_or(1);
1216 u32::try_from(quorum).map_err(|_| {
1217 VmError::Runtime(format!(
1218 "invalid quorum in HITL request '{}'",
1219 request.request_id
1220 ))
1221 })
1222}
1223
1224fn approval_reviewers_from_request(
1225 kind: HitlRequestKind,
1226 request: &HitlRequestEnvelope,
1227) -> Vec<String> {
1228 let key = match kind {
1229 HitlRequestKind::DualControl => "approvers",
1230 _ => "reviewers",
1231 };
1232 request
1233 .payload
1234 .get(key)
1235 .or_else(|| {
1236 request
1237 .payload
1238 .get("approval_request")
1239 .and_then(|approval| approval.get("reviewers"))
1240 })
1241 .and_then(JsonValue::as_array)
1242 .map(|values| {
1243 values
1244 .iter()
1245 .filter_map(JsonValue::as_str)
1246 .map(str::to_string)
1247 .collect()
1248 })
1249 .unwrap_or_default()
1250}
1251
1252fn approval_record_json(progress: &ApprovalProgress) -> JsonValue {
1253 json!({
1254 "request_id": progress.request_id.clone(),
1255 "approved": true,
1256 "reviewers": progress.reviewers.iter().cloned().collect::<Vec<_>>(),
1257 "approved_at": progress.approved_at.clone().unwrap_or_else(now_rfc3339),
1258 "reason": progress.reason,
1259 "signatures": progress.signatures,
1260 })
1261}
1262
1263fn approval_receipt_signature(
1264 request_id: &str,
1265 reviewer: &str,
1266 signed_at: &str,
1267 approved: bool,
1268 reason: Option<&str>,
1269) -> String {
1270 let material = format!(
1271 "harn-hitl-approval-v1\nrequest_id:{request_id}\nreviewer:{reviewer}\nsigned_at:{signed_at}\napproved:{approved}\nreason:{}\n",
1272 reason.unwrap_or("")
1273 );
1274 let hash = sha2::Sha256::digest(material.as_bytes());
1275 let hex: String = hash.iter().map(|byte| format!("{byte:02x}")).collect();
1276 format!("sha256:{hex}")
1277}
1278
1279fn approval_record_from_waitpoint(
1280 record: &WaitpointRecord,
1281 builtin: &str,
1282) -> Result<VmValue, VmError> {
1283 record
1284 .value
1285 .as_ref()
1286 .map(crate::stdlib::json_to_vm_value)
1287 .ok_or_else(|| VmError::Runtime(format!("{builtin}: missing approval record")))
1288}
1289
1290async fn approval_wait_error(
1291 log: &Arc<AnyEventLog>,
1292 kind: HitlRequestKind,
1293 request_id: &str,
1294) -> VmError {
1295 if let Ok(Some(record)) = inspect_waitpoint_on(log, request_id).await {
1296 if record.status == WaitpointStatus::Cancelled
1297 && record.reason.as_deref() != Some("upstream_cancelled")
1298 {
1299 return approval_denied_error(
1300 request_id,
1301 HitlHostResponse {
1302 request_id: request_id.to_string(),
1303 answer: None,
1304 approved: Some(false),
1305 accepted: None,
1306 reviewer: record.cancelled_by.clone(),
1307 reason: record.reason.clone(),
1308 metadata: record.metadata.clone(),
1309 responded_at: record.cancelled_at,
1310 signature: None,
1311 },
1312 );
1313 }
1314 if record.status == WaitpointStatus::Cancelled {
1315 return hitl_cancelled_error(
1316 request_id,
1317 kind,
1318 "",
1319 &[request_id.to_string()],
1320 record.reason,
1321 );
1322 }
1323 }
1324 hitl_cancelled_error(
1325 request_id,
1326 kind,
1327 "",
1328 &[request_id.to_string()],
1329 Some("upstream_cancelled".to_string()),
1330 )
1331}
1332
1333async fn append_timeout_once(
1334 log: &Arc<AnyEventLog>,
1335 kind: HitlRequestKind,
1336 request_id: &str,
1337 trace_id: &str,
1338) -> Result<(), VmError> {
1339 if hitl_event_exists(log, kind, request_id, "hitl.timeout").await? {
1340 return Ok(());
1341 }
1342 append_timeout(log, kind, request_id, trace_id).await
1343}
1344
1345async fn hitl_event_exists(
1346 log: &Arc<AnyEventLog>,
1347 kind: HitlRequestKind,
1348 request_id: &str,
1349 event_kind: &str,
1350) -> Result<bool, VmError> {
1351 let topic = topic(kind)?;
1352 let events = log
1353 .read_range(&topic, None, usize::MAX)
1354 .await
1355 .map_err(log_error)?;
1356 Ok(events
1357 .into_iter()
1358 .any(|(_, event)| event.kind == event_kind && event_matches_request(&event, request_id)))
1359}
1360
1361fn approved_event_kind(kind: HitlRequestKind) -> &'static str {
1362 match kind {
1363 HitlRequestKind::DualControl => "hitl.dual_control_approved",
1364 _ => "hitl.approval_approved",
1365 }
1366}
1367
1368fn denied_event_kind(kind: HitlRequestKind) -> &'static str {
1369 match kind {
1370 HitlRequestKind::DualControl => "hitl.dual_control_denied",
1371 _ => "hitl.approval_denied",
1372 }
1373}
1374
1375async fn append_request(
1376 log: &Arc<AnyEventLog>,
1377 request: &HitlRequestEnvelope,
1378) -> Result<(), VmError> {
1379 let topic = topic(request.kind)?;
1380 log.append(
1381 &topic,
1382 LogEvent::new(
1383 request.kind.request_event_kind(),
1384 serde_json::to_value(request).map_err(|error| VmError::Runtime(error.to_string()))?,
1385 )
1386 .with_headers(request_headers(request)),
1387 )
1388 .await
1389 .map(|_| ())
1390 .map_err(log_error)
1391}
1392
1393async fn append_named_event(
1394 log: &Arc<AnyEventLog>,
1395 kind: HitlRequestKind,
1396 event_kind: &str,
1397 request_id: &str,
1398 trace_id: &str,
1399 payload: JsonValue,
1400) -> Result<(), VmError> {
1401 let topic = topic(kind)?;
1402 let headers = headers_with_trace(request_id, trace_id);
1403 log.append(
1404 &topic,
1405 LogEvent::new(event_kind, payload).with_headers(headers),
1406 )
1407 .await
1408 .map(|_| ())
1409 .map_err(log_error)
1410}
1411
1412async fn append_timeout(
1413 log: &Arc<AnyEventLog>,
1414 kind: HitlRequestKind,
1415 request_id: &str,
1416 trace_id: &str,
1417) -> Result<(), VmError> {
1418 append_named_event(
1419 log,
1420 kind,
1421 "hitl.timeout",
1422 request_id,
1423 trace_id,
1424 serde_json::to_value(HitlTimeoutRecord {
1425 request_id: request_id.to_string(),
1426 kind,
1427 trace_id: trace_id.to_string(),
1428 timed_out_at: now_rfc3339(),
1429 })
1430 .map_err(|error| VmError::Runtime(error.to_string()))?,
1431 )
1432 .await
1433}
1434
1435fn parse_hitl_response_dict(
1436 request_id: &str,
1437 response_dict: &crate::value::DictMap,
1438) -> Result<HitlHostResponse, VmError> {
1439 Ok(HitlHostResponse {
1440 request_id: request_id.to_string(),
1441 answer: response_dict
1442 .get("answer")
1443 .map(crate::llm::vm_value_to_json),
1444 approved: response_dict.get("approved").and_then(vm_bool),
1445 accepted: response_dict.get("accepted").and_then(vm_bool),
1446 reviewer: response_dict.get("reviewer").map(VmValue::display),
1447 reason: response_dict.get("reason").map(VmValue::display),
1448 metadata: response_dict
1449 .get("metadata")
1450 .map(crate::llm::vm_value_to_json),
1451 responded_at: response_dict.get("responded_at").map(VmValue::display),
1452 signature: response_dict.get("signature").map(VmValue::display),
1453 })
1454}
1455
1456fn maybe_notify_host(ctx: Option<&AsyncBuiltinCtx>, request: &HitlRequestEnvelope) {
1457 let Some(bridge) = ctx.and_then(|ctx| ctx.child_vm().bridge.clone()) else {
1458 return;
1459 };
1460 bridge.notify(
1461 "harn.hitl.requested",
1462 serde_json::to_value(request).unwrap_or(JsonValue::Null),
1463 );
1464}
1465
1466fn emit_hitl_requested(request: &HitlRequestEnvelope) {
1473 let Some(session_id) = crate::agent_sessions::current_session_id() else {
1474 return;
1475 };
1476 crate::agent_events::emit_event(&crate::agent_events::AgentEvent::HitlRequested {
1477 session_id,
1478 request_id: request.request_id.clone(),
1479 kind: request.kind.as_str().to_string(),
1480 payload: request.payload.clone(),
1481 });
1482}
1483
1484fn emit_hitl_resolved(request_id: &str, kind: HitlRequestKind, outcome: &str) {
1489 let Some(session_id) = crate::agent_sessions::current_session_id() else {
1490 return;
1491 };
1492 crate::agent_events::emit_event(&crate::agent_events::AgentEvent::HitlResolved {
1493 session_id,
1494 request_id: request_id.to_string(),
1495 kind: kind.as_str().to_string(),
1496 outcome: outcome.to_string(),
1497 });
1498}
1499
1500async fn wait_for_request_waitpoint_with_events(
1507 request_id: &str,
1508 kind: HitlRequestKind,
1509 timeout: Option<StdDuration>,
1510) -> Result<WaitpointOutcome, VmError> {
1511 let outcome = wait_for_request_waitpoint(request_id, timeout).await;
1512 let label = match &outcome {
1513 Ok(WaitpointOutcome::Completed(_)) => "answered",
1514 Ok(WaitpointOutcome::Timeout) => "timeout",
1515 Ok(WaitpointOutcome::Cancelled { .. }) => "cancelled",
1516 Err(_) => "error",
1517 };
1518 emit_hitl_resolved(request_id, kind, label);
1519 outcome
1520}
1521
1522fn parse_ask_user_options(value: Option<&VmValue>) -> Result<AskUserOptions, VmError> {
1523 let Some(value) = value else {
1524 return Ok(AskUserOptions {
1525 schema: None,
1526 timeout: Some(default_question_timeout()),
1527 default: None,
1528 });
1529 };
1530 let dict = value
1531 .as_dict()
1532 .ok_or_else(|| VmError::Runtime("ask_user: options must be a dict".to_string()))?;
1533 Ok(AskUserOptions {
1534 schema: dict
1535 .get("schema")
1536 .cloned()
1537 .filter(|value| !matches!(value, VmValue::Nil)),
1538 timeout: dict
1539 .get("timeout")
1540 .map(parse_duration_value)
1541 .transpose()?
1542 .or_else(|| Some(default_question_timeout())),
1543 default: dict
1544 .get("default")
1545 .cloned()
1546 .filter(|value| !matches!(value, VmValue::Nil)),
1547 })
1548}
1549
1550fn default_question_timeout() -> StdDuration {
1551 StdDuration::from_millis(HITL_QUESTION_TIMEOUT_MS)
1552}
1553
1554fn escalation_capability_policy() -> JsonValue {
1555 crate::orchestration::current_execution_policy()
1556 .and_then(|policy| serde_json::to_value(policy).ok())
1557 .unwrap_or(JsonValue::Null)
1558}
1559
1560fn parse_approval_options(
1561 value: Option<&VmValue>,
1562 builtin: &str,
1563) -> Result<ApprovalOptions, VmError> {
1564 let dict = match value {
1565 None => None,
1566 Some(VmValue::Dict(dict)) => Some(dict),
1567 Some(_) => {
1568 return Err(VmError::Runtime(format!(
1569 "{builtin}: options must be a dict"
1570 )))
1571 }
1572 };
1573 let quorum = dict
1574 .and_then(|dict| dict.get("quorum"))
1575 .and_then(VmValue::as_int)
1576 .unwrap_or(1);
1577 if quorum <= 0 {
1578 return Err(VmError::Runtime(format!(
1579 "{builtin}: quorum must be positive"
1580 )));
1581 }
1582 let reviewers = optional_string_list(dict.and_then(|dict| dict.get("reviewers")), builtin)?;
1583 let capabilities_requested = optional_string_list(
1584 dict.and_then(|dict| dict.get("capabilities_requested")),
1585 builtin,
1586 )?;
1587 let evidence_refs = dict
1588 .and_then(|dict| dict.get("evidence_refs"))
1589 .map(|value| match value {
1590 VmValue::List(items) => Ok(items
1591 .iter()
1592 .map(crate::llm::vm_value_to_json)
1593 .collect::<Vec<_>>()),
1594 _ => Err(VmError::Runtime(format!(
1595 "{builtin}: evidence_refs must be a list"
1596 ))),
1597 })
1598 .transpose()?
1599 .unwrap_or_default();
1600 let deadline = dict
1601 .and_then(|dict| dict.get("deadline"))
1602 .map(parse_duration_value)
1603 .transpose()?
1604 .unwrap_or_else(|| StdDuration::from_millis(HITL_APPROVAL_TIMEOUT_MS));
1605 Ok(ApprovalOptions {
1606 detail: dict.and_then(|dict| dict.get("detail")).cloned(),
1607 args: dict.and_then(|dict| dict.get("args")).cloned(),
1608 quorum: quorum as u32,
1609 reviewers,
1610 deadline,
1611 principal: dict
1612 .and_then(|dict| dict.get("principal"))
1613 .map(VmValue::display)
1614 .filter(|value| !value.is_empty()),
1615 evidence_refs,
1616 undo_metadata: dict
1617 .and_then(|dict| dict.get("undo_metadata"))
1618 .map(crate::llm::vm_value_to_json),
1619 capabilities_requested,
1620 })
1621}
1622
1623fn quorum_arg(args: &Args<'_, '_>, idx: usize, name: &str) -> Result<i64, VmError> {
1625 match args.int(idx, name)? {
1626 value if value > 0 => Ok(value),
1627 _ => Err(ArgError::constraint(
1628 args.fn_name(),
1629 args.kind(),
1630 name,
1631 "must be > 0",
1632 )),
1633 }
1634}
1635
1636fn optional_string_list(value: Option<&VmValue>, builtin: &str) -> Result<Vec<String>, VmError> {
1637 let Some(value) = value else {
1638 return Ok(Vec::new());
1639 };
1640 match value {
1641 VmValue::List(list) => Ok(list.iter().map(VmValue::display).collect()),
1642 _ => Err(VmError::Runtime(format!(
1643 "{builtin}: expected list<string>"
1644 ))),
1645 }
1646}
1647
1648fn parse_duration_value(value: &VmValue) -> Result<StdDuration, VmError> {
1649 Args::single("hitl", ErrorKind::Runtime, Some(value)).duration(0, "timeout")
1650}
1651
1652fn ensure_hitl_event_log() -> Arc<AnyEventLog> {
1653 active_event_log()
1654 .unwrap_or_else(|| install_memory_for_current_thread(HITL_EVENT_LOG_QUEUE_DEPTH))
1655}
1656
1657fn ensure_hitl_event_log_for(base_dir: Option<&Path>) -> Result<Arc<AnyEventLog>, String> {
1658 if let Some(log) = active_event_log() {
1659 return Ok(log);
1660 }
1661 let Some(base_dir) = base_dir else {
1662 return Ok(install_memory_for_current_thread(
1663 HITL_EVENT_LOG_QUEUE_DEPTH,
1664 ));
1665 };
1666 install_default_for_base_dir(base_dir).map_err(|error| error.to_string())
1667}
1668
1669fn current_dispatch_keys() -> Option<DispatchKeys> {
1670 let context = current_dispatch_context()?;
1671 let stable_base = context
1672 .replay_of_event_id
1673 .clone()
1674 .unwrap_or_else(|| context.trigger_event.id.0.clone());
1675 let instance_key = format!(
1676 "{}::{}",
1677 context.trigger_event.id.0,
1678 context.replay_of_event_id.as_deref().unwrap_or("live")
1679 );
1680 Some(DispatchKeys {
1681 instance_key,
1682 stable_base,
1683 agent: context.agent_id,
1684 trace_id: context.trigger_event.trace_id.0,
1685 })
1686}
1687
1688fn next_request_id(kind: HitlRequestKind, dispatch_keys: Option<&DispatchKeys>) -> String {
1689 if let Some(keys) = dispatch_keys {
1690 let seq = REQUEST_SEQUENCE.with(|slot| {
1691 let mut state = slot.borrow_mut();
1692 if state.instance_key != keys.instance_key {
1693 state.instance_key = keys.instance_key.clone();
1694 state.next_seq = 0;
1695 }
1696 state.next_seq += 1;
1697 state.next_seq
1698 });
1699 return format!("hitl_{}_{}_{}", kind.as_str(), keys.stable_base, seq);
1700 }
1701 format!("hitl_{}_{}", kind.as_str(), Uuid::now_v7())
1702}
1703
1704fn request_headers(request: &HitlRequestEnvelope) -> BTreeMap<String, String> {
1705 let mut headers = headers_with_trace(&request.request_id, &request.trace_id);
1706 if let Some(run_id) = request.run_id.as_ref() {
1707 headers.insert("run_id".to_string(), run_id.clone());
1708 }
1709 headers
1710}
1711
1712fn response_headers(request_id: &str) -> BTreeMap<String, String> {
1713 let mut headers = std::collections::BTreeMap::new();
1714 headers.insert("request_id".to_string(), request_id.to_string());
1715 headers
1716}
1717
1718fn headers_with_trace(request_id: &str, trace_id: &str) -> BTreeMap<String, String> {
1719 let mut headers = response_headers(request_id);
1720 headers.insert("trace_id".to_string(), trace_id.to_string());
1721 headers
1722}
1723
1724fn topic(kind: HitlRequestKind) -> Result<Topic, VmError> {
1725 Topic::new(kind.topic()).map_err(|error| VmError::Runtime(error.to_string()))
1726}
1727
1728fn event_matches_request(event: &LogEvent, request_id: &str) -> bool {
1729 event
1730 .headers
1731 .get("request_id")
1732 .is_some_and(|value| value == request_id)
1733 || event
1734 .payload
1735 .get("request_id")
1736 .and_then(JsonValue::as_str)
1737 .is_some_and(|value| value == request_id)
1738}
1739
1740fn approval_denied_error(request_id: &str, response: HitlHostResponse) -> VmError {
1741 VmError::Thrown(crate::stdlib::json_to_vm_value(&json!({
1742 "name": "ApprovalDeniedError",
1743 "category": "generic",
1744 "message": response.reason.clone().unwrap_or_else(|| "approval was denied".to_string()),
1745 "request_id": request_id,
1746 "reviewers": response.reviewer.into_iter().collect::<Vec<_>>(),
1747 "reason": response.reason,
1748 })))
1749}
1750
1751fn hitl_cancelled_error(
1752 request_id: &str,
1753 kind: HitlRequestKind,
1754 wait_id: &str,
1755 waitpoint_ids: &[String],
1756 reason: Option<String>,
1757) -> VmError {
1758 let _ = categorized_error("HITL cancelled", ErrorCategory::Cancelled);
1759 let message = reason
1760 .clone()
1761 .unwrap_or_else(|| format!("{} cancelled", kind.as_str()));
1762 VmError::Thrown(crate::stdlib::json_to_vm_value(&json!({
1763 "name": "HumanCancelledError",
1764 "category": ErrorCategory::Cancelled.as_str(),
1765 "message": message,
1766 "request_id": request_id,
1767 "kind": kind.as_str(),
1768 "wait_id": wait_id,
1769 "waitpoint_ids": waitpoint_ids,
1770 "reason": reason,
1771 })))
1772}
1773
1774fn timeout_error(request_id: &str, kind: HitlRequestKind) -> VmError {
1775 let _ = categorized_error("HITL timed out", ErrorCategory::Timeout);
1776 VmError::Thrown(crate::stdlib::json_to_vm_value(&json!({
1777 "name": "HumanTimeoutError",
1778 "category": ErrorCategory::Timeout.as_str(),
1779 "message": format!("{} timed out", kind.as_str()),
1780 "request_id": request_id,
1781 "kind": kind.as_str(),
1782 })))
1783}
1784
1785fn coerce_like_default(value: &VmValue, default: &VmValue) -> VmValue {
1786 match default {
1787 VmValue::Int(_) => match value {
1788 VmValue::Int(_) => value.clone(),
1789 VmValue::Float(number) => VmValue::Int(*number as i64),
1790 VmValue::String(text) => text
1791 .parse::<i64>()
1792 .map(VmValue::Int)
1793 .unwrap_or_else(|_| default.clone()),
1794 _ => default.clone(),
1795 },
1796 VmValue::Float(_) => match value {
1797 VmValue::Float(_) => value.clone(),
1798 VmValue::Int(number) => VmValue::Float(*number as f64),
1799 VmValue::String(text) => text
1800 .parse::<f64>()
1801 .map(VmValue::Float)
1802 .unwrap_or_else(|_| default.clone()),
1803 _ => default.clone(),
1804 },
1805 VmValue::Bool(_) => match value {
1806 VmValue::Bool(_) => value.clone(),
1807 VmValue::String(text) if text.eq_ignore_ascii_case("true") => VmValue::Bool(true),
1808 VmValue::String(text) if text.eq_ignore_ascii_case("false") => VmValue::Bool(false),
1809 _ => default.clone(),
1810 },
1811 VmValue::String(_) => VmValue::String(arcstr::ArcStr::from(value.display())),
1812 VmValue::Duration(_) => match value {
1813 VmValue::Duration(_) => value.clone(),
1814 VmValue::Int(ms) => VmValue::Duration(*ms),
1815 _ => default.clone(),
1816 },
1817 VmValue::Nil => value.clone(),
1818 _ => {
1819 if value.type_name() == default.type_name() {
1820 value.clone()
1821 } else {
1822 default.clone()
1823 }
1824 }
1825 }
1826}
1827
1828fn log_error(error: impl std::fmt::Display) -> VmError {
1829 VmError::Runtime(error.to_string())
1830}
1831
1832fn now_rfc3339() -> String {
1833 format_rfc3339(OffsetDateTime::now_utc())
1834}
1835
1836fn format_rfc3339(timestamp: OffsetDateTime) -> String {
1837 harn_clock::format_rfc3339(timestamp)
1838}
1839
1840fn deadline_after(requested_at: OffsetDateTime, duration: StdDuration) -> Option<String> {
1841 time::Duration::try_from(duration)
1842 .ok()
1843 .map(|duration| format_rfc3339(requested_at + duration))
1844}
1845
1846fn new_trace_id() -> String {
1847 format!("trace_{}", Uuid::now_v7())
1848}
1849
1850fn vm_bool(value: &VmValue) -> Option<bool> {
1851 match value {
1852 VmValue::Bool(flag) => Some(*flag),
1853 _ => None,
1854 }
1855}
1856
1857fn vm_string(value: &VmValue) -> Option<&str> {
1858 match value {
1859 VmValue::String(text) => Some(text.as_ref()),
1860 _ => None,
1861 }
1862}
1863
1864fn vm_string_list(value: &VmValue) -> Option<Vec<String>> {
1865 match value {
1866 VmValue::List(values) => Some(values.iter().map(VmValue::display).collect()),
1867 _ => None,
1868 }
1869}
1870
1871#[cfg(test)]
1872mod tests {
1873 use std::sync::OnceLock;
1874
1875 use tokio::sync::Mutex;
1876
1877 use super::{
1878 HITL_APPROVALS_TOPIC, HITL_DUAL_CONTROL_TOPIC, HITL_ESCALATIONS_TOPIC, HITL_QUESTIONS_TOPIC,
1879 };
1880 use crate::event_log::{
1881 install_active_event_log, install_memory_for_current_thread, EventLog, Topic,
1882 };
1883 use crate::{register_vm_stdlib, reset_thread_local_state, Vm, VmError};
1884
1885 fn hitl_lock() -> &'static Mutex<()> {
1886 static LOCK: OnceLock<Mutex<()>> = OnceLock::new();
1887 LOCK.get_or_init(|| Mutex::new(()))
1888 }
1889
1890 fn compile_hitl_fixture(source: &str) -> crate::Chunk {
1891 let program =
1892 harn_parser::parse_source(source).expect("trusted HITL test fixture should parse");
1893 crate::Compiler::with_options(crate::CompilerOptions::privileged_wire())
1894 .compile(&program)
1895 .expect("trusted HITL test fixture should compile")
1896 }
1897
1898 async fn execute_hitl_script(
1899 base_dir: &std::path::Path,
1900 source: &str,
1901 ) -> Result<(String, Vec<String>, Vec<String>, Vec<String>, Vec<String>), VmError> {
1902 reset_thread_local_state();
1903 let log = install_memory_for_current_thread(super::HITL_EVENT_LOG_QUEUE_DEPTH);
1904 let chunk = compile_hitl_fixture(source);
1905 let mut vm = Vm::new();
1906 register_vm_stdlib(&mut vm);
1907 vm.set_source_dir(base_dir);
1908 vm.execute(&chunk).await?;
1909 let output = vm.output().trim_end().to_string();
1910 let question_events = event_kinds(log.clone(), HITL_QUESTIONS_TOPIC).await;
1911 let approval_events = event_kinds(log.clone(), HITL_APPROVALS_TOPIC).await;
1912 let dual_control_events = event_kinds(log.clone(), HITL_DUAL_CONTROL_TOPIC).await;
1913 let escalation_events = event_kinds(log, HITL_ESCALATIONS_TOPIC).await;
1914 Ok((
1915 output,
1916 question_events,
1917 approval_events,
1918 dual_control_events,
1919 escalation_events,
1920 ))
1921 }
1922
1923 async fn event_kinds(
1924 log: std::sync::Arc<crate::event_log::AnyEventLog>,
1925 topic: &str,
1926 ) -> Vec<String> {
1927 log.read_range(&Topic::new(topic).expect("valid topic"), None, usize::MAX)
1928 .await
1929 .expect("read topic")
1930 .into_iter()
1931 .map(|(_, event)| event.kind)
1932 .collect()
1933 }
1934
1935 async fn event_payloads(
1936 log: std::sync::Arc<crate::event_log::AnyEventLog>,
1937 topic: &str,
1938 ) -> Vec<serde_json::Value> {
1939 log.read_range(&Topic::new(topic).expect("valid topic"), None, usize::MAX)
1940 .await
1941 .expect("read topic")
1942 .into_iter()
1943 .map(|(_, event)| event.payload)
1944 .collect()
1945 }
1946
1947 #[tokio::test(flavor = "current_thread")]
1948 async fn ask_user_coerces_to_default_type_and_logs_events() {
1949 tokio::task::LocalSet::new()
1950 .run_until(async {
1951 let dir = tempfile::tempdir().expect("tempdir");
1952 let source = r#"
1953pipeline test(harness: Harness, task: unknown) {
1954 host_mock("hitl", "question", {answer: "9"})
1955 const answer: int = harness.interaction.ask_user("Pick a number", {default: 0})
1956 harness.stdio.println(answer)
1957}
1958"#;
1959 let (
1960 output,
1961 question_events,
1962 approval_events,
1963 dual_control_events,
1964 escalation_events,
1965 ) = execute_hitl_script(dir.path(), source)
1966 .await
1967 .expect("script succeeds");
1968 assert_eq!(output, "9");
1969 assert_eq!(
1970 question_events,
1971 vec![
1972 "hitl.question_asked".to_string(),
1973 "hitl.response_received".to_string()
1974 ]
1975 );
1976 assert!(approval_events.is_empty());
1977 assert!(dual_control_events.is_empty());
1978 assert!(escalation_events.is_empty());
1979 })
1980 .await;
1981 }
1982
1983 #[tokio::test(flavor = "current_thread")]
1984 async fn request_approval_waits_for_quorum_and_emits_a_record() {
1985 let _guard = hitl_lock().lock().await;
1986 tokio::task::LocalSet::new()
1987 .run_until(async {
1988 reset_thread_local_state();
1989 let dir = tempfile::tempdir().expect("tempdir");
1990 let source = r#"
1991pipeline test(harness: Harness, task: unknown) {
1992 host_mock("hitl", "approval", [
1993 {approved: true, reviewer: "alice", reason: "ok"},
1994 {approved: true, reviewer: "bob", reason: "ship it"},
1995 ])
1996 const record = harness.interaction.request_approval(
1997 "deploy production",
1998 {quorum: 2, reviewers: ["alice", "bob", "carol"]},
1999 )
2000 harness.stdio.println(record.approved)
2001 harness.stdio.println(len(record.reviewers))
2002 harness.stdio.println(record.reviewers[0])
2003 harness.stdio.println(record.reviewers[1])
2004}
2005"#;
2006 let (_, _, approval_events, _, _) = execute_hitl_script(dir.path(), source)
2007 .await
2008 .expect("script succeeds");
2009 assert_eq!(
2010 approval_events,
2011 vec![
2012 "hitl.approval_requested".to_string(),
2013 "hitl.response_received".to_string(),
2014 "hitl.response_received".to_string(),
2015 "hitl.approval_approved".to_string(),
2016 ]
2017 );
2018 })
2019 .await;
2020 }
2021
2022 #[tokio::test(flavor = "current_thread")]
2023 async fn request_approval_emits_canonical_approval_request_payload() {
2024 let _guard = hitl_lock().lock().await;
2025 tokio::task::LocalSet::new()
2026 .run_until(async {
2027 reset_thread_local_state();
2028 let dir = tempfile::tempdir().expect("tempdir");
2029 let log = install_memory_for_current_thread(super::HITL_EVENT_LOG_QUEUE_DEPTH);
2030 let source = r#"
2031pipeline test(harness: Harness, task: unknown) {
2032 host_mock("hitl", "approval", {approved: true, reviewer: "alice", reason: "ok"})
2033 harness.interaction.request_approval("deploy production", {
2034 args: {environment: "prod"},
2035 quorum: 1,
2036 reviewers: ["alice"],
2037 evidence_refs: [{kind: "run", uri: "run_123"}],
2038 undo_metadata: {strategy: "rollback"},
2039 capabilities_requested: ["deploy.production"],
2040 })
2041}
2042"#;
2043 let chunk = compile_hitl_fixture(source);
2044 let mut vm = Vm::new();
2045 register_vm_stdlib(&mut vm);
2046 vm.set_source_dir(dir.path());
2047 vm.execute(&chunk).await.expect("script succeeds");
2048
2049 let payloads = event_payloads(log, HITL_APPROVALS_TOPIC).await;
2050 let request_payload = &payloads[0]["payload"];
2051 let approval_request = &request_payload["approval_request"];
2052 assert_eq!(approval_request["id"], request_payload["id"]);
2053 assert_eq!(approval_request["action"], "deploy production");
2054 assert_eq!(approval_request["args"]["environment"], "prod");
2055 assert_eq!(approval_request["approvers_required"], 1);
2056 assert_eq!(approval_request["evidence_refs"][0]["uri"], "run_123");
2057 assert_eq!(approval_request["undo_metadata"]["strategy"], "rollback");
2058 assert_eq!(
2059 approval_request["capabilities_requested"][0],
2060 "deploy.production"
2061 );
2062 assert!(approval_request["requested_at"].as_str().is_some());
2063 assert!(approval_request["deadline"].as_str().is_some());
2064 })
2065 .await;
2066 }
2067
2068 #[tokio::test(flavor = "current_thread")]
2069 async fn request_approval_surfaces_denials_as_typed_errors() {
2070 let _guard = hitl_lock().lock().await;
2071 tokio::task::LocalSet::new()
2072 .run_until(async {
2073 reset_thread_local_state();
2074 let dir = tempfile::tempdir().expect("tempdir");
2075 let source = r#"
2076pipeline test(harness: Harness, task: unknown) {
2077 host_mock("hitl", "approval", {approved: false, reviewer: "alice", reason: "unsafe"})
2078 const denied = try {
2079 harness.interaction.request_approval("drop table", {reviewers: ["alice"]})
2080 }
2081 harness.stdio.println(is_err(denied))
2082 harness.stdio.println(unwrap_err(denied).name)
2083 harness.stdio.println(unwrap_err(denied).reason)
2084}
2085"#;
2086 let (output, _, approval_events, _, _) = execute_hitl_script(dir.path(), source)
2087 .await
2088 .expect("script succeeds");
2089 assert_eq!(output, "true\nApprovalDeniedError\nunsafe");
2090 assert_eq!(
2091 approval_events,
2092 vec![
2093 "hitl.approval_requested".to_string(),
2094 "hitl.response_received".to_string(),
2095 "hitl.approval_denied".to_string(),
2096 ]
2097 );
2098 })
2099 .await;
2100 }
2101
2102 #[tokio::test(flavor = "current_thread")]
2103 async fn dual_control_executes_action_after_quorum() {
2104 tokio::task::LocalSet::new()
2105 .run_until(async {
2106 let dir = tempfile::tempdir().expect("tempdir");
2107 let source = r#"
2108pipeline test(harness: Harness, task: unknown) {
2109 host_mock("hitl", "dual_control", [
2110 {approved: true, reviewer: "alice"},
2111 {approved: true, reviewer: "bob"},
2112 ])
2113 const result = harness.interaction.dual_control(
2114 2,
2115 3,
2116 { -> "launched" },
2117 ["alice", "bob", "carol"],
2118 )
2119 harness.stdio.println(result)
2120}
2121"#;
2122 let (output, _, _, dual_control_events, _) =
2123 execute_hitl_script(dir.path(), source)
2124 .await
2125 .expect("script succeeds");
2126 assert_eq!(output, "launched");
2127 assert_eq!(
2128 dual_control_events,
2129 vec![
2130 "hitl.dual_control_requested".to_string(),
2131 "hitl.response_received".to_string(),
2132 "hitl.response_received".to_string(),
2133 "hitl.dual_control_approved".to_string(),
2134 "hitl.dual_control_executed".to_string(),
2135 ]
2136 );
2137 })
2138 .await;
2139 }
2140
2141 #[tokio::test(flavor = "current_thread")]
2142 async fn escalate_to_waits_for_acceptance_event() {
2143 tokio::task::LocalSet::new()
2144 .run_until(async {
2145 let dir = tempfile::tempdir().expect("tempdir");
2146 let source = r#"
2147pipeline test(harness: Harness, task: unknown) {
2148 host_mock("hitl", "escalation", {accepted: true, reviewer: "lead", reason: "taking over"})
2149 const handle = harness.interaction.escalate_to("admin", "need override")
2150 harness.stdio.println(handle.status)
2151 harness.stdio.println(handle.reviewer)
2152}
2153"#;
2154 let (output, _, _, _, escalation_events) = execute_hitl_script(dir.path(), source)
2155 .await
2156 .expect("script succeeds");
2157 assert_eq!(output, "accepted\nlead");
2158 assert_eq!(
2159 escalation_events,
2160 vec![
2161 "hitl.escalation_issued".to_string(),
2162 "hitl.escalation_accepted".to_string(),
2163 ]
2164 );
2165 })
2166 .await;
2167 }
2168
2169 #[tokio::test(flavor = "current_thread")]
2175 async fn ask_user_emits_hitl_request_and_resolution_to_agent_event_sinks() {
2176 use std::sync::Mutex as StdMutex;
2177
2178 tokio::task::LocalSet::new()
2179 .run_until(async {
2180 let dir = tempfile::tempdir().expect("tempdir");
2181 let session_id = "hitl-session".to_string();
2182 let captured: std::sync::Arc<StdMutex<Vec<crate::agent_events::AgentEvent>>> =
2183 std::sync::Arc::new(StdMutex::new(Vec::new()));
2184
2185 struct CaptureSink(std::sync::Arc<StdMutex<Vec<crate::agent_events::AgentEvent>>>);
2186 impl crate::agent_events::AgentEventSink for CaptureSink {
2187 fn handle_event(&self, event: &crate::agent_events::AgentEvent) {
2188 self.0.lock().expect("captured").push(event.clone());
2189 }
2190 }
2191
2192 crate::reset_thread_local_state();
2198 let log = std::sync::Arc::new(crate::event_log::AnyEventLog::Memory(
2199 crate::event_log::MemoryEventLog::new(super::HITL_EVENT_LOG_QUEUE_DEPTH),
2200 ));
2201 install_active_event_log(log);
2202
2203 crate::agent_events::reset_all_sinks();
2204 let sink: std::sync::Arc<dyn crate::agent_events::AgentEventSink> =
2205 std::sync::Arc::new(CaptureSink(captured.clone()));
2206 crate::agent_events::register_sink(session_id.clone(), sink);
2207 crate::agent_sessions::open_or_create(Some(session_id.clone()));
2208 let _guard = crate::agent_sessions::enter_current_session(session_id.clone());
2209
2210 let source = r#"
2211pipeline test(harness: Harness, task: unknown) {
2212 host_mock("hitl", "question", {answer: "ok"})
2213 const answer: string = harness.interaction.ask_user("Are you sure?", {default: "no"})
2214 harness.stdio.println(answer)
2215}
2216"#;
2217 let chunk = compile_hitl_fixture(source);
2218 let mut vm = Vm::new();
2219 register_vm_stdlib(&mut vm);
2220 vm.set_source_dir(dir.path());
2221 vm.execute(&chunk).await.expect("script runs");
2222 assert_eq!(vm.output().trim_end(), "ok");
2223
2224 let events = captured.lock().expect("captured");
2225 let mut iter = events.iter().filter(|event| {
2226 matches!(
2227 event,
2228 crate::agent_events::AgentEvent::HitlRequested { .. }
2229 | crate::agent_events::AgentEvent::HitlResolved { .. }
2230 )
2231 });
2232 let requested = iter.next().expect("HitlRequested emitted");
2233 let resolved = iter.next().expect("HitlResolved emitted");
2234 assert!(iter.next().is_none(), "exactly one pair: {events:?}");
2235
2236 let crate::agent_events::AgentEvent::HitlRequested {
2237 session_id: req_session,
2238 request_id: req_id,
2239 kind: req_kind,
2240 payload,
2241 } = requested
2242 else {
2243 panic!("expected HitlRequested, got: {requested:?}");
2244 };
2245 assert_eq!(req_session, &session_id);
2246 assert_eq!(req_kind, "question");
2247 assert!(req_id.starts_with("hitl_question_"));
2248 assert_eq!(payload["prompt"], "Are you sure?");
2249
2250 let crate::agent_events::AgentEvent::HitlResolved {
2251 request_id: res_id,
2252 kind: res_kind,
2253 outcome,
2254 ..
2255 } = resolved
2256 else {
2257 panic!("expected HitlResolved, got: {resolved:?}");
2258 };
2259 assert_eq!(res_id, req_id);
2260 assert_eq!(res_kind, "question");
2261 assert_eq!(outcome, "answered");
2262
2263 drop(_guard);
2264 crate::agent_events::reset_all_sinks();
2265 })
2266 .await;
2267 }
2268}