Skip to main content

beam_core/workflow_run/
bootstrap.rs

1use 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    // Validate / normalize params before creating any run artifacts
77    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(&params_json);
91    let input_path = run_dir.join("blobs").join(&params_hash);
92    fs::write(&input_path, &params_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
131/// Normalize and validate workflow parameters against the definition's params
132/// schema. Returns a canonical `BTreeMap<String, Value>` suitable for writing
133/// as the run input blob.
134///
135/// If the workflow has no `params` definition (or an empty one), any supplied
136/// params are rejected — no parameters are allowed without a declaration.
137
138pub(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}