Skip to main content

harn_vm/stdlib/
workflow_messages.rs

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