Skip to main content

kranz_engine/gate_evaluation/
lifecycle.rs

1//! Pure audit transitions. A checker result is evidence; stage disposition,
2//! consent and consumption are separately authored by the engine.
3use super::protocol::{
4    Binding, Digest, EvaluationResult, Id, Request, Response, Stage, Status, SuccessResponse,
5    Verdict, WirePath,
6};
7use crate::pack::evaluator::{Enforcement, Kind};
8use chrono::{DateTime, Utc};
9use serde::{Deserialize, Serialize};
10
11pub const MAX_ATTEMPTS: usize = 2048;
12
13#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
14#[serde(rename_all = "camelCase", deny_unknown_fields)]
15pub struct Policy {
16    pub kind: Kind,
17    pub enforcement: Enforcement,
18    pub mission_policy_digest: Digest,
19    /// Only nonwaivable prerequisites; advisory diagnostics do not become floors.
20    pub mechanical_prerequisites_passed: bool,
21}
22impl Policy {
23    pub fn digest(&self) -> Digest {
24        Digest::of(&serde_json::to_vec(self).expect("typed policy is serializable"))
25    }
26}
27
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
29#[serde(rename_all = "camelCase", deny_unknown_fields)]
30pub struct RetainedArtifact {
31    /// Mission-relative, engine-selected path. Missing bytes remain unresolved.
32    pub path: WirePath,
33    pub raw_digest: Digest,
34    pub retained_digest: Digest,
35    pub retained_bytes: u64,
36    pub transformation: String,
37}
38
39#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
40#[serde(rename_all = "camelCase", deny_unknown_fields)]
41pub struct Requested {
42    pub request: Request,
43    pub policy: Policy,
44    pub retained_inputs: Vec<RetainedArtifact>,
45    #[serde(default, skip_serializing_if = "Option::is_none")]
46    pub permission_request_id: Option<String>,
47}
48
49#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
50#[serde(tag = "status", rename_all = "kebab-case", deny_unknown_fields)]
51pub enum Outcome {
52    Evaluated {
53        result: Box<EvaluationResult>,
54        #[serde(rename = "rawStdoutDigest")]
55        raw_stdout_digest: Digest,
56    },
57    Error {
58        message: String,
59    },
60}
61
62#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
63#[serde(rename_all = "camelCase", deny_unknown_fields)]
64pub struct Finished {
65    pub attempt_id: Id,
66    pub outcome: Outcome,
67    pub exit_code: Option<i32>,
68    pub cleanup_confirmed: bool,
69    pub artifacts: Vec<RetainedArtifact>,
70}
71
72#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
73#[serde(rename_all = "kebab-case")]
74pub enum Disposition {
75    Proceed,
76    Block,
77    RequireHuman,
78}
79
80#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
81#[serde(rename_all = "camelCase", deny_unknown_fields)]
82pub struct Consent {
83    pub actor: crate::live_permission::Actor,
84    pub allow: bool,
85    /// Existing authenticated stage decision, separate from checker output.
86    pub reference: String,
87}
88
89#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
90#[serde(rename_all = "camelCase", deny_unknown_fields)]
91pub struct Resolution {
92    pub id: Id,
93    pub attempt_id: Id,
94    pub binding: Binding,
95    pub disposition: Disposition,
96    pub rationale: String,
97    #[serde(default, skip_serializing_if = "Option::is_none")]
98    pub consent: Option<Consent>,
99}
100
101#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
102#[serde(rename_all = "kebab-case")]
103pub enum Action {
104    ApprovePlan,
105    AnswerPermission,
106    AcceptMilestone,
107    AcceptDeliverable,
108    AdvanceLocalBase,
109}
110impl Action {
111    fn stage(&self) -> Stage {
112        match self {
113            Self::ApprovePlan => Stage::PlanApproval,
114            Self::AnswerPermission => Stage::CommandPermission,
115            Self::AcceptMilestone => Stage::MilestoneValidation,
116            Self::AcceptDeliverable => Stage::FinalGate,
117            Self::AdvanceLocalBase => Stage::Merge,
118        }
119    }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
123#[serde(rename_all = "camelCase", deny_unknown_fields)]
124pub struct Consumed {
125    pub attempt_id: Id,
126    pub resolution_id: Id,
127    pub rechecked_binding: Binding,
128    pub action: Action,
129}
130
131#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
132#[serde(rename_all = "camelCase")]
133pub struct Record {
134    pub requested: Requested,
135    pub requested_at: DateTime<Utc>,
136    pub requested_seq: u64,
137    #[serde(default, skip_serializing_if = "Option::is_none")]
138    pub finished: Option<Finished>,
139    #[serde(default, skip_serializing_if = "Option::is_none")]
140    pub finished_at: Option<DateTime<Utc>>,
141    #[serde(default, skip_serializing_if = "Option::is_none")]
142    pub resolution: Option<Resolution>,
143    #[serde(default, skip_serializing_if = "Option::is_none")]
144    pub resolved_at: Option<DateTime<Utc>>,
145    #[serde(default, skip_serializing_if = "Option::is_none")]
146    pub consumed: Option<Consumed>,
147    #[serde(default, skip_serializing_if = "Option::is_none")]
148    pub closed: Option<String>,
149}
150impl Record {
151    pub fn new(requested: Requested, at: DateTime<Utc>, seq: u64) -> Result<Self, String> {
152        requested.request.validate()?;
153        if requested.policy.digest() != requested.request.params.binding.policy_digest {
154            return Err("gate policy bytes differ from the request binding".into());
155        }
156        if requested.policy.kind == Kind::Judgment
157            && !requested.policy.mechanical_prerequisites_passed
158        {
159            return Err("judgment requires successful mechanical prerequisites".into());
160        }
161        if at >= deadline(&requested.request)? {
162            return Err("gate request is already expired".into());
163        }
164        validate_artifacts(&requested.retained_inputs)?;
165        Ok(Self {
166            requested,
167            requested_at: at,
168            requested_seq: seq,
169            finished: None,
170            finished_at: None,
171            resolution: None,
172            resolved_at: None,
173            consumed: None,
174            closed: None,
175        })
176    }
177
178    pub fn finish(&mut self, finished: Finished, at: DateTime<Utc>) -> Result<(), String> {
179        let request = &self.requested.request;
180        if self.closed.is_some()
181            || self.finished.is_some()
182            || finished.attempt_id != request.params.attempt_id
183            || at < self.requested_at
184        {
185            return Err("duplicate, foreign or out-of-order gate result".into());
186        }
187        validate_artifacts(&finished.artifacts)?;
188        match &finished.outcome {
189            Outcome::Evaluated { result, .. } => {
190                if at >= deadline(request)?
191                    || !finished.cleanup_confirmed
192                    || finished.exit_code != Some(0)
193                {
194                    return Err("evaluation expired or process completion was not proved".into());
195                }
196                let wire = Response::Result(Box::new(SuccessResponse {
197                    jsonrpc: "2.0".into(),
198                    id: finished.attempt_id.clone(),
199                    result: *result.clone(),
200                }));
201                Response::from_bytes(&serde_json::to_vec(&wire).map_err(|e| e.to_string())?)?
202                    .correlate(request)?;
203            }
204            Outcome::Error { message } => bounded(message, 8192)?,
205        }
206        self.finished = Some(finished);
207        self.finished_at = Some(at);
208        Ok(())
209    }
210
211    pub fn disposition(&self, consent: Option<&Consent>) -> Result<Disposition, String> {
212        if self.closed.is_some() {
213            return Ok(Disposition::Block);
214        }
215        let finished = self
216            .finished
217            .as_ref()
218            .ok_or("gate has no terminal attempt result")?;
219        if let Some(consent) = consent {
220            bounded(&consent.reference, 1024)?;
221            match &consent.actor {
222                crate::live_permission::Actor::Policy => {
223                    return Err("policy is not operator consent".into())
224                }
225                crate::live_permission::Actor::SlackUser(id) => bounded(id, 256)?,
226                _ => {}
227            }
228            if !consent.allow {
229                return Ok(Disposition::Block);
230            }
231        }
232        if !self.requested.policy.mechanical_prerequisites_passed || !finished.cleanup_confirmed {
233            return Ok(Disposition::Block);
234        }
235        match &finished.outcome {
236            Outcome::Evaluated { result, .. } if result.status == Status::Escalate => {
237                return Ok(Disposition::RequireHuman)
238            }
239            Outcome::Evaluated { result, .. }
240                if result.verdict == Some(Verdict::Fail)
241                    && self.requested.policy.enforcement == Enforcement::Blocking =>
242            {
243                return Ok(Disposition::Block)
244            }
245            Outcome::Error { .. } if self.requested.policy.enforcement == Enforcement::Blocking => {
246                return Ok(Disposition::Block)
247            }
248            _ => {}
249        }
250        if matches!(
251            self.requested.request.params.stage,
252            Stage::PlanApproval | Stage::CommandPermission | Stage::Merge
253        ) && consent.is_none()
254        {
255            return Ok(Disposition::RequireHuman);
256        }
257        Ok(Disposition::Proceed)
258    }
259
260    pub fn resolve(&mut self, resolution: Resolution, at: DateTime<Utc>) -> Result<(), String> {
261        let request = &self.requested.request;
262        if self.closed.is_some()
263            || self.finished_at.is_none_or(|finished| at < finished)
264            || self.resolution.is_some()
265            || resolution.attempt_id != request.params.attempt_id
266            || resolution.binding != request.params.binding
267        {
268            return Err("duplicate, foreign or changed gate resolution".into());
269        }
270        bounded(&resolution.rationale, 8192)?;
271        if resolution.disposition != self.disposition(resolution.consent.as_ref())? {
272            return Err("gate disposition violates recorded policy or required consent".into());
273        }
274        self.resolution = Some(resolution);
275        self.resolved_at = Some(at);
276        Ok(())
277    }
278
279    pub fn consume(&mut self, consumed: Consumed, at: DateTime<Utc>) -> Result<(), String> {
280        let resolution = self.resolution.as_ref().ok_or("gate has no resolution")?;
281        let request = &self.requested.request;
282        if self.closed.is_some()
283            || self.consumed.is_some()
284            || resolution.disposition != Disposition::Proceed
285            || at >= deadline(request)?
286            || self.resolved_at.is_none_or(|resolved| at < resolved)
287            || consumed.attempt_id != request.params.attempt_id
288            || consumed.resolution_id != resolution.id
289            || consumed.rechecked_binding != request.params.binding
290            || consumed.action.stage() != request.params.stage
291        {
292            return Err("stale, duplicate or mismatched gate consumption".into());
293        }
294        self.consumed = Some(consumed);
295        Ok(())
296    }
297
298    pub fn close(&mut self, reason: String) -> Result<(), String> {
299        if self.closed.is_some() || self.consumed.is_some() {
300            return Err("gate attempt is already closed or consumed".into());
301        }
302        bounded(&reason, 8192)?;
303        self.closed = Some(reason);
304        Ok(())
305    }
306}
307
308fn deadline(request: &Request) -> Result<DateTime<Utc>, String> {
309    DateTime::parse_from_rfc3339(&request.params.deadline)
310        .map(|d| d.with_timezone(&Utc))
311        .map_err(|e| e.to_string())
312}
313fn bounded(value: &str, max: usize) -> Result<(), String> {
314    if value.trim().is_empty() || value.len() > max {
315        Err("invalid bounded gate text".into())
316    } else {
317        Ok(())
318    }
319}
320fn validate_artifacts(artifacts: &[RetainedArtifact]) -> Result<(), String> {
321    if artifacts.len() > 10_000 {
322        return Err("gate retention inventory exceeds limit".into());
323    }
324    super::protocol::validate_paths(artifacts.iter().map(|a| a.path.as_str()))?;
325    let mut names = std::collections::BTreeSet::new();
326    for artifact in artifacts {
327        if artifact.retained_bytes > 64 * 1024 * 1024
328            || !artifact.path.as_str().starts_with("runs/gates/")
329            || !names.insert(artifact.path.as_str().to_lowercase())
330        {
331            return Err("invalid retained gate artifact".into());
332        }
333        bounded(&artifact.transformation, 512)?;
334    }
335    Ok(())
336}
337
338pub(crate) fn fold(
339    state: &mut crate::types::MissionState,
340    event: &crate::events::Event,
341) -> crate::error::Result<()> {
342    use crate::events::EventKind;
343    let invalid = |message: String| crate::error::EngineError::InvalidState(message);
344    if event.mission_id != state.mission.id {
345        return Err(invalid("foreign gate mission".into()));
346    }
347    match &event.kind {
348        EventKind::GateEvaluationClosed { attempt_id, reason } => {
349            state
350                .gate_evaluations
351                .get_mut(attempt_id.as_str())
352                .ok_or_else(|| invalid("closed gate has no request".into()))?
353                .close(reason.clone())
354                .map_err(invalid)?;
355        }
356        EventKind::GateEvaluationRequested { evaluation } => {
357            let request = &evaluation.request;
358            validate_stage_status(state, request.params.stage).map_err(invalid)?;
359            if state.gate_evaluations.len() >= MAX_ATTEMPTS
360                || request.params.mission_id.as_str() != event.mission_id
361                || state
362                    .gate_evaluations
363                    .contains_key(request.params.attempt_id.as_str())
364            {
365                return Err(invalid("duplicate, foreign or excess gate request".into()));
366            }
367            validate_permission_join(state, evaluation).map_err(invalid)?;
368            let record = Record::new(*evaluation.clone(), event.ts, event.seq).map_err(invalid)?;
369            state
370                .gate_evaluations
371                .insert(request.params.attempt_id.as_str().into(), record);
372        }
373        EventKind::GateEvaluationFinished { evaluation } => {
374            state
375                .gate_evaluations
376                .get_mut(evaluation.attempt_id.as_str())
377                .ok_or_else(|| invalid("gate result has no request".into()))?
378                .finish(*evaluation.clone(), event.ts)
379                .map_err(invalid)?;
380        }
381        EventKind::GateResolutionRecorded { resolution } => {
382            let record = state
383                .gate_evaluations
384                .get(resolution.attempt_id.as_str())
385                .ok_or_else(|| invalid("gate resolution has no request".into()))?;
386            if let (Some(id), Some(consent)) =
387                (&record.requested.permission_request_id, &resolution.consent)
388            {
389                let actual = state
390                    .permissions
391                    .get(id)
392                    .and_then(|permission| permission.resolution.as_ref())
393                    .ok_or_else(|| invalid("permission consent is absent".into()))?;
394                if actual.actor != consent.actor
395                    || actual.allow != consent.allow
396                    || consent.reference != *id
397                {
398                    return Err(invalid(
399                        "gate cannot replace or relabel permission consent".into(),
400                    ));
401                }
402            }
403            if state.gate_evaluations.values().any(|record| {
404                record
405                    .resolution
406                    .as_ref()
407                    .is_some_and(|r| r.id == resolution.id)
408            }) {
409                return Err(invalid("gate resolution ID was reused".into()));
410            }
411            state
412                .gate_evaluations
413                .get_mut(resolution.attempt_id.as_str())
414                .ok_or_else(|| invalid("gate resolution has no request".into()))?
415                .resolve(resolution.clone(), event.ts)
416                .map_err(invalid)?;
417        }
418        EventKind::GateResolutionConsumed { consumption } => {
419            let record = state
420                .gate_evaluations
421                .get(consumption.attempt_id.as_str())
422                .ok_or_else(|| invalid("gate consumption has no request".into()))?;
423            validate_stage_status(state, record.requested.request.params.stage).map_err(invalid)?;
424            validate_permission_join(state, &record.requested).map_err(invalid)?;
425            if let Some(id) = &record.requested.permission_request_id {
426                let permission = &state.permissions[id];
427                let actual = permission
428                    .resolution
429                    .as_ref()
430                    .ok_or_else(|| invalid("permission consent is absent".into()))?;
431                let consent = record
432                    .resolution
433                    .as_ref()
434                    .and_then(|r| r.consent.as_ref())
435                    .ok_or_else(|| invalid("permission consent join is absent".into()))?;
436                if !actual.allow
437                    || !consent.allow
438                    || actual.actor != consent.actor
439                    || consent.reference != *id
440                {
441                    return Err(invalid(
442                        "gate cannot replace or relabel permission consent".into(),
443                    ));
444                }
445            }
446            if state
447                .consumed_gate_resolutions
448                .contains(consumption.resolution_id.as_str())
449            {
450                return Err(invalid("gate resolution was already consumed".into()));
451            }
452            state
453                .gate_evaluations
454                .get_mut(consumption.attempt_id.as_str())
455                .ok_or_else(|| invalid("gate consumption has no request".into()))?
456                .consume(consumption.clone(), event.ts)
457                .map_err(invalid)?;
458            state
459                .consumed_gate_resolutions
460                .insert(consumption.resolution_id.as_str().into());
461        }
462        _ => return Err(invalid("not a gate lifecycle event".into())),
463    }
464    Ok(())
465}
466
467fn validate_permission_join(
468    state: &crate::types::MissionState,
469    evaluation: &Requested,
470) -> Result<(), String> {
471    use super::protocol::Subject;
472    match (
473        &evaluation.request.params.subject,
474        &evaluation.permission_request_id,
475    ) {
476        (
477            Subject::Invocation {
478                run_id,
479                peer_session_id,
480                tool_call_id,
481                peer_request_id,
482                action_digest,
483                options_digest,
484                cwd_id,
485            },
486            Some(id),
487        ) => {
488            let record = state
489                .permissions
490                .get(id)
491                .ok_or("gate invocation has no live permission request")?;
492            let permission = &record.request;
493            let run = state
494                .runs
495                .get(&permission.binding.run_id)
496                .ok_or("permission run is absent")?;
497            let workspace = workspace_id(&permission.binding.workspace);
498            let peer = serde_json::to_value(peer_request_id).map_err(|e| e.to_string())?;
499            if record.closed.is_some()
500                || record.delivery.is_some()
501                || run.ended_at.is_some()
502                || evaluation.request.params.binding.workspace_id != workspace
503                || *cwd_id != workspace
504                || evaluation.request.params.binding.plan_digest.as_str()
505                    != format!("sha256:{}", permission.binding.plan_digest)
506                || evaluation.policy.mission_policy_digest.as_str()
507                    != format!("sha256:{}", permission.binding.policy_digest)
508                || deadline(&evaluation.request)? > permission.proposal.deadline
509                || permission.binding.run_id != run_id.as_str()
510                || permission.proposal.peer_session_id != peer_session_id.as_str()
511                || permission.proposal.tool_call_id != tool_call_id.as_str()
512                || permission.proposal.peer_request_id != peer
513                || action_digest.as_str() != format!("sha256:{}", permission.proposal.action_digest)
514                || options_digest.as_str()
515                    != format!("sha256:{}", permission.proposal.options_digest)
516            {
517                return Err("gate invocation does not match its permission request".into());
518            }
519        }
520        (Subject::Invocation { .. }, None) => {
521            return Err("gate invocation requires a permission join".into())
522        }
523        (_, Some(_)) => return Err("only invocation evaluations may join a permission".into()),
524        _ => {}
525    }
526    Ok(())
527}
528
529pub fn workspace_id(path: &str) -> Id {
530    Id::try_from(format!(
531        "workspace-{}",
532        &Digest::of(path.as_bytes()).as_str()[7..]
533    ))
534    .expect("a SHA-256 workspace label is a valid opaque ID")
535}
536
537fn validate_stage_status(state: &crate::types::MissionState, stage: Stage) -> Result<(), String> {
538    use crate::types::MissionStatus;
539    let ready = match stage {
540        Stage::PlanApproval => {
541            state.mission.status == MissionStatus::Planning
542                || (state.pending_revision.is_some()
543                    && matches!(
544                        state.mission.status,
545                        MissionStatus::Running | MissionStatus::Blocked
546                    ))
547        }
548        Stage::CommandPermission | Stage::MilestoneValidation => matches!(
549            state.mission.status,
550            MissionStatus::Running | MissionStatus::Validating
551        ),
552        Stage::FinalGate => state.mission.status == MissionStatus::Validating,
553        Stage::Merge => state.mission.status == MissionStatus::Complete,
554    };
555    if ready {
556        Ok(())
557    } else {
558        Err("gate stage does not match the mission lifecycle".into())
559    }
560}