1use 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 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 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 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_ascii_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}