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::{Deserialize, 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, Deserialize)]
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                "baselineCandidate": crate::contract_controls::pair::descriptor(report.outcome.artefact.detail.as_deref()).map(|d| d.summary),
295                "authority":"existing-stage-policy; diagnostic is not an execution receipt"}),
296        )?;
297    }
298    let policy_digest = parts.json("policy", ArtifactRole::Policy, &policy)?;
299    let subject_digest = parts.json("subject", ArtifactRole::Subject, &subject)?;
300    let binding = Binding {
301        subject_digest,
302        plan_digest,
303        policy_digest,
304        registration_digest: input.registration.digest(),
305        workspace_id: input.workspace_id,
306    };
307    let manifest_bytes = serde_json::to_vec(&Manifest {
308        schema_version: 1,
309        mission_id: input.mission_id.clone(),
310        binding: binding.clone(),
311        artifacts: parts.artifacts.into_values().collect(),
312        source_log_range: input.source_log_range,
313    })
314    .map_err(|e| e.to_string())?;
315    let request = Request {
316        jsonrpc: "2.0".into(),
317        id: input.attempt_id.clone(),
318        method: "gate/evaluate".into(),
319        params: Parameters {
320            schema_version: 1,
321            evaluation_id: input.evaluation_id,
322            attempt_id: input.attempt_id,
323            gate_id: input.registration.declaration().name.clone(),
324            mission_id: input.mission_id,
325            binding,
326            stage: subject.stage(),
327            subject,
328            evidence: ArtifactRef {
329                path: WirePath::try_from("inputs/manifest.json".to_string())?,
330                digest: Digest::of(&manifest_bytes),
331                bytes: manifest_bytes.len() as u64,
332            },
333            deadline: input.deadline.to_rfc3339_opts(SecondsFormat::Secs, true),
334            limits: input.limits,
335        },
336    };
337    let evidence = FrozenEvidence::new(
338        serde_json::to_vec(&request).map_err(|e| e.to_string())?,
339        manifest_bytes,
340        parts.bytes,
341        input.registration,
342    )?;
343    Ok(BuiltInput {
344        evidence,
345        policy,
346        permission_request_id,
347    })
348}
349
350#[derive(Default)]
351struct Parts {
352    artifacts: BTreeMap<Id, InputArtifact>,
353    bytes: BTreeMap<Id, Vec<u8>>,
354    total: usize,
355}
356impl Parts {
357    fn add(
358        &mut self,
359        name: &str,
360        role: ArtifactRole,
361        bytes: Vec<u8>,
362        producer: Option<Producer>,
363    ) -> Result<Digest, String> {
364        if self.artifacts.len() >= 1000 || bytes.len() > 64 * 1024 * 1024 {
365            return Err("stage input exceeds evidence limits".into());
366        }
367        self.total = self
368            .total
369            .checked_add(bytes.len())
370            .ok_or("input byte overflow")?;
371        if self.total > 128 * 1024 * 1024 {
372            return Err("stage input exceeds aggregate byte limit".into());
373        }
374        let id = Id::try_from(name.to_string())?;
375        if self.artifacts.contains_key(&id) {
376            return Err("duplicate stage input ID".into());
377        }
378        let digest = Digest::of(&bytes);
379        self.artifacts.insert(
380            id.clone(),
381            InputArtifact {
382                id: id.clone(),
383                role,
384                content: ArtifactRef {
385                    path: WirePath::try_from(format!("inputs/{name}"))?,
386                    digest: digest.clone(),
387                    bytes: bytes.len() as u64,
388                },
389                producer: producer.unwrap_or(Producer {
390                    kind: ProducerKind::Engine,
391                    run_id: None,
392                }),
393            },
394        );
395        self.bytes.insert(id, bytes);
396        Ok(digest)
397    }
398    fn json(
399        &mut self,
400        name: &str,
401        role: ArtifactRole,
402        value: &impl Serialize,
403    ) -> Result<Digest, String> {
404        self.add(
405            name,
406            role,
407            serde_json::to_vec(value).map_err(|e| e.to_string())?,
408            None,
409        )
410    }
411    fn snapshot(&mut self, snapshot: &SourceSnapshot) -> Result<Digest, String> {
412        if Digest::of(&snapshot.inventory) != snapshot.identity.inventory_digest
413            || Digest::of(&snapshot.selection) != snapshot.identity.exclusions_digest
414        {
415            return Err("source snapshot identity changed".into());
416        }
417        self.add(
418            "snapshot-inventory",
419            ArtifactRole::SnapshotInventory,
420            snapshot.inventory.clone(),
421            None,
422        )?;
423        self.add(
424            "source-selection",
425            ArtifactRole::Context,
426            snapshot.selection.clone(),
427            None,
428        )?;
429        for (id, bytes) in &snapshot.files {
430            self.add(id.as_str(), ArtifactRole::Source, bytes.clone(), None)?;
431        }
432        self.json("source-identity", ArtifactRole::Context, &snapshot.identity)
433    }
434    fn checks(&mut self, checks: &Checks<'_>, content: &Digest) -> Result<bool, String> {
435        if checks.observed.len() > 256 || checks.required.len() > 256 {
436            return Err("too many stage checks".into());
437        }
438        let mut required_ids = BTreeSet::new();
439        let mut seen = BTreeSet::new();
440        for receipt in checks.observed {
441            if receipt
442                .output_summary
443                .as_ref()
444                .is_some_and(|s| s.len() > 32768)
445                || receipt.command.trim().is_empty()
446                || receipt.command.len() > 8192
447                || !seen.insert((&receipt.check_id, receipt.sequence))
448            {
449                return Err("invalid or duplicate observed check".into());
450            }
451            self.add(&label("check", &(&receipt.check_id, receipt.sequence))?, ArtifactRole::CheckReceipt,
452                serde_json::to_vec(&serde_json::json!({"receipt":receipt,"current":receipt.checked_content == *content && receipt.environment == checks.environment})).map_err(|e| e.to_string())?,
453                Some(Producer { kind: ProducerKind::Engine, run_id: Some(receipt.run_id.clone()) }))?;
454        }
455        let mut passed = true;
456        for required in checks.required {
457            if !required_ids.insert(&required.id)
458                || required.command.trim().is_empty()
459                || required.command.len() > 8192
460            {
461                return Err("invalid required check set".into());
462            }
463            let latest = checks
464                .observed
465                .iter()
466                .filter(|r| r.check_id == required.id)
467                .max_by_key(|r| r.sequence);
468            passed &= latest.is_some_and(|r| {
469                r.command == required.command
470                    && r.checked_content == *content
471                    && r.environment == checks.environment
472                    && r.exit_code == Some(0)
473                    && (!required.require_assertions
474                        || r.assertions_executed.is_some_and(|n| n > 0))
475            });
476        }
477        self.json("check-requirements", ArtifactRole::Context, &serde_json::json!({
478            "checkedContent":content,"environment":checks.environment,
479            "required":checks.required.iter().map(|r| serde_json::json!({"id":r.id,"command":r.command,"requireAssertions":r.require_assertions})).collect::<Vec<_>>(),
480            "passed":passed,
481        }))?;
482        Ok(passed)
483    }
484    fn prior_findings(&mut self, records: &[&Record], mission: &Id) -> Result<(), String> {
485        if records.len() > 128 {
486            return Err("too many prior finding records".into());
487        }
488        for record in records {
489            let params = &record.requested.request.params;
490            if params.mission_id != *mission || record.requested.policy.kind != Kind::Judgment {
491                return Err(
492                    "prior findings require an independent result from this mission".into(),
493                );
494            }
495            let finished = record
496                .finished
497                .as_ref()
498                .ok_or("prior findings have no finished evaluation")?;
499            let Outcome::Evaluated { result, .. } = &finished.outcome else {
500                return Err("failed process is not a prior finding".into());
501            };
502            let mut verified = Record::new(
503                record.requested.clone(),
504                record.requested_at,
505                record.requested_seq,
506            )?;
507            verified.finish(
508                finished.clone(),
509                record
510                    .finished_at
511                    .ok_or("prior result timestamp is missing")?,
512            )?;
513            self.add(&label("prior", &params.attempt_id)?, ArtifactRole::PriorFinding,
514                serde_json::to_vec(&serde_json::json!({"evaluationId":params.evaluation_id,"attemptId":params.attempt_id,
515                    "binding":result.binding,"evidenceDigest":result.evidence_digest,"findings":result.findings})).map_err(|e| e.to_string())?,
516                Some(Producer { kind: ProducerKind::IndependentChecker, run_id: None }))?;
517        }
518        Ok(())
519    }
520}
521
522pub(crate) fn git_object(value: &str) -> Result<GitObject, String> {
523    let algorithm = match value.len() {
524        40 => "sha1",
525        64 => "sha256",
526        _ => return Err("invalid source Git identity".into()),
527    };
528    if !value
529        .bytes()
530        .all(|c| c.is_ascii_digit() || (b'a'..=b'f').contains(&c))
531    {
532        return Err("invalid source Git identity".into());
533    }
534    Ok(GitObject {
535        algorithm: algorithm.into(),
536        value: value.into(),
537    })
538}
539
540fn label(prefix: &str, value: &impl Serialize) -> Result<String, String> {
541    let bytes = serde_json::to_vec(value).map_err(|e| e.to_string())?;
542    Ok(format!("{prefix}-{}", &Digest::of(&bytes).as_str()[7..]))
543}