Skip to main content

harn_vm/stdlib/
hitl.rs

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// `Completed` carries the full `WaitpointRecord`, which dominates the
257// enum's size — boxing it would force every match arm to indirect even
258// though the enum is dropped within nanoseconds of being constructed
259// (it's a local return type for the waitpoint poll loop, never stored).
260// Surfaced by the host-target compile of `harn-vm` introduced when
261// `harn-cli`'s build script gained `harn-vm` as a build-dep for the
262// AOT bytecode embedding pass (G7 / harn#2300).
263#[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
1466/// Emit a `HitlRequested` `AgentEvent` so transport adapters
1467/// (currently the A2A `A2aWorkerSink`) can flip a task into
1468/// `input-required` while the script is suspended on the waitpoint.
1469/// No-op when there is no current agent session — the bridge-level
1470/// `harn.hitl.requested` notification still fires for hosts that drive
1471/// HITL UX through the bridge.
1472fn 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
1484/// Companion to `emit_hitl_requested`: notifies sinks that the
1485/// suspended waitpoint has resolved so a paused task can flip back
1486/// out of `input-required`. `outcome` is one of `"answered"`,
1487/// `"timeout"`, `"cancelled"`, or `"error"`.
1488fn 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
1500/// Wrapper around `wait_for_request_waitpoint` that emits the
1501/// canonical `HitlResolved` `AgentEvent` regardless of which terminal
1502/// branch the waitpoint takes (response / timeout / cancellation /
1503/// error). Pair-emitted with `emit_hitl_requested` so transport
1504/// adapters can bracket the `input-required` pause cleanly without
1505/// each `*_impl` having to duplicate the emission at every match arm.
1506async 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
1623/// A quorum bound: an int that has to be at least one.
1624fn 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    /// `harn-serve` adapters (A2A `input-required`, ACP `hitl_request`)
2170    /// rely on the canonical `AgentEvent::HitlRequested` /
2171    /// `AgentEvent::HitlResolved` pair to bracket every HITL pause.
2172    /// Pin the contract here so future HITL primitives keep emitting
2173    /// the event around their waitpoint blocks.
2174    #[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                // Inline the script setup rather than using the
2193                // `execute_hitl_script` helper: that helper calls
2194                // `reset_thread_local_state` (which wipes the session
2195                // store), so any session pushed before it would be
2196                // gone by the time `ask_user` runs.
2197                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}