Skip to main content

a3s_code_core/research/
workflow.rs

1//! Research workflow-plan bridge over existing execution result receipts.
2
3use super::{
4    digest, validate_digest_field, validate_id, ResearchContractError, RESEARCH_MAX_DIGESTS,
5};
6use crate::execution_identity::ExecutionResultReceiptV1;
7use serde::{Deserialize, Serialize};
8use std::collections::{BTreeMap, BTreeSet};
9
10pub const RESEARCH_WORKFLOW_STEP_SCHEMA_V1: &str = "a3s.code.research-workflow-step.v1";
11pub const RESEARCH_WORKFLOW_PLAN_SCHEMA_V1: &str = "a3s.code.research-workflow-plan.v1";
12pub const RESEARCH_RERUN_LINEAGE_SCHEMA_V1: &str = "a3s.code.research-rerun-lineage.v1";
13pub const RESEARCH_MAX_WORKFLOW_STEPS: usize = 512;
14const RESEARCH_WORKFLOW_STEP_DIGEST_DOMAIN: &str = "a3s.code.research-workflow-step.identity.v1";
15const RESEARCH_WORKFLOW_PLAN_DIGEST_DOMAIN: &str = "a3s.code.research-workflow-plan.identity.v1";
16const RESEARCH_WORKFLOW_RECEIPT_DIGEST_DOMAIN: &str =
17    "a3s.code.research-workflow-receipt.identity.v1";
18const RESEARCH_RERUN_LINEAGE_DIGEST_DOMAIN: &str = "a3s.code.research-rerun-lineage.identity.v1";
19
20/// One research-visible workflow step bound to digests and optional receipts.
21///
22/// Code does not schedule the step. It only records the exact identity and
23/// artifact/input digests so a later review finding can select a minimal
24/// re-run set without rewriting the parent evidence ledger.
25#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
26#[serde(rename_all = "camelCase", deny_unknown_fields)]
27pub struct ResearchWorkflowStepV1 {
28    pub schema: String,
29    pub step_id: String,
30    pub project_id: String,
31    pub run_id: String,
32    pub workflow_digest: String,
33    pub step_identity_digest: String,
34    #[serde(default, skip_serializing_if = "Option::is_none")]
35    pub result_receipt_digest: Option<String>,
36    pub input_digests: Vec<String>,
37    pub output_artifact_digests: Vec<String>,
38    pub depends_on: Vec<String>,
39    pub step_digest: String,
40}
41
42impl ResearchWorkflowStepV1 {
43    #[allow(clippy::too_many_arguments)]
44    pub fn new(
45        step_id: impl Into<String>,
46        project_id: impl Into<String>,
47        run_id: impl Into<String>,
48        workflow_digest: impl Into<String>,
49        step_identity_digest: impl Into<String>,
50        mut input_digests: Vec<String>,
51        mut output_artifact_digests: Vec<String>,
52        mut depends_on: Vec<String>,
53    ) -> Result<Self, ResearchContractError> {
54        input_digests.sort();
55        input_digests.dedup();
56        output_artifact_digests.sort();
57        output_artifact_digests.dedup();
58        depends_on.sort();
59        depends_on.dedup();
60        let mut step = Self {
61            schema: RESEARCH_WORKFLOW_STEP_SCHEMA_V1.to_owned(),
62            step_id: step_id.into(),
63            project_id: project_id.into(),
64            run_id: run_id.into(),
65            workflow_digest: workflow_digest.into(),
66            step_identity_digest: step_identity_digest.into(),
67            result_receipt_digest: None,
68            input_digests,
69            output_artifact_digests,
70            depends_on,
71            step_digest: String::new(),
72        };
73        step.validate_without_digest()?;
74        step.step_digest = step.expected_digest()?;
75        Ok(step)
76    }
77
78    /// Bind the exact execution result receipt produced for this step.
79    ///
80    /// The receipt identity digest must match the admitted step identity. The
81    /// receipt evidence digest must appear in the step inputs so a reviewer
82    /// cannot attach an unrelated terminal outcome.
83    pub fn bind_result_receipt(
84        mut self,
85        receipt: &ExecutionResultReceiptV1,
86    ) -> Result<Self, ResearchContractError> {
87        self.validate()?;
88        receipt
89            .validate()
90            .map_err(|_| ResearchContractError::InvalidField("resultReceipt"))?;
91        if receipt.identity.digest != self.step_identity_digest {
92            return Err(ResearchContractError::InvalidField("stepIdentityDigest"));
93        }
94        if !self
95            .input_digests
96            .iter()
97            .any(|digest| digest == &receipt.evidence_digest)
98        {
99            return Err(ResearchContractError::InvalidField("evidenceDigest"));
100        }
101        let receipt_digest = digest(RESEARCH_WORKFLOW_RECEIPT_DIGEST_DOMAIN, receipt)
102            .map_err(|error| ResearchContractError::Serialization(error.to_string()))?;
103        if let Some(existing) = &self.result_receipt_digest {
104            if existing != &receipt_digest {
105                return Err(ResearchContractError::InvalidField("resultReceiptDigest"));
106            }
107            return Ok(self);
108        }
109        self.result_receipt_digest = Some(receipt_digest);
110        self.step_digest = self.expected_digest()?;
111        self.validate()?;
112        Ok(self)
113    }
114
115    pub fn validate_for_run(
116        &self,
117        run: &crate::research::ResearchRunV1,
118    ) -> Result<(), ResearchContractError> {
119        self.validate()?;
120        if self.project_id != run.project_id {
121            return Err(ResearchContractError::InvalidField("projectId"));
122        }
123        if self.run_id != run.run_id {
124            return Err(ResearchContractError::InvalidField("runId"));
125        }
126        Ok(())
127    }
128
129    pub fn validate(&self) -> Result<(), ResearchContractError> {
130        self.validate_without_digest()?;
131        validate_digest_field("stepDigest", &self.step_digest)?;
132        if self.step_digest != self.expected_digest()? {
133            return Err(ResearchContractError::DigestMismatch("stepDigest"));
134        }
135        Ok(())
136    }
137
138    pub fn from_slice(bytes: &[u8]) -> Result<Self, ResearchContractError> {
139        let step: Self = super::decode_json_slice(bytes)?;
140        step.validate()?;
141        Ok(step)
142    }
143
144    pub fn to_vec(&self) -> Result<Vec<u8>, ResearchContractError> {
145        self.validate()?;
146        super::encode_json(self)
147    }
148
149    fn validate_without_digest(&self) -> Result<(), ResearchContractError> {
150        if self.schema != RESEARCH_WORKFLOW_STEP_SCHEMA_V1 {
151            return Err(ResearchContractError::UnsupportedSchema);
152        }
153        validate_id("stepId", &self.step_id)?;
154        validate_id("projectId", &self.project_id)?;
155        validate_id("runId", &self.run_id)?;
156        validate_digest_field("workflowDigest", &self.workflow_digest)?;
157        validate_digest_field("stepIdentityDigest", &self.step_identity_digest)?;
158        if let Some(receipt_digest) = &self.result_receipt_digest {
159            validate_digest_field("resultReceiptDigest", receipt_digest)?;
160        }
161        if self.input_digests.is_empty() || self.input_digests.len() > RESEARCH_MAX_DIGESTS {
162            return Err(ResearchContractError::InvalidField("inputDigests"));
163        }
164        if self.output_artifact_digests.len() > RESEARCH_MAX_DIGESTS {
165            return Err(ResearchContractError::InvalidField("outputArtifactDigests"));
166        }
167        if self.depends_on.len() > RESEARCH_MAX_DIGESTS {
168            return Err(ResearchContractError::InvalidField("dependsOn"));
169        }
170        validate_sorted_unique_digests("inputDigests", &self.input_digests)?;
171        validate_sorted_unique_digests("outputArtifactDigests", &self.output_artifact_digests)?;
172        for pair in self.depends_on.windows(2) {
173            if pair[0] >= pair[1] {
174                return Err(ResearchContractError::InvalidField("dependsOn"));
175            }
176        }
177        for dependency in &self.depends_on {
178            validate_id("dependsOn", dependency)?;
179            if dependency == &self.step_id {
180                return Err(ResearchContractError::InvalidField("dependsOn"));
181            }
182        }
183        Ok(())
184    }
185
186    fn expected_digest(&self) -> Result<String, ResearchContractError> {
187        #[derive(Serialize)]
188        struct Identity<'a> {
189            schema: &'a str,
190            step_id: &'a str,
191            project_id: &'a str,
192            run_id: &'a str,
193            workflow_digest: &'a str,
194            step_identity_digest: &'a str,
195            result_receipt_digest: Option<&'a str>,
196            input_digests: &'a [String],
197            output_artifact_digests: &'a [String],
198            depends_on: &'a [String],
199        }
200        digest(
201            RESEARCH_WORKFLOW_STEP_DIGEST_DOMAIN,
202            &Identity {
203                schema: &self.schema,
204                step_id: &self.step_id,
205                project_id: &self.project_id,
206                run_id: &self.run_id,
207                workflow_digest: &self.workflow_digest,
208                step_identity_digest: &self.step_identity_digest,
209                result_receipt_digest: self.result_receipt_digest.as_deref(),
210                input_digests: &self.input_digests,
211                output_artifact_digests: &self.output_artifact_digests,
212                depends_on: &self.depends_on,
213            },
214        )
215    }
216}
217
218/// Bounded workflow projection for one research Run.
219#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
220#[serde(rename_all = "camelCase", deny_unknown_fields)]
221pub struct ResearchWorkflowPlanV1 {
222    pub schema: String,
223    pub plan_id: String,
224    pub project_id: String,
225    pub run_id: String,
226    pub workflow_digest: String,
227    #[serde(default, skip_serializing_if = "Option::is_none")]
228    pub random_seed: Option<u64>,
229    pub steps: Vec<ResearchWorkflowStepV1>,
230    pub plan_digest: String,
231}
232
233impl ResearchWorkflowPlanV1 {
234    pub fn new_for_run(
235        plan_id: impl Into<String>,
236        run: &crate::research::ResearchRunV1,
237        workflow_digest: impl Into<String>,
238        steps: Vec<ResearchWorkflowStepV1>,
239    ) -> Result<Self, ResearchContractError> {
240        let plan = Self::new(
241            plan_id,
242            run.project_id.clone(),
243            run.run_id.clone(),
244            workflow_digest,
245            run.random_seed,
246            steps,
247        )?;
248        plan.validate_for_run(run)?;
249        Ok(plan)
250    }
251
252    pub fn new(
253        plan_id: impl Into<String>,
254        project_id: impl Into<String>,
255        run_id: impl Into<String>,
256        workflow_digest: impl Into<String>,
257        random_seed: Option<u64>,
258        mut steps: Vec<ResearchWorkflowStepV1>,
259    ) -> Result<Self, ResearchContractError> {
260        steps.sort_unstable_by(|left, right| left.step_id.cmp(&right.step_id));
261        let mut plan = Self {
262            schema: RESEARCH_WORKFLOW_PLAN_SCHEMA_V1.to_owned(),
263            plan_id: plan_id.into(),
264            project_id: project_id.into(),
265            run_id: run_id.into(),
266            workflow_digest: workflow_digest.into(),
267            random_seed,
268            steps,
269            plan_digest: String::new(),
270        };
271        plan.validate_without_digest()?;
272        plan.plan_digest = plan.expected_digest()?;
273        Ok(plan)
274    }
275
276    pub fn validate_for_run(
277        &self,
278        run: &crate::research::ResearchRunV1,
279    ) -> Result<(), ResearchContractError> {
280        self.validate()?;
281        if self.project_id != run.project_id {
282            return Err(ResearchContractError::InvalidField("projectId"));
283        }
284        if self.run_id != run.run_id {
285            return Err(ResearchContractError::InvalidField("runId"));
286        }
287        if self.random_seed != run.random_seed {
288            return Err(ResearchContractError::InvalidField("randomSeed"));
289        }
290        for step in &self.steps {
291            step.validate_for_run(run)?;
292            if step.workflow_digest != self.workflow_digest {
293                return Err(ResearchContractError::InvalidField("workflowDigest"));
294            }
295        }
296        Ok(())
297    }
298
299    pub fn validate(&self) -> Result<(), ResearchContractError> {
300        self.validate_without_digest()?;
301        validate_digest_field("planDigest", &self.plan_digest)?;
302        if self.plan_digest != self.expected_digest()? {
303            return Err(ResearchContractError::DigestMismatch("planDigest"));
304        }
305        Ok(())
306    }
307
308    pub fn from_slice(bytes: &[u8]) -> Result<Self, ResearchContractError> {
309        let plan: Self = super::decode_json_slice(bytes)?;
310        plan.validate()?;
311        Ok(plan)
312    }
313
314    pub fn to_vec(&self) -> Result<Vec<u8>, ResearchContractError> {
315        self.validate()?;
316        super::encode_json(self)
317    }
318
319    /// Return affected step ids in dependency-first order for a re-run.
320    pub fn affected_step_ids_for_findings(
321        &self,
322        findings: &[crate::research::ResearchReviewFindingV1],
323    ) -> Result<Vec<String>, ResearchContractError> {
324        self.validate()?;
325        let mut directly_affected = BTreeSet::new();
326        for finding in findings {
327            finding.validate()?;
328            if finding.project_id != self.project_id {
329                return Err(ResearchContractError::InvalidField("finding.projectId"));
330            }
331            if finding.run_id != self.run_id {
332                return Err(ResearchContractError::InvalidField("finding.runId"));
333            }
334            for step in &self.steps {
335                if step
336                    .output_artifact_digests
337                    .iter()
338                    .any(|digest| digest == &finding.artifact_digest)
339                    || finding
340                        .evidence_digests
341                        .iter()
342                        .any(|digest| step.input_digests.contains(digest))
343                    || finding.evidence_digests.iter().any(|digest| {
344                        step.output_artifact_digests
345                            .iter()
346                            .any(|output| output == digest)
347                    })
348                {
349                    directly_affected.insert(step.step_id.clone());
350                }
351            }
352        }
353        close_dependents(&self.steps, directly_affected)
354    }
355
356    fn validate_without_digest(&self) -> Result<(), ResearchContractError> {
357        if self.schema != RESEARCH_WORKFLOW_PLAN_SCHEMA_V1 {
358            return Err(ResearchContractError::UnsupportedSchema);
359        }
360        validate_id("planId", &self.plan_id)?;
361        validate_id("projectId", &self.project_id)?;
362        validate_id("runId", &self.run_id)?;
363        validate_digest_field("workflowDigest", &self.workflow_digest)?;
364        if self.steps.len() > RESEARCH_MAX_WORKFLOW_STEPS {
365            return Err(ResearchContractError::InvalidField("steps"));
366        }
367        for pair in self.steps.windows(2) {
368            if pair[0].step_id >= pair[1].step_id {
369                return Err(ResearchContractError::InvalidField("steps"));
370            }
371        }
372        let mut step_ids = BTreeMap::new();
373        for step in &self.steps {
374            step.validate()?;
375            if step.project_id != self.project_id {
376                return Err(ResearchContractError::InvalidField("step.projectId"));
377            }
378            if step.run_id != self.run_id {
379                return Err(ResearchContractError::InvalidField("step.runId"));
380            }
381            if step.workflow_digest != self.workflow_digest {
382                return Err(ResearchContractError::InvalidField("workflowDigest"));
383            }
384            if step_ids.insert(step.step_id.as_str(), ()).is_some() {
385                return Err(ResearchContractError::InvalidField("stepId"));
386            }
387        }
388        for step in &self.steps {
389            for dependency in &step.depends_on {
390                if !step_ids.contains_key(dependency.as_str()) {
391                    return Err(ResearchContractError::InvalidField("dependsOn"));
392                }
393            }
394        }
395        detect_dependency_cycles(&self.steps)?;
396        Ok(())
397    }
398
399    fn expected_digest(&self) -> Result<String, ResearchContractError> {
400        #[derive(Serialize)]
401        struct Identity<'a> {
402            schema: &'a str,
403            plan_id: &'a str,
404            project_id: &'a str,
405            run_id: &'a str,
406            workflow_digest: &'a str,
407            random_seed: Option<u64>,
408            step_digests: Vec<&'a str>,
409        }
410        digest(
411            RESEARCH_WORKFLOW_PLAN_DIGEST_DOMAIN,
412            &Identity {
413                schema: &self.schema,
414                plan_id: &self.plan_id,
415                project_id: &self.project_id,
416                run_id: &self.run_id,
417                workflow_digest: &self.workflow_digest,
418                random_seed: self.random_seed,
419                step_digests: self
420                    .steps
421                    .iter()
422                    .map(|step| step.step_digest.as_str())
423                    .collect(),
424            },
425        )
426    }
427}
428
429/// Immutable projection of finding-triggered affected-step re-run lineage.
430#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
431#[serde(rename_all = "camelCase", deny_unknown_fields)]
432pub struct ResearchRerunLineageV1 {
433    pub schema: String,
434    pub lineage_id: String,
435    pub project_id: String,
436    pub run_id: String,
437    pub plan_digest: String,
438    pub finding_ids: Vec<String>,
439    pub affected_step_ids: Vec<String>,
440    pub lineage_digest: String,
441}
442
443impl ResearchRerunLineageV1 {
444    pub fn from_findings_for_plan(
445        lineage_id: impl Into<String>,
446        plan: &ResearchWorkflowPlanV1,
447        findings: &[crate::research::ResearchReviewFindingV1],
448    ) -> Result<Self, ResearchContractError> {
449        plan.validate()?;
450        let mut finding_ids: Vec<String> = findings
451            .iter()
452            .map(|finding| finding.finding_id.clone())
453            .collect();
454        finding_ids.sort();
455        finding_ids.dedup();
456        if finding_ids.len() != findings.len() {
457            return Err(ResearchContractError::InvalidField("findingId"));
458        }
459        let affected_step_ids = plan.affected_step_ids_for_findings(findings)?;
460        let mut lineage = Self {
461            schema: RESEARCH_RERUN_LINEAGE_SCHEMA_V1.to_owned(),
462            lineage_id: lineage_id.into(),
463            project_id: plan.project_id.clone(),
464            run_id: plan.run_id.clone(),
465            plan_digest: plan.plan_digest.clone(),
466            finding_ids,
467            affected_step_ids,
468            lineage_digest: String::new(),
469        };
470        lineage.validate_without_digest()?;
471        lineage.lineage_digest = lineage.expected_digest()?;
472        Ok(lineage)
473    }
474
475    pub fn validate(&self) -> Result<(), ResearchContractError> {
476        self.validate_without_digest()?;
477        validate_digest_field("lineageDigest", &self.lineage_digest)?;
478        if self.lineage_digest != self.expected_digest()? {
479            return Err(ResearchContractError::DigestMismatch("lineageDigest"));
480        }
481        Ok(())
482    }
483
484    pub fn from_slice(bytes: &[u8]) -> Result<Self, ResearchContractError> {
485        let lineage: Self = super::decode_json_slice(bytes)?;
486        lineage.validate()?;
487        Ok(lineage)
488    }
489
490    pub fn to_vec(&self) -> Result<Vec<u8>, ResearchContractError> {
491        self.validate()?;
492        super::encode_json(self)
493    }
494
495    fn validate_without_digest(&self) -> Result<(), ResearchContractError> {
496        if self.schema != RESEARCH_RERUN_LINEAGE_SCHEMA_V1 {
497            return Err(ResearchContractError::UnsupportedSchema);
498        }
499        validate_id("lineageId", &self.lineage_id)?;
500        validate_id("projectId", &self.project_id)?;
501        validate_id("runId", &self.run_id)?;
502        validate_digest_field("planDigest", &self.plan_digest)?;
503        if self.finding_ids.is_empty() || self.finding_ids.len() > RESEARCH_MAX_DIGESTS {
504            return Err(ResearchContractError::InvalidField("findingIds"));
505        }
506        for pair in self.finding_ids.windows(2) {
507            if pair[0] >= pair[1] {
508                return Err(ResearchContractError::InvalidField("findingIds"));
509            }
510        }
511        for finding_id in &self.finding_ids {
512            validate_id("findingId", finding_id)?;
513        }
514        if self.affected_step_ids.len() > RESEARCH_MAX_WORKFLOW_STEPS {
515            return Err(ResearchContractError::InvalidField("affectedStepIds"));
516        }
517        for step_id in &self.affected_step_ids {
518            validate_id("affectedStepIds", step_id)?;
519        }
520        Ok(())
521    }
522
523    fn expected_digest(&self) -> Result<String, ResearchContractError> {
524        #[derive(Serialize)]
525        struct Identity<'a> {
526            schema: &'a str,
527            lineage_id: &'a str,
528            project_id: &'a str,
529            run_id: &'a str,
530            plan_digest: &'a str,
531            finding_ids: &'a [String],
532            affected_step_ids: &'a [String],
533        }
534        digest(
535            RESEARCH_RERUN_LINEAGE_DIGEST_DOMAIN,
536            &Identity {
537                schema: &self.schema,
538                lineage_id: &self.lineage_id,
539                project_id: &self.project_id,
540                run_id: &self.run_id,
541                plan_digest: &self.plan_digest,
542                finding_ids: &self.finding_ids,
543                affected_step_ids: &self.affected_step_ids,
544            },
545        )
546    }
547}
548
549fn validate_sorted_unique_digests(
550    field: &'static str,
551    digests: &[String],
552) -> Result<(), ResearchContractError> {
553    for pair in digests.windows(2) {
554        if pair[0] >= pair[1] {
555            return Err(ResearchContractError::InvalidField(field));
556        }
557    }
558    for value in digests {
559        validate_digest_field(field, value)?;
560    }
561    Ok(())
562}
563
564fn detect_dependency_cycles(steps: &[ResearchWorkflowStepV1]) -> Result<(), ResearchContractError> {
565    let by_id: BTreeMap<&str, &ResearchWorkflowStepV1> = steps
566        .iter()
567        .map(|step| (step.step_id.as_str(), step))
568        .collect();
569    let mut visiting = BTreeSet::new();
570    let mut visited = BTreeSet::new();
571    for step in steps {
572        visit_cycle(step.step_id.as_str(), &by_id, &mut visiting, &mut visited)?;
573    }
574    Ok(())
575}
576
577fn visit_cycle(
578    step_id: &str,
579    by_id: &BTreeMap<&str, &ResearchWorkflowStepV1>,
580    visiting: &mut BTreeSet<String>,
581    visited: &mut BTreeSet<String>,
582) -> Result<(), ResearchContractError> {
583    if visited.contains(step_id) {
584        return Ok(());
585    }
586    if !visiting.insert(step_id.to_owned()) {
587        return Err(ResearchContractError::InvalidField("dependsOn"));
588    }
589    let Some(step) = by_id.get(step_id) else {
590        return Err(ResearchContractError::InvalidField("dependsOn"));
591    };
592    for dependency in &step.depends_on {
593        visit_cycle(dependency.as_str(), by_id, visiting, visited)?;
594    }
595    visiting.remove(step_id);
596    visited.insert(step_id.to_owned());
597    Ok(())
598}
599
600fn close_dependents(
601    steps: &[ResearchWorkflowStepV1],
602    mut affected: BTreeSet<String>,
603) -> Result<Vec<String>, ResearchContractError> {
604    let mut changed = true;
605    while changed {
606        changed = false;
607        for step in steps {
608            if affected.contains(&step.step_id) {
609                continue;
610            }
611            if step
612                .depends_on
613                .iter()
614                .any(|dependency| affected.contains(dependency))
615            {
616                affected.insert(step.step_id.clone());
617                changed = true;
618            }
619        }
620    }
621    // Dependency-first order: emit a step only after every dependency that is
622    // also affected has already been emitted.
623    let mut remaining = affected;
624    let mut ordered = Vec::new();
625    while !remaining.is_empty() {
626        let ready: Vec<String> = remaining
627            .iter()
628            .filter(|step_id| {
629                steps
630                    .iter()
631                    .find(|step| &step.step_id == *step_id)
632                    .map(|step| {
633                        step.depends_on
634                            .iter()
635                            .all(|dependency| !remaining.contains(dependency))
636                    })
637                    .unwrap_or(false)
638            })
639            .cloned()
640            .collect();
641        if ready.is_empty() {
642            return Err(ResearchContractError::InvalidField("dependsOn"));
643        }
644        for step_id in ready {
645            remaining.remove(&step_id);
646            ordered.push(step_id);
647        }
648    }
649    Ok(ordered)
650}
651
652#[cfg(test)]
653mod tests {
654    use super::*;
655    use crate::capability::{
656        CapabilityCeiling, CapabilityContribution, CapabilityDescriptor,
657        CapabilityExecutionCeiling, CapabilityKind, CapabilitySet, CapabilitySource,
658        CodeCatalogGeneration, GovernanceCapabilityCeiling, RunCapabilityBindingV1, Sha256Digest,
659        WorkspaceCapabilityCeiling,
660    };
661    use crate::execution_identity::{ExecutionIdentityV1, ExecutionResultOutcomeV1};
662    use crate::research::{
663        ResearchReproducibilityV1, ResearchReviewCategoryV1, ResearchReviewFindingV1,
664        ResearchReviewSeverityV1, ResearchRunStatusV1, ResearchRunV1,
665    };
666
667    fn digest(ch: char) -> String {
668        format!("sha256:{}", ch.to_string().repeat(64))
669    }
670
671    fn binding() -> RunCapabilityBindingV1 {
672        let source =
673            CapabilitySource::builtin("test", Sha256Digest::new(digest('c')).unwrap()).unwrap();
674        let descriptor = CapabilityDescriptor::new(
675            &source,
676            CapabilityKind::Tool,
677            "tool",
678            "tool",
679            Sha256Digest::new(digest('d')).unwrap(),
680            [],
681        )
682        .unwrap();
683        let contribution = CapabilityContribution::new(source, [descriptor]).unwrap();
684        let set = CapabilitySet::from_contributions(CodeCatalogGeneration::new(1), [contribution])
685            .unwrap();
686        let ceiling = CapabilityCeiling::all(
687            &set,
688            WorkspaceCapabilityCeiling::default(),
689            GovernanceCapabilityCeiling::default(),
690            CapabilityExecutionCeiling::new(1, 1, None, None, None).unwrap(),
691        )
692        .unwrap();
693        RunCapabilityBindingV1::from_set_and_ceiling(&set, &ceiling).unwrap()
694    }
695
696    fn admitted_run() -> ResearchRunV1 {
697        let mut run = ResearchRunV1::new(
698            "run-1",
699            "project-1",
700            1,
701            digest('1'),
702            digest('2'),
703            binding(),
704            "provider-1",
705            "model-1",
706            ResearchReproducibilityV1::Deterministic,
707            Some(7),
708        )
709        .unwrap();
710        run.transition_to(ResearchRunStatusV1::Admitted).unwrap();
711        run
712    }
713
714    fn receipt(identity_digest: &str, evidence: &str, result: &str) -> ExecutionResultReceiptV1 {
715        let identity = ExecutionIdentityV1 {
716            schema: crate::execution_identity::EXECUTION_IDENTITY_SCHEMA_V1.to_owned(),
717            domain: crate::execution_identity::WORKFLOW_STEP_IDENTITY_DOMAIN_V1.to_owned(),
718            digest: identity_digest.to_owned(),
719        };
720        ExecutionResultReceiptV1::new(
721            identity,
722            evidence.to_owned(),
723            ExecutionResultOutcomeV1::Succeeded,
724            Some(result.to_owned()),
725            16,
726        )
727        .unwrap()
728    }
729
730    #[test]
731    fn plan_binds_receipts_and_selects_transitive_rerun_steps() {
732        let run = admitted_run();
733        let analyze = ResearchWorkflowStepV1::new(
734            "analyze",
735            "project-1",
736            "run-1",
737            digest('7'),
738            digest('a'),
739            vec![digest('2')],
740            vec![digest('3')],
741            Vec::new(),
742        )
743        .unwrap()
744        .bind_result_receipt(&receipt(&digest('a'), &digest('2'), &digest('9')))
745        .unwrap();
746        let report = ResearchWorkflowStepV1::new(
747            "report",
748            "project-1",
749            "run-1",
750            digest('7'),
751            digest('b'),
752            vec![digest('3')],
753            vec![digest('4')],
754            vec!["analyze".to_owned()],
755        )
756        .unwrap();
757        let plan = ResearchWorkflowPlanV1::new_for_run(
758            "plan-1",
759            &run,
760            digest('7'),
761            vec![report.clone(), analyze],
762        )
763        .unwrap();
764        assert_eq!(plan.random_seed, Some(7));
765        let finding = ResearchReviewFindingV1::new(
766            "finding-1",
767            "project-1",
768            "run-1",
769            digest('3'),
770            ResearchReviewCategoryV1::Numeric,
771            ResearchReviewSeverityV1::Error,
772            "numeric mismatch",
773            None,
774            vec![digest('2')],
775            "reviewer",
776            1,
777        )
778        .unwrap();
779        assert_eq!(
780            plan.affected_step_ids_for_findings(&[finding.clone()])
781                .unwrap(),
782            vec!["analyze".to_owned(), "report".to_owned()]
783        );
784        let encoded = plan.to_vec().unwrap();
785        assert_eq!(ResearchWorkflowPlanV1::from_slice(&encoded).unwrap(), plan);
786        let lineage =
787            ResearchRerunLineageV1::from_findings_for_plan("lineage-1", &plan, &[finding]).unwrap();
788        assert_eq!(
789            lineage.affected_step_ids,
790            vec!["analyze".to_owned(), "report".to_owned()]
791        );
792        assert_eq!(lineage.plan_digest, plan.plan_digest);
793        let lineage_bytes = lineage.to_vec().unwrap();
794        assert_eq!(
795            ResearchRerunLineageV1::from_slice(&lineage_bytes).unwrap(),
796            lineage
797        );
798    }
799
800    #[test]
801    fn receipt_and_seed_drift_fail_closed() {
802        let run = admitted_run();
803        let step = ResearchWorkflowStepV1::new(
804            "analyze",
805            "project-1",
806            "run-1",
807            digest('7'),
808            digest('a'),
809            vec![digest('2')],
810            vec![digest('3')],
811            Vec::new(),
812        )
813        .unwrap();
814        assert!(matches!(
815            step.clone()
816                .bind_result_receipt(&receipt(&digest('b'), &digest('2'), &digest('9'))),
817            Err(ResearchContractError::InvalidField("stepIdentityDigest"))
818        ));
819        assert!(matches!(
820            step.bind_result_receipt(&receipt(&digest('a'), &digest('5'), &digest('9'))),
821            Err(ResearchContractError::InvalidField("evidenceDigest"))
822        ));
823        let plan = ResearchWorkflowPlanV1::new(
824            "plan-1",
825            "project-1",
826            "run-1",
827            digest('7'),
828            Some(99),
829            Vec::new(),
830        )
831        .unwrap();
832        assert!(matches!(
833            plan.validate_for_run(&run),
834            Err(ResearchContractError::InvalidField("randomSeed"))
835        ));
836    }
837}