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 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#[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#[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#[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#[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#[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#[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#[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#[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#[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#[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#[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}