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 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(¬ifier));
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#[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#[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#[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#[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#[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#[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#[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#[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#[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#[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#[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}