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::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// `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 = 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
1460/// Emit a `HitlRequested` `AgentEvent` so transport adapters
1461/// (currently the A2A `A2aWorkerSink`) can flip a task into
1462/// `input-required` while the script is suspended on the waitpoint.
1463/// No-op when there is no current agent session — the bridge-level
1464/// `harn.hitl.requested` notification still fires for hosts that drive
1465/// HITL UX through the bridge.
1466fn 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
1478/// Companion to `emit_hitl_requested`: notifies sinks that the
1479/// suspended waitpoint has resolved so a paused task can flip back
1480/// out of `input-required`. `outcome` is one of `"answered"`,
1481/// `"timeout"`, `"cancelled"`, or `"error"`.
1482fn 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
1494/// Wrapper around `wait_for_request_waitpoint` that emits the
1495/// canonical `HitlResolved` `AgentEvent` regardless of which terminal
1496/// branch the waitpoint takes (response / timeout / cancellation /
1497/// error). Pair-emitted with `emit_hitl_requested` so transport
1498/// adapters can bracket the `input-required` pause cleanly without
1499/// each `*_impl` having to duplicate the emission at every match arm.
1500async 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    /// `harn-serve` adapters (A2A `input-required`, ACP `hitl_request`)
2171    /// rely on the canonical `AgentEvent::HitlRequested` /
2172    /// `AgentEvent::HitlResolved` pair to bracket every HITL pause.
2173    /// Pin the contract here so future HITL primitives keep emitting
2174    /// the event around their waitpoint blocks.
2175    #[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                // Inline the script setup rather than using the
2194                // `execute_hitl_script` helper: that helper calls
2195                // `reset_thread_local_state` (which wipes the session
2196                // store), so any session pushed before it would be
2197                // gone by the time `ask_user` runs.
2198                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}