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