1use crate::error::{EngineError, Result};
7use chrono::{DateTime, Utc};
8use serde::{Deserialize, Serialize};
9use serde_json::Value;
10use sha2::{Digest, Sha256};
11use tokio::sync::mpsc;
12
13pub const MAX_PENDING: usize = 16;
14pub const MAX_MISSION_REQUESTS: usize = 1024;
15pub const MAX_REQUEST_BYTES: usize = 64 * 1024;
16pub const REQUEST_TTL_SECS: i64 = 300;
17
18pub fn digest(value: &impl Serialize) -> Result<String> {
19 Ok(Sha256::digest(serde_json::to_vec(value)?)
20 .iter()
21 .map(|b| format!("{b:02x}"))
22 .collect())
23}
24
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
27#[serde(rename_all = "camelCase")]
28pub struct Proposal {
29 pub id: String,
30 pub engine_session_id: String,
31 pub peer_session_id: String,
32 pub peer_request_id: Value,
33 pub tool_call_id: String,
34 pub action: Value,
35 pub options: Vec<Value>,
36 pub action_digest: String,
37 pub options_digest: String,
38 pub observed_at: DateTime<Utc>,
39 pub deadline: DateTime<Utc>,
40 pub prohibition: Option<String>,
42}
43
44impl Proposal {
45 pub fn validate(&self) -> Result<()> {
46 if self.id.is_empty()
47 || self.engine_session_id.is_empty()
48 || self.peer_session_id.is_empty()
49 || (self.tool_call_id.is_empty() && self.prohibition.is_none())
50 || self.action_digest != digest(&self.action)?
51 || self.options_digest != digest(&self.options)?
52 || self.deadline <= self.observed_at
53 || self.deadline > self.observed_at + chrono::Duration::seconds(REQUEST_TTL_SECS)
54 || serde_json::to_vec(self)?.len() > MAX_REQUEST_BYTES
55 {
56 return Err(EngineError::Backend(
57 "invalid live permission proposal".into(),
58 ));
59 }
60 Ok(())
61 }
62
63 pub fn option(&self, allow: bool) -> Option<String> {
65 let mut ids = std::collections::BTreeSet::new();
66 if self.options.is_empty() || self.options.len() > 32 {
67 return None;
68 }
69 for option in &self.options {
70 let id = option.get("optionId")?.as_str()?;
71 if id.trim().is_empty() || id.len() > 256 || !ids.insert(id) {
72 return None;
73 }
74 }
75 let kind = if allow { "allow_once" } else { "reject_once" };
76 let mut candidates = self.options.iter().filter(|o| o["kind"] == kind);
77 let selected = candidates.next()?.get("optionId")?.as_str()?;
78 candidates.next().is_none().then(|| selected.to_owned())
79 }
80}
81
82#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
84#[serde(rename_all = "camelCase")]
85pub struct Binding {
86 pub mission_id: String,
87 pub run_id: String,
88 pub workspace: String,
89 pub plan_digest: String,
90 pub policy_digest: String,
91}
92
93impl Binding {
94 fn validate(&self) -> Result<()> {
95 let hash = |value: &str| {
96 value.len() == 64
97 && value
98 .bytes()
99 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
100 };
101 if self.mission_id.trim().is_empty()
102 || self.run_id.trim().is_empty()
103 || !std::path::Path::new(&self.workspace).is_absolute()
104 || !hash(&self.plan_digest)
105 || !hash(&self.policy_digest)
106 {
107 return Err(EngineError::Backend(
108 "incomplete live permission authority binding".into(),
109 ));
110 }
111 Ok(())
112 }
113}
114
115#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
116#[serde(rename_all = "camelCase")]
117pub struct Request {
118 pub proposal: Proposal,
119 pub binding: Binding,
120 pub binding_digest: String,
121}
122
123impl Request {
124 pub fn new(proposal: Proposal, binding: Binding) -> Result<Self> {
125 proposal.validate()?;
126 binding.validate()?;
127 let binding_digest = digest(&(&proposal, &binding))?;
128 Ok(Self {
129 proposal,
130 binding,
131 binding_digest,
132 })
133 }
134
135 pub fn validate(&self) -> Result<()> {
136 self.proposal.validate()?;
137 self.binding.validate()?;
138 if self.binding_digest != digest(&(&self.proposal, &self.binding))? {
139 return Err(EngineError::Backend("permission binding changed".into()));
140 }
141 Ok(())
142 }
143}
144
145#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
147#[serde(rename_all = "kebab-case", tag = "kind", content = "id")]
148pub enum Actor {
149 Policy,
150 LocalRepositoryAuthority,
151 LocalMutationCapability,
152 SlackUser(String),
153}
154
155#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
156#[serde(rename_all = "camelCase")]
157pub struct Resolution {
158 pub request_id: String,
159 pub binding_digest: String,
160 pub allow: bool,
161 pub actor: Actor,
162 pub reason: String,
163}
164
165#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
166#[serde(rename_all = "kebab-case")]
167pub enum Delivery {
168 Sent,
169 Uncertain,
170}
171
172#[derive(Debug, Clone, Serialize, Deserialize)]
173#[serde(rename_all = "camelCase")]
174pub struct Record {
175 pub request: Request,
176 #[serde(default, skip_serializing_if = "Option::is_none")]
177 pub resolution: Option<Resolution>,
178 #[serde(default, skip_serializing_if = "Option::is_none")]
179 pub resolved_at: Option<DateTime<Utc>>,
180 #[serde(default, skip_serializing_if = "Option::is_none")]
181 pub delivery: Option<Delivery>,
182 #[serde(default, skip_serializing_if = "Option::is_none")]
183 pub responded_at: Option<DateTime<Utc>>,
184 #[serde(default, skip_serializing_if = "Option::is_none")]
185 pub closed: Option<String>,
186}
187
188impl Record {
189 pub fn pending(&self, now: DateTime<Utc>) -> bool {
190 self.resolution.is_none() && self.closed.is_none() && now < self.request.proposal.deadline
191 }
192
193 pub fn validate_answer(
194 &self,
195 binding_digest: &str,
196 allow: bool,
197 now: DateTime<Utc>,
198 ) -> Result<()> {
199 if !self.pending(now)
200 || self.request.binding_digest != binding_digest
201 || (allow
202 && (self.request.proposal.prohibition.is_some()
203 || self.request.proposal.option(true).is_none()))
204 {
205 return Err(EngineError::InvalidState(
206 "permission is stale, expired, already resolved or prohibited".into(),
207 ));
208 }
209 Ok(())
210 }
211}
212
213pub(crate) fn fold(
214 state: &mut crate::types::MissionState,
215 event: &crate::events::Event,
216) -> Result<()> {
217 use crate::events::EventKind;
218 let invalid =
219 || EngineError::InvalidState("invalid or stale one-call permission transition".into());
220 match &event.kind {
221 EventKind::PermissionRequested { request } => {
222 request.validate()?;
223 let run = state
224 .runs
225 .get(&request.binding.run_id)
226 .ok_or_else(invalid)?;
227 if request.binding.mission_id != state.mission.id
228 || request.proposal.engine_session_id != run.sdk_session_id
229 || run.role != crate::types::Role::Worker
230 || run.ended_at.is_some()
231 || event.ts < request.proposal.observed_at
232 || event.ts >= request.proposal.deadline
233 || state.permissions.contains_key(&request.proposal.id)
234 || state.permissions.len() >= MAX_MISSION_REQUESTS
235 || state
236 .permissions
237 .values()
238 .filter(|r| r.closed.is_none() && r.delivery.is_none())
239 .count()
240 >= MAX_PENDING
241 {
242 return Err(invalid());
243 }
244 state.permissions.insert(
245 request.proposal.id.clone(),
246 Record {
247 request: request.clone(),
248 resolution: None,
249 resolved_at: None,
250 delivery: None,
251 responded_at: None,
252 closed: None,
253 },
254 );
255 }
256 EventKind::PermissionResolved { resolution } => {
257 let record = state
258 .permissions
259 .get_mut(&resolution.request_id)
260 .ok_or_else(invalid)?;
261 if !record.pending(event.ts)
262 || resolution.binding_digest != record.request.binding_digest
263 || (resolution.allow
264 && (record.request.proposal.prohibition.is_some()
265 || record.request.proposal.option(true).is_none()
266 || resolution.actor == Actor::Policy))
267 || resolution.reason.trim().is_empty()
268 || matches!(&resolution.actor, Actor::SlackUser(id) if id.trim().is_empty())
269 {
270 return Err(invalid());
271 }
272 record.resolution = Some(resolution.clone());
273 record.resolved_at = Some(event.ts);
274 }
275 EventKind::PermissionResponseRecorded {
276 request_id,
277 delivery,
278 } => {
279 let record = state.permissions.get_mut(request_id).ok_or_else(invalid)?;
280 if record.resolution.is_none() || record.delivery.is_some() || record.closed.is_some() {
281 return Err(invalid());
282 }
283 record.delivery = Some(delivery.clone());
284 record.responded_at = Some(event.ts);
285 }
286 EventKind::PermissionClosed { request_id, reason } => {
287 let record = state.permissions.get_mut(request_id).ok_or_else(invalid)?;
288 if record.closed.is_some() || reason.trim().is_empty() {
289 return Err(invalid());
290 }
291 record.closed = Some(reason.clone());
292 }
293 _ => return Err(invalid()),
294 }
295 Ok(())
296}
297
298#[derive(Debug)]
299pub(crate) struct Answer {
300 pub proposal: Proposal,
301 pub allow: bool,
302}
303
304#[derive(Debug, Clone)]
305pub struct PermissionResponder(mpsc::Sender<Answer>);
306
307impl PermissionResponder {
308 pub fn respond(&self, proposal: &Proposal, allow: bool) -> Result<()> {
311 self.0
312 .try_send(Answer {
313 proposal: proposal.clone(),
314 allow,
315 })
316 .map_err(|_| {
317 EngineError::Backend("permission response channel unavailable or full".into())
318 })
319 }
320
321 pub(crate) fn channel() -> (Self, mpsc::Receiver<Answer>) {
322 let (tx, rx) = mpsc::channel(MAX_PENDING);
323 (Self(tx), rx)
324 }
325}