Skip to main content

kranz_engine/gate_evaluation/
input_builder.rs

1//! Host-selected stage evidence. This API accepts observed receipts and approved
2//! objects, never paths to worker manifests, transcripts or provider state.
3use super::evidence::FrozenEvidence;
4use super::lifecycle::{Outcome, Policy, Record};
5use super::protocol::*;
6use super::snapshot::SourceSnapshot;
7use crate::pack::evaluator::{Kind, PinnedRegistration};
8use crate::types::Plan;
9use chrono::{DateTime, SecondsFormat, Utc};
10use serde::Serialize;
11use std::collections::{BTreeMap, BTreeSet};
12
13/// Engine-observed execution metadata. Do not construct these from a worker's
14/// report or infer an assertion count from the subprocess exit code.
15#[derive(Debug, Clone, Serialize)]
16#[serde(rename_all = "camelCase")]
17pub struct ObservedCheck {
18    pub check_id: Id,
19    pub run_id: Id,
20    pub sequence: u64,
21    pub checked_content: Digest,
22    pub environment: Digest,
23    pub command: String,
24    pub exit_code: Option<i32>,
25    pub assertions_executed: Option<u64>,
26    /// Bounded, scrubbed process output; never an assertion-count claim.
27    #[serde(skip_serializing_if = "Option::is_none")]
28    pub output_summary: Option<String>,
29}
30
31/// Selected by stage policy, separately from whichever receipts happen to exist.
32pub struct RequiredCheck {
33    pub id: Id,
34    pub command: String,
35    pub require_assertions: bool,
36}
37
38#[derive(Clone)]
39pub struct Checks<'a> {
40    pub environment: Digest,
41    pub required: &'a [RequiredCheck],
42    pub observed: &'a [ObservedCheck],
43}
44
45/// Identifies an independently accepted feature without importing its worker
46/// report. The stage driver resolves these from existing engine events.
47#[derive(Debug, Clone, Serialize)]
48#[serde(rename_all = "camelCase")]
49pub struct FeatureReceipt {
50    pub feature_id: Id,
51    pub worker_run_id: Id,
52    pub candidate_commit: GitObject,
53    pub independent_validation_attempt: Id,
54}
55
56#[derive(Clone)]
57pub enum StageInput<'a> {
58    Plan {
59        revision: u64,
60        base_commit: GitObject,
61    },
62    Invocation(&'a crate::live_permission::Request),
63    Milestone {
64        milestone_id: Id,
65        plan_index: usize,
66        snapshot: &'a SourceSnapshot,
67    },
68    Deliverable {
69        snapshot: &'a SourceSnapshot,
70        features: &'a [FeatureReceipt],
71    },
72    Integration {
73        snapshot: &'a SourceSnapshot,
74        live_base_commit: GitObject,
75        candidate_commit: GitObject,
76        integration_tree: GitObject,
77    },
78}
79
80impl StageInput<'_> {
81    pub(crate) fn stage(&self) -> Stage {
82        match self {
83            Self::Plan { .. } => Stage::PlanApproval,
84            Self::Invocation(_) => Stage::CommandPermission,
85            Self::Milestone { .. } => Stage::MilestoneValidation,
86            Self::Deliverable { .. } => Stage::FinalGate,
87            Self::Integration { .. } => Stage::Merge,
88        }
89    }
90}
91
92pub struct BuildInput<'a> {
93    pub mission_id: Id,
94    pub evaluation_id: Id,
95    pub attempt_id: Id,
96    pub workspace_id: Id,
97    /// The exact engine-approved serialization, not a freshly rendered plan.
98    pub plan_bytes: &'a [u8],
99    pub policy: Policy,
100    pub registration: &'a PinnedRegistration,
101    pub stage: StageInput<'a>,
102    pub checks: Checks<'a>,
103    /// Existing stage diagnostics retain their own verdict and posture. They
104    /// are not command exits and never satisfy required command receipts.
105    pub diagnostics: &'a [crate::gate::GateReport],
106    /// Previously validated independent results selected from the engine log.
107    pub prior_findings: &'a [&'a Record],
108    pub source_log_range: Option<LogRange>,
109    pub deadline: DateTime<Utc>,
110    pub limits: Limits,
111}
112
113pub struct BuiltInput {
114    pub evidence: FrozenEvidence,
115    pub policy: Policy,
116    pub permission_request_id: Option<String>,
117}
118
119/// Freeze all hashes from the supplied bytes in one operation. Retrying a
120/// checker reuses the returned evidence rather than rebuilding this input.
121pub fn build(input: BuildInput<'_>) -> Result<BuiltInput, String> {
122    if input.policy.kind != input.registration.declaration().kind
123        || input.policy.enforcement != input.registration.declaration().enforcement
124    {
125        return Err("gate policy differs from the approved registration".into());
126    }
127    let plan: Plan = super::protocol::decode(input.plan_bytes)?;
128    let mut parts = Parts::default();
129    parts.add("plan", ArtifactRole::Plan, input.plan_bytes.to_vec(), None)?;
130    parts.json(
131        "scope",
132        ArtifactRole::Scope,
133        &serde_json::json!({"goal":plan.goal,"touchSet":plan.touch_set}),
134    )?;
135    parts.add(
136        "registration",
137        ArtifactRole::Registration,
138        input.registration.bytes().to_vec(),
139        None,
140    )?;
141    let plan_digest = Digest::of(input.plan_bytes);
142    let mut permission_request_id = None;
143    let (subject, checked_content) = match input.stage {
144        StageInput::Plan {
145            revision,
146            base_commit,
147        } => {
148            parts.json(
149                "criteria",
150                ArtifactRole::Criteria,
151                &plan.validation_contract,
152            )?;
153            (
154                Subject::Plan {
155                    revision,
156                    plan_digest: plan_digest.clone(),
157                    base_commit,
158                },
159                plan_digest.clone(),
160            )
161        }
162        StageInput::Invocation(permission) => {
163            permission.validate().map_err(|e| e.to_string())?;
164            let workspace = super::lifecycle::workspace_id(&permission.binding.workspace);
165            if permission.binding.mission_id != input.mission_id.as_str()
166                || workspace != input.workspace_id
167                || plan_digest.as_str() != format!("sha256:{}", permission.binding.plan_digest)
168                || input.policy.mission_policy_digest.as_str()
169                    != format!("sha256:{}", permission.binding.policy_digest)
170                || input.deadline > permission.proposal.deadline
171                || permission.proposal.prohibition.is_some()
172            {
173                return Err("invocation evidence does not match live permission authority".into());
174            }
175            let action_digest =
176                parts.json("action", ArtifactRole::Action, &permission.proposal.action)?;
177            let options_digest = parts.json(
178                "permission-options",
179                ArtifactRole::PermissionOptions,
180                &permission.proposal.options,
181            )?;
182            permission_request_id = Some(permission.proposal.id.clone());
183            (
184                Subject::Invocation {
185                    run_id: Id::try_from(permission.binding.run_id.clone())?,
186                    peer_session_id: permission.proposal.peer_session_id.clone(),
187                    tool_call_id: permission.proposal.tool_call_id.clone(),
188                    peer_request_id: serde_json::from_value(
189                        permission.proposal.peer_request_id.clone(),
190                    )
191                    .map_err(|e| e.to_string())?,
192                    action_digest: action_digest.clone(),
193                    options_digest,
194                    cwd_id: workspace,
195                },
196                action_digest,
197            )
198        }
199        StageInput::Milestone {
200            milestone_id,
201            plan_index,
202            snapshot,
203        } => {
204            let milestone = plan
205                .milestones
206                .get(plan_index)
207                .ok_or("milestone is absent from approved plan")?;
208            let criteria = milestone
209                .features
210                .iter()
211                .map(|f| &f.validation_criteria)
212                .collect::<Vec<_>>();
213            let criteria_digest = parts.json("criteria", ArtifactRole::Criteria, &criteria)?;
214            let content = parts.snapshot(snapshot)?;
215            (
216                Subject::Milestone {
217                    milestone_id,
218                    start_commit: git_object(&snapshot.identity.base)?,
219                    snapshot_digest: snapshot.identity.inventory_digest.clone(),
220                    criteria_digest,
221                },
222                content,
223            )
224        }
225        StageInput::Deliverable { snapshot, features } => {
226            let mut ids = BTreeSet::new();
227            for feature in features {
228                if !ids.insert(&feature.feature_id) {
229                    return Err("duplicate feature receipt".into());
230                }
231                if git_object(&feature.candidate_commit.value)? != feature.candidate_commit {
232                    return Err("feature receipt Git algorithm does not match its object".into());
233                }
234            }
235            if features.is_empty() {
236                return Err("deliverable evidence requires accepted feature receipts".into());
237            }
238            let feature_receipt_digest =
239                parts.json("feature-receipts", ArtifactRole::FeatureReceipts, &features)?;
240            parts.json(
241                "criteria",
242                ArtifactRole::Criteria,
243                &plan.validation_contract,
244            )?;
245            let content = parts.snapshot(snapshot)?;
246            (
247                Subject::Deliverable {
248                    base_commit: git_object(&snapshot.identity.base)?,
249                    snapshot_digest: snapshot.identity.inventory_digest.clone(),
250                    feature_receipt_digest,
251                },
252                content,
253            )
254        }
255        StageInput::Integration {
256            snapshot,
257            live_base_commit,
258            candidate_commit,
259            integration_tree,
260        } => {
261            if live_base_commit != git_object(&snapshot.identity.base)? {
262                return Err("integration snapshot does not use the live base".into());
263            }
264            parts.json(
265                "criteria",
266                ArtifactRole::Criteria,
267                &plan.validation_contract,
268            )?;
269            let content = parts.snapshot(snapshot)?;
270            (
271                Subject::Integration {
272                    live_base_commit,
273                    candidate_commit,
274                    integration_tree,
275                    snapshot_digest: snapshot.identity.inventory_digest.clone(),
276                },
277                content,
278            )
279        }
280    };
281    let mut policy = input.policy;
282    policy.mechanical_prerequisites_passed &= parts.checks(&input.checks, &checked_content)?;
283    if policy.kind == Kind::Judgment && !policy.mechanical_prerequisites_passed {
284        return Err("judgment requires current, nonvacuous mechanical evidence".into());
285    }
286    parts.prior_findings(input.prior_findings, &input.mission_id)?;
287    for (index, report) in input.diagnostics.iter().enumerate() {
288        parts.json(
289            &format!("diagnostic-{index}"),
290            ArtifactRole::CheckReceipt,
291            &serde_json::json!({"name":report.name,"kind":report.kind,
292                "verdict":report.outcome.verdict,"reference":report.outcome.artefact.reference,
293                "detail":report.outcome.artefact.detail,"ruleIds":report.outcome.rule_ids,
294                "authority":"existing-stage-policy; diagnostic is not an execution receipt"}),
295        )?;
296    }
297    let policy_digest = parts.json("policy", ArtifactRole::Policy, &policy)?;
298    let subject_digest = parts.json("subject", ArtifactRole::Subject, &subject)?;
299    let binding = Binding {
300        subject_digest,
301        plan_digest,
302        policy_digest,
303        registration_digest: input.registration.digest(),
304        workspace_id: input.workspace_id,
305    };
306    let manifest_bytes = serde_json::to_vec(&Manifest {
307        schema_version: 1,
308        mission_id: input.mission_id.clone(),
309        binding: binding.clone(),
310        artifacts: parts.artifacts.into_values().collect(),
311        source_log_range: input.source_log_range,
312    })
313    .map_err(|e| e.to_string())?;
314    let request = Request {
315        jsonrpc: "2.0".into(),
316        id: input.attempt_id.clone(),
317        method: "gate/evaluate".into(),
318        params: Parameters {
319            schema_version: 1,
320            evaluation_id: input.evaluation_id,
321            attempt_id: input.attempt_id,
322            gate_id: input.registration.declaration().name.clone(),
323            mission_id: input.mission_id,
324            binding,
325            stage: subject.stage(),
326            subject,
327            evidence: ArtifactRef {
328                path: WirePath::try_from("inputs/manifest.json".to_string())?,
329                digest: Digest::of(&manifest_bytes),
330                bytes: manifest_bytes.len() as u64,
331            },
332            deadline: input.deadline.to_rfc3339_opts(SecondsFormat::Secs, true),
333            limits: input.limits,
334        },
335    };
336    let evidence = FrozenEvidence::new(
337        serde_json::to_vec(&request).map_err(|e| e.to_string())?,
338        manifest_bytes,
339        parts.bytes,
340        input.registration,
341    )?;
342    Ok(BuiltInput {
343        evidence,
344        policy,
345        permission_request_id,
346    })
347}
348
349#[derive(Default)]
350struct Parts {
351    artifacts: BTreeMap<Id, InputArtifact>,
352    bytes: BTreeMap<Id, Vec<u8>>,
353    total: usize,
354}
355impl Parts {
356    fn add(
357        &mut self,
358        name: &str,
359        role: ArtifactRole,
360        bytes: Vec<u8>,
361        producer: Option<Producer>,
362    ) -> Result<Digest, String> {
363        if self.artifacts.len() >= 1000 || bytes.len() > 64 * 1024 * 1024 {
364            return Err("stage input exceeds evidence limits".into());
365        }
366        self.total = self
367            .total
368            .checked_add(bytes.len())
369            .ok_or("input byte overflow")?;
370        if self.total > 128 * 1024 * 1024 {
371            return Err("stage input exceeds aggregate byte limit".into());
372        }
373        let id = Id::try_from(name.to_string())?;
374        if self.artifacts.contains_key(&id) {
375            return Err("duplicate stage input ID".into());
376        }
377        let digest = Digest::of(&bytes);
378        self.artifacts.insert(
379            id.clone(),
380            InputArtifact {
381                id: id.clone(),
382                role,
383                content: ArtifactRef {
384                    path: WirePath::try_from(format!("inputs/{name}"))?,
385                    digest: digest.clone(),
386                    bytes: bytes.len() as u64,
387                },
388                producer: producer.unwrap_or(Producer {
389                    kind: ProducerKind::Engine,
390                    run_id: None,
391                }),
392            },
393        );
394        self.bytes.insert(id, bytes);
395        Ok(digest)
396    }
397    fn json(
398        &mut self,
399        name: &str,
400        role: ArtifactRole,
401        value: &impl Serialize,
402    ) -> Result<Digest, String> {
403        self.add(
404            name,
405            role,
406            serde_json::to_vec(value).map_err(|e| e.to_string())?,
407            None,
408        )
409    }
410    fn snapshot(&mut self, snapshot: &SourceSnapshot) -> Result<Digest, String> {
411        if Digest::of(&snapshot.inventory) != snapshot.identity.inventory_digest
412            || Digest::of(&snapshot.selection) != snapshot.identity.exclusions_digest
413        {
414            return Err("source snapshot identity changed".into());
415        }
416        self.add(
417            "snapshot-inventory",
418            ArtifactRole::SnapshotInventory,
419            snapshot.inventory.clone(),
420            None,
421        )?;
422        self.add(
423            "source-selection",
424            ArtifactRole::Context,
425            snapshot.selection.clone(),
426            None,
427        )?;
428        for (id, bytes) in &snapshot.files {
429            self.add(id.as_str(), ArtifactRole::Source, bytes.clone(), None)?;
430        }
431        self.json("source-identity", ArtifactRole::Context, &snapshot.identity)
432    }
433    fn checks(&mut self, checks: &Checks<'_>, content: &Digest) -> Result<bool, String> {
434        if checks.observed.len() > 256 || checks.required.len() > 256 {
435            return Err("too many stage checks".into());
436        }
437        let mut required_ids = BTreeSet::new();
438        let mut seen = BTreeSet::new();
439        for receipt in checks.observed {
440            if receipt
441                .output_summary
442                .as_ref()
443                .is_some_and(|s| s.len() > 32768)
444                || receipt.command.trim().is_empty()
445                || receipt.command.len() > 8192
446                || !seen.insert((&receipt.check_id, receipt.sequence))
447            {
448                return Err("invalid or duplicate observed check".into());
449            }
450            self.add(&label("check", &(&receipt.check_id, receipt.sequence))?, ArtifactRole::CheckReceipt,
451                serde_json::to_vec(&serde_json::json!({"receipt":receipt,"current":receipt.checked_content == *content && receipt.environment == checks.environment})).map_err(|e| e.to_string())?,
452                Some(Producer { kind: ProducerKind::Engine, run_id: Some(receipt.run_id.clone()) }))?;
453        }
454        let mut passed = true;
455        for required in checks.required {
456            if !required_ids.insert(&required.id)
457                || required.command.trim().is_empty()
458                || required.command.len() > 8192
459            {
460                return Err("invalid required check set".into());
461            }
462            let latest = checks
463                .observed
464                .iter()
465                .filter(|r| r.check_id == required.id)
466                .max_by_key(|r| r.sequence);
467            passed &= latest.is_some_and(|r| {
468                r.command == required.command
469                    && r.checked_content == *content
470                    && r.environment == checks.environment
471                    && r.exit_code == Some(0)
472                    && (!required.require_assertions
473                        || r.assertions_executed.is_some_and(|n| n > 0))
474            });
475        }
476        self.json("check-requirements", ArtifactRole::Context, &serde_json::json!({
477            "checkedContent":content,"environment":checks.environment,
478            "required":checks.required.iter().map(|r| serde_json::json!({"id":r.id,"command":r.command,"requireAssertions":r.require_assertions})).collect::<Vec<_>>(),
479            "passed":passed,
480        }))?;
481        Ok(passed)
482    }
483    fn prior_findings(&mut self, records: &[&Record], mission: &Id) -> Result<(), String> {
484        if records.len() > 128 {
485            return Err("too many prior finding records".into());
486        }
487        for record in records {
488            let params = &record.requested.request.params;
489            if params.mission_id != *mission || record.requested.policy.kind != Kind::Judgment {
490                return Err(
491                    "prior findings require an independent result from this mission".into(),
492                );
493            }
494            let finished = record
495                .finished
496                .as_ref()
497                .ok_or("prior findings have no finished evaluation")?;
498            let Outcome::Evaluated { result, .. } = &finished.outcome else {
499                return Err("failed process is not a prior finding".into());
500            };
501            let mut verified = Record::new(
502                record.requested.clone(),
503                record.requested_at,
504                record.requested_seq,
505            )?;
506            verified.finish(
507                finished.clone(),
508                record
509                    .finished_at
510                    .ok_or("prior result timestamp is missing")?,
511            )?;
512            self.add(&label("prior", &params.attempt_id)?, ArtifactRole::PriorFinding,
513                serde_json::to_vec(&serde_json::json!({"evaluationId":params.evaluation_id,"attemptId":params.attempt_id,
514                    "binding":result.binding,"evidenceDigest":result.evidence_digest,"findings":result.findings})).map_err(|e| e.to_string())?,
515                Some(Producer { kind: ProducerKind::IndependentChecker, run_id: None }))?;
516        }
517        Ok(())
518    }
519}
520
521pub(crate) fn git_object(value: &str) -> Result<GitObject, String> {
522    let algorithm = match value.len() {
523        40 => "sha1",
524        64 => "sha256",
525        _ => return Err("invalid source Git identity".into()),
526    };
527    if !value
528        .bytes()
529        .all(|c| c.is_ascii_digit() || (b'a'..=b'f').contains(&c))
530    {
531        return Err("invalid source Git identity".into());
532    }
533    Ok(GitObject {
534        algorithm: algorithm.into(),
535        value: value.into(),
536    })
537}
538
539fn label(prefix: &str, value: &impl Serialize) -> Result<String, String> {
540    let bytes = serde_json::to_vec(value).map_err(|e| e.to_string())?;
541    Ok(format!("{prefix}-{}", &Digest::of(&bytes).as_str()[7..]))
542}