beam_core/workflow_run/
bootstrap.rs1use std::collections::BTreeMap;
2use std::fs;
3use std::path::Path;
4
5use anyhow::{Context, Result};
6use serde::{Deserialize, Serialize};
7use serde_json::Value;
8use sha2::{Digest, Sha256};
9
10use super::validation::normalize_workflow_params;
11use crate::{BeamPaths, EventDraft, EventLog, WorkflowActor, parse_workflow_definition};
12
13#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
14#[serde(rename_all = "camelCase")]
15pub struct RunChatBinding {
16 pub chat_id: String,
17 pub lark_app_id: String,
18}
19
20#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
21#[serde(rename_all = "camelCase")]
22pub struct WorkflowOutputRef {
23 pub output_hash: String,
24 pub output_path: String,
25 pub output_bytes: usize,
26 pub output_schema_version: u32,
27 #[serde(default, skip_serializing_if = "Option::is_none")]
28 pub content_type: Option<String>,
29}
30
31#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
32pub struct WorkflowRunBootstrap {
33 pub run_id: String,
34 pub workflow_id: String,
35 pub revision_id: String,
36 pub input_ref: WorkflowOutputRef,
37}
38
39#[derive(Debug, Clone)]
40pub struct BootstrapWorkflowRunInput<'a> {
41 pub run_id: &'a str,
42 pub workflow_json: &'a str,
43 pub expected_workflow_id: Option<&'a str>,
44 pub params: &'a BTreeMap<String, Value>,
45 pub initiator: &'a str,
46 pub chat_binding: Option<RunChatBinding>,
47}
48
49pub fn mint_workflow_run_id(workflow_id: &str, now_ms: u64) -> String {
50 let safe = workflow_id
51 .chars()
52 .map(|ch| match ch {
53 'a'..='z' | 'A'..='Z' | '0'..='9' | '.' | '_' | '-' => ch,
54 _ => '_',
55 })
56 .collect::<String>();
57 format!("{}-{}", safe, now_ms)
58}
59
60pub fn bootstrap_workflow_run(
61 paths: &BeamPaths,
62 input: BootstrapWorkflowRunInput<'_>,
63) -> Result<WorkflowRunBootstrap> {
64 let workflow = parse_workflow_definition(input.workflow_json)?;
65 let workflow_id = workflow.workflow_id.clone();
66 if let Some(expected) = input.expected_workflow_id {
67 if expected != workflow_id {
68 anyhow::bail!(
69 "workflowId mismatch: requested={} file={}",
70 expected,
71 workflow_id
72 );
73 }
74 }
75
76 let normalized_params = normalize_workflow_params(&workflow, input.params)?;
78
79 let run_dir = paths.workflow_run_dir(input.run_id);
80 fs::create_dir_all(run_dir.join("blobs"))?;
81 fs::write(run_dir.join("workflow.json"), input.workflow_json)?;
82 if let Some(binding) = input.chat_binding {
83 fs::write(
84 run_dir.join("chat-binding.json"),
85 serde_json::to_vec_pretty(&binding)?,
86 )?;
87 }
88
89 let params_json = serde_json::to_vec(&normalized_params)?;
90 let params_hash = sha256_hex(¶ms_json);
91 let input_path = run_dir.join("blobs").join(¶ms_hash);
92 fs::write(&input_path, ¶ms_json)?;
93 let input_ref = WorkflowOutputRef {
94 output_hash: format!("sha256:{}", params_hash),
95 output_path: input_path.display().to_string(),
96 output_bytes: params_json.len(),
97 output_schema_version: 1,
98 content_type: Some("application/json".to_string()),
99 };
100
101 let mut log = EventLog::new(input.run_id.to_string(), paths.workflow_runs_dir())?;
102 let revision_id = sha256_hex(&serde_json::to_vec(&workflow)?);
103 let _run_created = log.append(EventDraft {
104 event_type: "runCreated".to_string(),
105 actor: WorkflowActor::System,
106 payload: serde_json::json!({
107 "workflowId": workflow_id,
108 "revisionId": revision_id,
109 "inputRef": input_ref,
110 "initiator": input.initiator,
111 }),
112 timestamp: None,
113 payload_hash: None,
114 })?;
115 let _run_started = log.append(EventDraft {
116 event_type: "runStarted".to_string(),
117 actor: WorkflowActor::Scheduler,
118 payload: serde_json::json!({}),
119 timestamp: None,
120 payload_hash: None,
121 })?;
122
123 Ok(WorkflowRunBootstrap {
124 run_id: input.run_id.to_string(),
125 workflow_id,
126 revision_id,
127 input_ref,
128 })
129}
130
131pub(crate) fn sha256_hex(bytes: &[u8]) -> String {
139 let mut hasher = Sha256::new();
140 hasher.update(bytes);
141 lower_hex(&hasher.finalize())
142}
143
144pub(crate) fn lower_hex(bytes: &[u8]) -> String {
145 const HEX: &[u8; 16] = b"0123456789abcdef";
146 let mut out = String::with_capacity(bytes.len() * 2);
147 for byte in bytes {
148 out.push(HEX[(byte >> 4) as usize] as char);
149 out.push(HEX[(byte & 0x0f) as usize] as char);
150 }
151 out
152}
153
154pub fn read_workflow_definition_from_path(path: &Path) -> Result<String> {
155 Ok(fs::read_to_string(path).with_context(|| format!("读取 {} 失败", path.display()))?)
156}