Skip to main content

harn_vm/stdlib/
workflow_messages.rs

1use std::collections::{BTreeMap, VecDeque};
2use std::future::Future;
3use std::path::{Path, PathBuf};
4use std::sync::{Arc, OnceLock, Weak};
5use std::time::Duration as StdDuration;
6
7use harn_clock::Clock;
8use serde::{Deserialize, Serialize};
9
10use crate::stdlib::macros::{harn_builtin, VmBuiltinDef};
11use crate::stdlib::process::runtime_root_base;
12use crate::value::{VmError, VmValue};
13use crate::vm::Vm;
14
15const DEFAULT_UPDATE_TIMEOUT_MS: u64 = 30_000;
16const WORKFLOW_POLL_INTERVAL_MS: u64 = 25;
17
18pub(crate) const MODULE_BUILTINS: &[&VmBuiltinDef] = &[
19    &WORKFLOW_SIGNAL_BUILTIN_DEF,
20    &WORKFLOW_QUERY_BUILTIN_DEF,
21    &WORKFLOW_PUBLISH_QUERY_BUILTIN_DEF,
22    &WORKFLOW_RECEIVE_BUILTIN_DEF,
23    &WORKFLOW_RESPOND_UPDATE_BUILTIN_DEF,
24    &WORKFLOW_PAUSE_BUILTIN_DEF,
25    &WORKFLOW_RESUME_BUILTIN_DEF,
26    &WORKFLOW_STATUS_BUILTIN_DEF,
27    &WORKFLOW_CONTINUE_AS_NEW_BUILTIN_DEF,
28    &CONTINUE_AS_NEW_BUILTIN_DEF,
29    &WORKFLOW_UPDATE_BUILTIN_DEF,
30];
31
32#[derive(Clone, Debug, Serialize, Deserialize)]
33pub struct WorkflowMessageRecord {
34    pub seq: u64,
35    pub kind: String,
36    pub name: String,
37    #[serde(skip_serializing_if = "Option::is_none")]
38    pub request_id: Option<String>,
39    pub payload: serde_json::Value,
40    pub enqueued_at: String,
41}
42
43#[derive(Clone, Debug, Serialize, Deserialize)]
44pub struct WorkflowQueryRecord {
45    pub value: serde_json::Value,
46    pub published_at: String,
47}
48
49#[derive(Clone, Debug, Serialize, Deserialize)]
50pub struct WorkflowUpdateResponseRecord {
51    pub request_id: String,
52    #[serde(skip_serializing_if = "Option::is_none")]
53    pub name: Option<String>,
54    pub value: serde_json::Value,
55    pub responded_at: String,
56}
57
58#[derive(Clone, Debug, Serialize, Deserialize)]
59pub struct WorkflowMailboxState {
60    #[serde(rename = "_type")]
61    pub type_name: String,
62    pub workflow_id: String,
63    #[serde(default = "default_generation")]
64    pub generation: u64,
65    #[serde(default)]
66    pub continue_as_new_count: u64,
67    #[serde(skip_serializing_if = "Option::is_none")]
68    pub last_continue_as_new_at: Option<String>,
69    #[serde(default)]
70    pub paused: bool,
71    #[serde(default)]
72    pub next_seq: u64,
73    #[serde(default)]
74    pub mailbox: VecDeque<WorkflowMessageRecord>,
75    #[serde(default)]
76    pub queries: BTreeMap<String, WorkflowQueryRecord>,
77    #[serde(default)]
78    pub responses: BTreeMap<String, WorkflowUpdateResponseRecord>,
79}
80
81impl Default for WorkflowMailboxState {
82    fn default() -> Self {
83        Self {
84            type_name: "workflow_mailbox".to_string(),
85            workflow_id: String::new(),
86            generation: default_generation(),
87            continue_as_new_count: 0,
88            last_continue_as_new_at: None,
89            paused: false,
90            next_seq: 0,
91            mailbox: VecDeque::new(),
92            queries: BTreeMap::new(),
93            responses: BTreeMap::new(),
94        }
95    }
96}
97
98fn default_generation() -> u64 {
99    1
100}
101
102#[derive(Clone, Debug, PartialEq, Eq)]
103struct WorkflowTarget {
104    workflow_id: String,
105    base_dir: PathBuf,
106}
107
108fn sanitize_workflow_id(raw: &str) -> String {
109    let trimmed = raw.trim();
110    let mut sanitized = String::with_capacity(trimmed.len());
111    for ch in trimmed.chars() {
112        if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
113            sanitized.push(ch);
114        } else {
115            sanitized.push('_');
116        }
117    }
118    if sanitized.is_empty() || sanitized == "." || sanitized == ".." {
119        "workflow".to_string()
120    } else {
121        sanitized
122    }
123}
124
125fn workflow_base_dir_from_persisted_path(path: &Path) -> PathBuf {
126    for ancestor in path.ancestors() {
127        if ancestor.file_name().and_then(|value| value.to_str()) == Some(".harn-runs") {
128            return non_empty_path_or_dot(ancestor.parent().unwrap_or_else(|| Path::new(".")));
129        }
130    }
131    non_empty_path_or_dot(path.parent().unwrap_or_else(|| Path::new(".")))
132}
133
134fn non_empty_path_or_dot(path: &Path) -> PathBuf {
135    if path.as_os_str().is_empty() {
136        PathBuf::from(".")
137    } else {
138        path.to_path_buf()
139    }
140}
141
142fn workflow_target_root(target: &WorkflowTarget) -> PathBuf {
143    crate::runtime_paths::workflow_dir(&target.base_dir).join(&target.workflow_id)
144}
145
146fn workflow_state_path(target: &WorkflowTarget) -> PathBuf {
147    workflow_target_root(target).join("state.json")
148}
149
150fn now_rfc3339() -> String {
151    crate::clock_mock::active_clock().map_or_else(harn_clock::system_now_rfc3339, |clock| {
152        harn_clock::now_rfc3339(clock.as_ref())
153    })
154}
155
156fn workflow_clock() -> Arc<dyn Clock> {
157    crate::clock_mock::active_clock().unwrap_or_else(harn_clock::RealClock::arc)
158}
159
160struct WorkflowDeadline {
161    clock: Arc<dyn Clock>,
162    started_ms: i64,
163    timeout_ms: u64,
164}
165
166impl WorkflowDeadline {
167    fn new(timeout: StdDuration) -> Self {
168        let clock = workflow_clock();
169        let started_ms = clock.monotonic_ms();
170        Self {
171            clock,
172            started_ms,
173            timeout_ms: u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX),
174        }
175    }
176
177    async fn wait_for_next_poll<F>(&self, changed: F) -> bool
178    where
179        F: Future<Output = ()>,
180    {
181        let elapsed_ms = self
182            .clock
183            .monotonic_ms()
184            .saturating_sub(self.started_ms)
185            .max(0) as u64;
186        let remaining_ms = self.timeout_ms.saturating_sub(elapsed_ms);
187        if remaining_ms == 0 {
188            return false;
189        }
190        tokio::select! {
191            () = changed => {}
192            () = self.clock.sleep(StdDuration::from_millis(
193                remaining_ms.min(WORKFLOW_POLL_INTERVAL_MS),
194            )) => {}
195        }
196        // Auto-advancing test clocks complete sleeps synchronously. Yielding
197        // here preserves worker progress without coupling it to physical time.
198        tokio::task::yield_now().await;
199        true
200    }
201}
202
203fn workflow_notifier(target: &WorkflowTarget) -> Arc<tokio::sync::Notify> {
204    static NOTIFIERS: OnceLock<parking_lot::Mutex<BTreeMap<PathBuf, Weak<tokio::sync::Notify>>>> =
205        OnceLock::new();
206    let mut notifiers = NOTIFIERS.get_or_init(Default::default).lock();
207    let path = workflow_state_path(target);
208    if let Some(notifier) = notifiers.get(&path).and_then(Weak::upgrade) {
209        return notifier;
210    }
211    notifiers.retain(|_, notifier| notifier.strong_count() != 0);
212    let notifier = Arc::new(tokio::sync::Notify::new());
213    notifiers.insert(path, Arc::downgrade(&notifier));
214    notifier
215}
216
217fn load_state(target: &WorkflowTarget) -> Result<WorkflowMailboxState, String> {
218    let path = workflow_state_path(target);
219    if !path.exists() {
220        return Ok(WorkflowMailboxState {
221            workflow_id: target.workflow_id.clone(),
222            ..WorkflowMailboxState::default()
223        });
224    }
225    let text = std::fs::read_to_string(&path)
226        .map_err(|error| format!("workflow state read error: {error}"))?;
227    let mut state: WorkflowMailboxState = serde_json::from_str(&text)
228        .map_err(|error| format!("workflow state parse error: {error}"))?;
229    if state.type_name.is_empty() {
230        state.type_name = "workflow_mailbox".to_string();
231    }
232    if state.workflow_id.is_empty() {
233        state.workflow_id = target.workflow_id.clone();
234    }
235    if state.generation == 0 {
236        state.generation = 1;
237    }
238    Ok(state)
239}
240
241fn save_state(target: &WorkflowTarget, state: &WorkflowMailboxState) -> Result<(), String> {
242    let path = workflow_state_path(target);
243    let json = serde_json::to_string_pretty(state)
244        .map_err(|error| format!("workflow state encode error: {error}"))?;
245    crate::atomic_io::atomic_write(&path, json.as_bytes())
246        .map_err(|error| format!("workflow state write error: {error}"))?;
247    workflow_notifier(target).notify_one();
248    Ok(())
249}
250
251fn parse_target_json(
252    value: &serde_json::Value,
253    fallback_base_dir: Option<&Path>,
254) -> Option<WorkflowTarget> {
255    match value {
256        serde_json::Value::String(text) => Some(WorkflowTarget {
257            workflow_id: sanitize_workflow_id(text),
258            base_dir: fallback_base_dir
259                .map(Path::to_path_buf)
260                .unwrap_or_else(runtime_root_base),
261        }),
262        serde_json::Value::Object(map) => {
263            let workflow_id = map
264                .get("workflow_id")
265                .and_then(|value| value.as_str())
266                .or_else(|| map.get("workflow").and_then(|value| value.as_str()))
267                .or_else(|| {
268                    map.get("run")
269                        .and_then(|value| value.get("workflow_id"))
270                        .and_then(|value| value.as_str())
271                })
272                .or_else(|| {
273                    map.get("result")
274                        .and_then(|value| value.get("run"))
275                        .and_then(|value| value.get("workflow_id"))
276                        .and_then(|value| value.as_str())
277                })?;
278            let explicit_base = map
279                .get("base_dir")
280                .and_then(|value| value.as_str())
281                .filter(|value| !value.trim().is_empty())
282                .map(PathBuf::from);
283            let persisted_path = map
284                .get("persisted_path")
285                .and_then(|value| value.as_str())
286                .or_else(|| map.get("path").and_then(|value| value.as_str()))
287                .or_else(|| {
288                    map.get("run")
289                        .and_then(|value| value.get("persisted_path"))
290                        .and_then(|value| value.as_str())
291                })
292                .or_else(|| {
293                    map.get("result")
294                        .and_then(|value| value.get("run"))
295                        .and_then(|value| value.get("persisted_path"))
296                        .and_then(|value| value.as_str())
297                });
298            let base_dir = explicit_base
299                .or_else(|| {
300                    persisted_path
301                        .map(|path| workflow_base_dir_from_persisted_path(Path::new(path)))
302                })
303                .or_else(|| fallback_base_dir.map(Path::to_path_buf))
304                .unwrap_or_else(runtime_root_base);
305            Some(WorkflowTarget {
306                workflow_id: sanitize_workflow_id(workflow_id),
307                base_dir,
308            })
309        }
310        _ => None,
311    }
312}
313
314fn parse_target_vm(
315    value: Option<&VmValue>,
316    fallback_base_dir: Option<&Path>,
317    builtin: &str,
318) -> Result<WorkflowTarget, VmError> {
319    let value = value.ok_or_else(|| VmError::Runtime(format!("{builtin}: missing target")))?;
320    parse_target_json(&crate::llm::vm_value_to_json(value), fallback_base_dir).ok_or_else(|| {
321        VmError::Runtime(format!(
322            "{builtin}: target must be a workflow id string or dict with workflow_id/workflow"
323        ))
324    })
325}
326
327fn workflow_status_json(
328    target: &WorkflowTarget,
329    state: &WorkflowMailboxState,
330) -> serde_json::Value {
331    serde_json::json!({
332        "workflow_id": target.workflow_id,
333        "base_dir": target.base_dir.to_string_lossy(),
334        "generation": state.generation,
335        "paused": state.paused,
336        "pending_count": state.mailbox.len(),
337        "query_count": state.queries.len(),
338        "response_count": state.responses.len(),
339        "continue_as_new_count": state.continue_as_new_count,
340        "last_continue_as_new_at": state.last_continue_as_new_at,
341    })
342}
343
344fn target_for_base(base_dir: &Path, workflow_id: &str) -> WorkflowTarget {
345    WorkflowTarget {
346        workflow_id: sanitize_workflow_id(workflow_id),
347        base_dir: base_dir.to_path_buf(),
348    }
349}
350
351fn enqueue_message(
352    target: &WorkflowTarget,
353    kind: &str,
354    name: &str,
355    payload: serde_json::Value,
356    request_id: Option<String>,
357) -> Result<serde_json::Value, String> {
358    let mut state = load_state(target)?;
359    let message = push_message(&mut state, kind, name, payload, request_id);
360    save_state(target, &state)?;
361    Ok(serde_json::json!({
362        "workflow_id": target.workflow_id,
363        "message": message,
364        "status": workflow_status_json(target, &state),
365    }))
366}
367
368fn push_message(
369    state: &mut WorkflowMailboxState,
370    kind: &str,
371    name: &str,
372    payload: serde_json::Value,
373    request_id: Option<String>,
374) -> WorkflowMessageRecord {
375    state.next_seq += 1;
376    let message = WorkflowMessageRecord {
377        seq: state.next_seq,
378        kind: kind.to_string(),
379        name: name.to_string(),
380        request_id,
381        payload,
382        enqueued_at: now_rfc3339(),
383    };
384    state.mailbox.push_back(message.clone());
385    message
386}
387
388fn receive_message(target: &WorkflowTarget) -> Result<Option<WorkflowMessageRecord>, String> {
389    let mut state = load_state(target)?;
390    let message = state.mailbox.pop_front();
391    if message.is_some() {
392        save_state(target, &state)?;
393    }
394    Ok(message)
395}
396
397async fn receive_message_with_timeout(
398    target: &WorkflowTarget,
399    timeout: Option<StdDuration>,
400) -> Result<Option<WorkflowMessageRecord>, String> {
401    let Some(timeout) = timeout else {
402        return receive_message(target);
403    };
404    let deadline = WorkflowDeadline::new(timeout);
405    let notifier = workflow_notifier(target);
406    loop {
407        if let Some(message) = receive_message(target)? {
408            return Ok(Some(message));
409        }
410        if !deadline.wait_for_next_poll(notifier.notified()).await {
411            return Ok(None);
412        }
413    }
414}
415
416pub fn workflow_signal_for_base(
417    base_dir: &Path,
418    workflow_id: &str,
419    name: &str,
420    payload: serde_json::Value,
421) -> Result<serde_json::Value, String> {
422    let target = target_for_base(base_dir, workflow_id);
423    enqueue_message(&target, "signal", name, payload, None)
424}
425
426pub fn workflow_query_for_base(
427    base_dir: &Path,
428    workflow_id: &str,
429    name: &str,
430) -> Result<serde_json::Value, String> {
431    let target = target_for_base(base_dir, workflow_id);
432    let state = load_state(&target)?;
433    Ok(state
434        .queries
435        .get(name)
436        .map(|record| record.value.clone())
437        .unwrap_or(serde_json::Value::Null))
438}
439
440pub fn workflow_publish_query_for_base(
441    base_dir: &Path,
442    workflow_id: &str,
443    name: &str,
444    value: serde_json::Value,
445) -> Result<serde_json::Value, String> {
446    let target = target_for_base(base_dir, workflow_id);
447    let mut state = load_state(&target)?;
448    state.queries.insert(
449        name.to_string(),
450        WorkflowQueryRecord {
451            value,
452            published_at: now_rfc3339(),
453        },
454    );
455    save_state(&target, &state)?;
456    Ok(workflow_status_json(&target, &state))
457}
458
459pub fn workflow_pause_for_base(
460    base_dir: &Path,
461    workflow_id: &str,
462) -> Result<serde_json::Value, String> {
463    let target = target_for_base(base_dir, workflow_id);
464    let mut state = load_state(&target)?;
465    state.paused = true;
466    push_message(&mut state, "control", "pause", serde_json::json!({}), None);
467    save_state(&target, &state)?;
468    Ok(workflow_status_json(&target, &state))
469}
470
471pub fn workflow_resume_for_base(
472    base_dir: &Path,
473    workflow_id: &str,
474) -> Result<serde_json::Value, String> {
475    let target = target_for_base(base_dir, workflow_id);
476    let mut state = load_state(&target)?;
477    state.paused = false;
478    push_message(&mut state, "control", "resume", serde_json::json!({}), None);
479    save_state(&target, &state)?;
480    Ok(workflow_status_json(&target, &state))
481}
482
483pub async fn workflow_update_for_base(
484    base_dir: &Path,
485    workflow_id: &str,
486    name: &str,
487    payload: serde_json::Value,
488    timeout: StdDuration,
489) -> Result<serde_json::Value, String> {
490    let target = target_for_base(base_dir, workflow_id);
491    let request_id = enqueue_update_request(&target, name, payload)?;
492    wait_for_update_response(&target, name, &request_id, timeout).await
493}
494
495fn enqueue_update_request(
496    target: &WorkflowTarget,
497    name: &str,
498    payload: serde_json::Value,
499) -> Result<String, String> {
500    let request_id = uuid::Uuid::now_v7().to_string();
501    enqueue_message(target, "update", name, payload, Some(request_id.clone()))?;
502    Ok(request_id)
503}
504
505async fn wait_for_update_response(
506    target: &WorkflowTarget,
507    name: &str,
508    request_id: &str,
509    timeout: StdDuration,
510) -> Result<serde_json::Value, String> {
511    let deadline = WorkflowDeadline::new(timeout);
512    let notifier = workflow_notifier(target);
513    loop {
514        match update_response_value(target, request_id) {
515            Ok(Some(value)) => return Ok(value),
516            Ok(None) => {}
517            Err(error) => return Err(error),
518        }
519
520        if !deadline.wait_for_next_poll(notifier.notified()).await {
521            break;
522        }
523    }
524    Err(format!(
525        "workflow update '{name}' timed out for '{}'",
526        target.workflow_id
527    ))
528}
529
530fn update_response_value(
531    target: &WorkflowTarget,
532    request_id: &str,
533) -> Result<Option<serde_json::Value>, String> {
534    Ok(load_state(target)?
535        .responses
536        .get(request_id)
537        .map(|response| response.value.clone()))
538}
539
540pub fn workflow_respond_update_for_base(
541    base_dir: &Path,
542    workflow_id: &str,
543    request_id: &str,
544    name: Option<&str>,
545    value: serde_json::Value,
546) -> Result<serde_json::Value, String> {
547    let target = target_for_base(base_dir, workflow_id);
548    record_update_response(&target, request_id, name, value)
549}
550
551fn record_update_response(
552    target: &WorkflowTarget,
553    request_id: &str,
554    name: Option<&str>,
555    value: serde_json::Value,
556) -> Result<serde_json::Value, String> {
557    let mut state = load_state(target)?;
558    state.responses.insert(
559        request_id.to_string(),
560        WorkflowUpdateResponseRecord {
561            request_id: request_id.to_string(),
562            name: name.map(ToString::to_string),
563            value,
564            responded_at: now_rfc3339(),
565        },
566    );
567    save_state(target, &state)?;
568    Ok(workflow_status_json(target, &state))
569}
570
571pub(crate) fn register_workflow_message_builtins(vm: &mut Vm) {
572    vm.set_global(
573        "workflow",
574        VmValue::dict(BTreeMap::from([
575            (
576                "signal".to_string(),
577                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.signal")),
578            ),
579            (
580                "query".to_string(),
581                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.query")),
582            ),
583            (
584                "update".to_string(),
585                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.update")),
586            ),
587            (
588                "publish_query".to_string(),
589                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.publish_query")),
590            ),
591            (
592                "receive".to_string(),
593                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.receive")),
594            ),
595            (
596                "respond_update".to_string(),
597                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.respond_update")),
598            ),
599            (
600                "pause".to_string(),
601                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.pause")),
602            ),
603            (
604                "resume".to_string(),
605                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.resume")),
606            ),
607            (
608                "status".to_string(),
609                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.status")),
610            ),
611            (
612                "continue_as_new".to_string(),
613                VmValue::BuiltinRef(arcstr::ArcStr::from("workflow.continue_as_new")),
614            ),
615        ])),
616    );
617
618    for def in MODULE_BUILTINS {
619        vm.register_builtin_def(def);
620    }
621}
622
623/// Enqueue a workflow signal message.
624#[harn_builtin(
625    exposure = "harness.runtime.workflow_signal",
626    effects = ["state.write@dynamic"],
627    sig = "workflow.signal(target: string|dict, name: string, payload?: any) -> dict",
628    category = "workflow.messages"
629)]
630fn workflow_signal_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
631    let target = parse_target_vm(args.first(), None, "workflow.signal")?;
632    let name = args
633        .get(1)
634        .map(|value| value.display())
635        .filter(|value| !value.is_empty())
636        .ok_or_else(|| VmError::Runtime("workflow.signal: missing name".to_string()))?;
637    let payload = args
638        .get(2)
639        .map(crate::llm::vm_value_to_json)
640        .unwrap_or(serde_json::Value::Null);
641    let result =
642        enqueue_message(&target, "signal", &name, payload, None).map_err(VmError::Runtime)?;
643    Ok(crate::stdlib::json_to_vm_value(&result))
644}
645
646/// Read the latest published workflow query value.
647#[harn_builtin(
648    exposure = "harness.runtime.workflow_query",
649    effects = ["state.read@dynamic"],
650    sig = "workflow.query(target: string|dict, name: string) -> any",
651    category = "workflow.messages"
652)]
653fn workflow_query_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
654    let target = parse_target_vm(args.first(), None, "workflow.query")?;
655    let name = args
656        .get(1)
657        .map(|value| value.display())
658        .filter(|value| !value.is_empty())
659        .ok_or_else(|| VmError::Runtime("workflow.query: missing name".to_string()))?;
660    let state = load_state(&target).map_err(VmError::Runtime)?;
661    Ok(crate::stdlib::json_to_vm_value(
662        &state
663            .queries
664            .get(&name)
665            .map(|record| record.value.clone())
666            .unwrap_or(serde_json::Value::Null),
667    ))
668}
669
670/// Enqueue a workflow update and wait for a response.
671#[harn_builtin(
672    exposure = "harness.runtime.workflow_update",
673    effects = ["state.mutate@dynamic"],
674    sig = "workflow.update(target: string|dict, name: string, payload?: any, options?: dict|nil) -> any",
675    kind = "async",
676    category = "workflow.messages"
677)]
678async fn workflow_update_builtin(
679    _ctx: crate::vm::AsyncBuiltinCtx,
680    args: Vec<VmValue>,
681) -> Result<VmValue, VmError> {
682    let target = parse_target_vm(args.first(), None, "workflow.update")?;
683    let name = args
684        .get(1)
685        .map(|value| value.display())
686        .filter(|value| !value.is_empty())
687        .ok_or_else(|| VmError::Runtime("workflow.update: missing name".to_string()))?;
688    let payload = args
689        .get(2)
690        .map(crate::llm::vm_value_to_json)
691        .unwrap_or(serde_json::Value::Null);
692    let timeout_ms = args
693        .get(3)
694        .and_then(|value| value.as_dict())
695        .and_then(|dict| dict.get("timeout_ms"))
696        .and_then(VmValue::as_int)
697        .unwrap_or(DEFAULT_UPDATE_TIMEOUT_MS as i64)
698        .max(1) as u64;
699    let result = workflow_update_for_base(
700        &target.base_dir,
701        &target.workflow_id,
702        &name,
703        payload,
704        StdDuration::from_millis(timeout_ms),
705    )
706    .await
707    .map_err(VmError::Runtime)?;
708    Ok(crate::stdlib::json_to_vm_value(&result))
709}
710
711/// Publish a workflow query value.
712#[harn_builtin(
713    exposure = "harness.runtime.workflow_publish_query",
714    effects = ["state.write@dynamic"],
715    sig = "workflow.publish_query(target: string|dict, name: string, value?: any) -> dict",
716    category = "workflow.messages"
717)]
718fn workflow_publish_query_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
719    let target = parse_target_vm(args.first(), None, "workflow.publish_query")?;
720    let name = args
721        .get(1)
722        .map(|value| value.display())
723        .filter(|value| !value.is_empty())
724        .ok_or_else(|| VmError::Runtime("workflow.publish_query: missing name".to_string()))?;
725    let value = args
726        .get(2)
727        .map(crate::llm::vm_value_to_json)
728        .unwrap_or(serde_json::Value::Null);
729    let result =
730        workflow_publish_query_for_base(&target.base_dir, &target.workflow_id, &name, value)
731            .map_err(VmError::Runtime)?;
732    Ok(crate::stdlib::json_to_vm_value(&result))
733}
734
735/// Receive the next workflow mailbox message, optionally waiting for one.
736#[harn_builtin(
737    exposure = "harness.runtime.workflow_receive",
738    effects = ["state.observe@dynamic", "state.mutate@dynamic"],
739    sig = "workflow.receive(target: string | dict, timeout_ms?: int) -> {workflow_id: string, seq: int, kind: \"signal\" | \"update\" | \"control\", name: string, request_id: string?, payload: unknown, enqueued_at: string}?",
740    kind = "async",
741    category = "workflow.messages"
742)]
743async fn workflow_receive_builtin(
744    _ctx: crate::vm::AsyncBuiltinCtx,
745    args: Vec<VmValue>,
746) -> Result<VmValue, VmError> {
747    let target = parse_target_vm(args.first(), None, "workflow.receive")?;
748    let timeout = match args.get(1) {
749        None | Some(VmValue::Nil) => None,
750        Some(value) => {
751            let timeout_ms = value.as_int().ok_or_else(|| {
752                VmError::Runtime("workflow.receive: `timeout_ms` must be an int".to_string())
753            })?;
754            if timeout_ms <= 0 {
755                return Err(VmError::Runtime(
756                    "workflow.receive: `timeout_ms` must be positive".to_string(),
757                ));
758            }
759            Some(StdDuration::from_millis(timeout_ms as u64))
760        }
761    };
762    let Some(message) = receive_message_with_timeout(&target, timeout)
763        .await
764        .map_err(VmError::Runtime)?
765    else {
766        return Ok(VmValue::Nil);
767    };
768    Ok(crate::stdlib::json_to_vm_value(&serde_json::json!({
769        "workflow_id": target.workflow_id,
770        "seq": message.seq,
771        "kind": message.kind,
772        "name": message.name,
773        "request_id": message.request_id,
774        "payload": message.payload,
775        "enqueued_at": message.enqueued_at,
776    })))
777}
778
779/// Respond to a pending workflow update request.
780#[harn_builtin(
781    exposure = "harness.runtime.workflow_respond_update",
782    effects = ["state.write@dynamic"],
783    sig = "workflow.respond_update(target: string|dict, request_id: string, value?: any, name?: string|nil) -> dict",
784    category = "workflow.messages"
785)]
786fn workflow_respond_update_builtin(
787    args: &[VmValue],
788    _out: &mut String,
789) -> Result<VmValue, VmError> {
790    let target = parse_target_vm(args.first(), None, "workflow.respond_update")?;
791    let request_id = args
792        .get(1)
793        .map(|value| value.display())
794        .filter(|value| !value.is_empty())
795        .ok_or_else(|| {
796            VmError::Runtime("workflow.respond_update: missing request id".to_string())
797        })?;
798    let value = args
799        .get(2)
800        .map(crate::llm::vm_value_to_json)
801        .unwrap_or(serde_json::Value::Null);
802    let name = args
803        .get(3)
804        .map(|value| value.display())
805        .filter(|value| !value.is_empty());
806    let result = workflow_respond_update_for_base(
807        &target.base_dir,
808        &target.workflow_id,
809        &request_id,
810        name.as_deref(),
811        value,
812    )
813    .map_err(VmError::Runtime)?;
814    Ok(crate::stdlib::json_to_vm_value(&result))
815}
816
817/// Pause a workflow mailbox.
818#[harn_builtin(
819    exposure = "harness.runtime.workflow_pause",
820    effects = ["state.mutate@dynamic"],
821    sig = "workflow.pause(target: string|dict) -> dict",
822    category = "workflow.messages"
823)]
824fn workflow_pause_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
825    let target = parse_target_vm(args.first(), None, "workflow.pause")?;
826    let result =
827        workflow_pause_for_base(&target.base_dir, &target.workflow_id).map_err(VmError::Runtime)?;
828    Ok(crate::stdlib::json_to_vm_value(&result))
829}
830
831/// Resume a workflow mailbox.
832#[harn_builtin(
833    exposure = "harness.runtime.workflow_resume",
834    effects = ["state.mutate@dynamic"],
835    sig = "workflow.resume(target: string|dict) -> dict",
836    category = "workflow.messages"
837)]
838fn workflow_resume_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
839    let target = parse_target_vm(args.first(), None, "workflow.resume")?;
840    let result = workflow_resume_for_base(&target.base_dir, &target.workflow_id)
841        .map_err(VmError::Runtime)?;
842    Ok(crate::stdlib::json_to_vm_value(&result))
843}
844
845/// Return workflow mailbox status.
846#[harn_builtin(
847    exposure = "harness.runtime.workflow_status",
848    effects = ["state.read@dynamic"],
849    sig = "workflow.status(target: string|dict) -> dict",
850    category = "workflow.messages"
851)]
852fn workflow_status_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
853    let target = parse_target_vm(args.first(), None, "workflow.status")?;
854    let state = load_state(&target).map_err(VmError::Runtime)?;
855    Ok(crate::stdlib::json_to_vm_value(&workflow_status_json(
856        &target, &state,
857    )))
858}
859
860/// Advance a workflow mailbox generation.
861#[harn_builtin(
862    exposure = "harness.runtime.workflow_continue_as_new",
863    effects = ["state.mutate@dynamic"],
864    sig = "workflow.continue_as_new(target: string|dict) -> dict",
865    category = "workflow.messages"
866)]
867fn workflow_continue_as_new_builtin(
868    args: &[VmValue],
869    _out: &mut String,
870) -> Result<VmValue, VmError> {
871    continue_as_new_for_label(args, "workflow.continue_as_new")
872}
873
874/// Advance a workflow mailbox generation (top-level alias).
875#[harn_builtin(
876    exposure = "harness.runtime.continue_as_new",
877    effects = ["state.mutate@dynamic"],
878    sig = "continue_as_new(target: string|dict) -> dict",
879    category = "workflow.messages"
880)]
881fn continue_as_new_builtin(args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
882    continue_as_new_for_label(args, "continue_as_new")
883}
884
885fn continue_as_new_for_label(args: &[VmValue], label: &str) -> Result<VmValue, VmError> {
886    let target = parse_target_vm(args.first(), None, label)?;
887    let mut state = load_state(&target).map_err(VmError::Runtime)?;
888    state.generation += 1;
889    state.continue_as_new_count += 1;
890    state.last_continue_as_new_at = Some(now_rfc3339());
891    state.responses.clear();
892    save_state(&target, &state).map_err(VmError::Runtime)?;
893    Ok(crate::stdlib::json_to_vm_value(&workflow_status_json(
894        &target, &state,
895    )))
896}
897
898#[cfg(test)]
899mod tests {
900    use super::*;
901    use std::task::Poll;
902
903    #[tokio::test]
904    async fn mailbox_timestamp_uses_active_harness_clock() {
905        let dir = tempfile::tempdir().expect("tempdir");
906        let target = target_for_base(dir.path(), "wf-clock");
907        let clock = crate::clock_mock::MockClock::at_wall_ms(1_700_000_000_000);
908
909        let message = crate::clock_mock::scope_capability_clock(clock, async {
910            enqueue_message(&target, "signal", "ready", serde_json::Value::Null, None)
911                .expect("enqueue signal");
912            receive_message(&target)
913                .expect("receive signal")
914                .expect("queued signal")
915        })
916        .await;
917
918        assert_eq!(message.enqueued_at, "2023-11-14T22:13:20Z");
919    }
920
921    #[tokio::test(start_paused = true)]
922    async fn update_round_trip_waits_for_response() {
923        let dir = tempfile::tempdir().expect("tempdir");
924        let workflow_id = "wf-update";
925        let base_dir = dir.path().to_path_buf();
926        let target = target_for_base(&base_dir, workflow_id);
927        let request_id =
928            enqueue_update_request(&target, "adjust_budget", serde_json::json!({"max_usd": 10}))
929                .expect("enqueue update");
930
931        let message = receive_message(&target)
932            .expect("receive queued update")
933            .expect("queued update");
934        assert_eq!(message.kind, "update");
935        assert_eq!(message.name, "adjust_budget");
936        assert_eq!(message.request_id.as_deref(), Some(request_id.as_str()));
937        assert_eq!(
938            update_response_value(&target, &request_id).expect("read response"),
939            None
940        );
941
942        let waiter = wait_for_update_response(
943            &target,
944            "adjust_budget",
945            &request_id,
946            StdDuration::from_millis(500),
947        );
948        tokio::pin!(waiter);
949        assert!(matches!(futures::poll!(&mut waiter), Poll::Pending));
950
951        workflow_respond_update_for_base(
952            &base_dir,
953            workflow_id,
954            &request_id,
955            Some("adjust_budget"),
956            serde_json::json!({"ok": true}),
957        )
958        .expect("save response");
959        assert_eq!(
960            update_response_value(&target, &request_id).expect("read response"),
961            Some(serde_json::json!({"ok": true}))
962        );
963        tokio::time::advance(StdDuration::from_millis(WORKFLOW_POLL_INTERVAL_MS)).await;
964
965        let result = waiter.await.expect("update result");
966        assert_eq!(result, serde_json::json!({"ok": true}));
967    }
968
969    #[tokio::test(start_paused = true)]
970    async fn update_wait_respects_short_timeout() {
971        let dir = tempfile::tempdir().expect("tempdir");
972        let target = target_for_base(dir.path(), "wf-timeout");
973        let request_id =
974            enqueue_update_request(&target, "adjust_budget", serde_json::json!({"max_usd": 10}))
975                .expect("enqueue update");
976
977        let waiter = wait_for_update_response(
978            &target,
979            "adjust_budget",
980            &request_id,
981            StdDuration::from_millis(10),
982        );
983        tokio::pin!(waiter);
984        assert!(matches!(futures::poll!(&mut waiter), Poll::Pending));
985
986        tokio::time::advance(StdDuration::from_millis(9)).await;
987        assert!(matches!(futures::poll!(&mut waiter), Poll::Pending));
988
989        tokio::time::advance(StdDuration::from_millis(1)).await;
990        let err = waiter.await.expect_err("update should time out");
991        assert!(err.contains("timed out"));
992    }
993
994    #[tokio::test]
995    async fn update_wait_propagates_state_errors() {
996        let dir = tempfile::tempdir().expect("tempdir");
997        let target = target_for_base(dir.path(), "wf-corrupt");
998        std::fs::create_dir_all(workflow_target_root(&target)).expect("state dir");
999        std::fs::write(workflow_state_path(&target), "{not json").expect("state write");
1000
1001        let err = wait_for_update_response(
1002            &target,
1003            "adjust_budget",
1004            "request-1",
1005            StdDuration::from_millis(10),
1006        )
1007        .await
1008        .expect_err("corrupt state should fail immediately");
1009        assert!(err.contains("workflow state parse error"));
1010    }
1011
1012    #[test]
1013    fn workflow_ids_preserve_namespace_without_path_segments() {
1014        assert_eq!(
1015            sanitize_workflow_id("workflow://local/start-my-day"),
1016            "workflow___local_start-my-day"
1017        );
1018        assert_eq!(sanitize_workflow_id("../start-my-day"), ".._start-my-day");
1019        assert_eq!(sanitize_workflow_id(".."), "workflow");
1020    }
1021
1022    #[test]
1023    fn persisted_path_drives_target_base_dir() {
1024        let base = parse_target_json(
1025            &serde_json::json!({
1026                "workflow_id": "wf",
1027                "persisted_path": "/tmp/demo/.harn-runs/run.json"
1028            }),
1029            None,
1030        )
1031        .expect("target");
1032        assert_eq!(base.workflow_id, "wf");
1033        assert_eq!(base.base_dir, PathBuf::from("/tmp/demo"));
1034    }
1035
1036    #[test]
1037    fn nested_persisted_path_drives_target_base_dir() {
1038        let base = parse_target_json(
1039            &serde_json::json!({
1040                "workflow_id": "wf",
1041                "persisted_path": "/tmp/demo/.harn-runs/session/run.json"
1042            }),
1043            None,
1044        )
1045        .expect("target");
1046        assert_eq!(base.base_dir, PathBuf::from("/tmp/demo"));
1047
1048        let relative = parse_target_json(
1049            &serde_json::json!({
1050                "workflow_id": "wf",
1051                "persisted_path": ".harn-runs/session/run.json"
1052            }),
1053            None,
1054        )
1055        .expect("target");
1056        assert_eq!(relative.base_dir, PathBuf::from("."));
1057    }
1058}