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 ambiguous_display(&self) -> bool {
66 serde_json::to_value(self)
67 .map(|value| crate::presentation::has_ambiguous_json(&value))
68 .unwrap_or(true)
69 }
70
71 pub fn option(&self, allow: bool) -> Option<String> {
73 let mut ids = std::collections::BTreeSet::new();
74 if self.options.is_empty() || self.options.len() > 32 {
75 return None;
76 }
77 for option in &self.options {
78 let id = option.get("optionId")?.as_str()?;
79 if id.trim().is_empty() || id.len() > 256 || !ids.insert(id) {
80 return None;
81 }
82 }
83 let kind = if allow { "allow_once" } else { "reject_once" };
84 let mut candidates = self.options.iter().filter(|o| o["kind"] == kind);
85 let selected = candidates.next()?.get("optionId")?.as_str()?;
86 candidates.next().is_none().then(|| selected.to_owned())
87 }
88}
89
90#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
92#[serde(rename_all = "camelCase")]
93pub struct Binding {
94 pub mission_id: String,
95 pub run_id: String,
96 pub workspace: String,
97 pub plan_digest: String,
98 pub policy_digest: String,
99}
100
101impl Binding {
102 fn validate(&self) -> Result<()> {
103 let hash = |value: &str| {
104 value.len() == 64
105 && value
106 .bytes()
107 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
108 };
109 if self.mission_id.trim().is_empty()
110 || self.run_id.trim().is_empty()
111 || !std::path::Path::new(&self.workspace).is_absolute()
112 || !hash(&self.plan_digest)
113 || !hash(&self.policy_digest)
114 {
115 return Err(EngineError::Backend(
116 "incomplete live permission authority binding".into(),
117 ));
118 }
119 Ok(())
120 }
121}
122
123#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
124#[serde(rename_all = "camelCase")]
125pub struct Request {
126 pub proposal: Proposal,
127 pub binding: Binding,
128 pub binding_digest: String,
129}
130
131impl Request {
132 pub fn new(proposal: Proposal, binding: Binding) -> Result<Self> {
133 proposal.validate()?;
134 binding.validate()?;
135 let binding_digest = digest(&(&proposal, &binding))?;
136 Ok(Self {
137 proposal,
138 binding,
139 binding_digest,
140 })
141 }
142
143 pub fn ambiguous_display(&self) -> bool {
144 self.proposal.ambiguous_display()
145 || serde_json::to_value(&self.binding)
146 .map(|value| crate::presentation::has_ambiguous_json(&value))
147 .unwrap_or(true)
148 }
149
150 pub fn validate(&self) -> Result<()> {
151 self.proposal.validate()?;
152 self.binding.validate()?;
153 if self.binding_digest != digest(&(&self.proposal, &self.binding))? {
154 return Err(EngineError::Backend("permission binding changed".into()));
155 }
156 Ok(())
157 }
158}
159
160#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
162#[serde(rename_all = "kebab-case", tag = "kind", content = "id")]
163pub enum Actor {
164 Policy,
165 LocalRepositoryAuthority,
166 LocalMutationCapability,
167 SlackUser(String),
168}
169
170#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
171#[serde(rename_all = "camelCase")]
172pub struct Resolution {
173 pub request_id: String,
174 pub binding_digest: String,
175 pub allow: bool,
176 pub actor: Actor,
177 pub reason: String,
178}
179
180#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
181#[serde(rename_all = "kebab-case")]
182pub enum Delivery {
183 Sent,
184 Uncertain,
185}
186
187#[derive(Debug, Clone, Serialize, Deserialize)]
188#[serde(rename_all = "camelCase")]
189pub struct Record {
190 pub request: Request,
191 #[serde(default, skip_serializing_if = "Option::is_none")]
192 pub resolution: Option<Resolution>,
193 #[serde(default, skip_serializing_if = "Option::is_none")]
194 pub resolved_at: Option<DateTime<Utc>>,
195 #[serde(default, skip_serializing_if = "Option::is_none")]
196 pub delivery: Option<Delivery>,
197 #[serde(default, skip_serializing_if = "Option::is_none")]
198 pub responded_at: Option<DateTime<Utc>>,
199 #[serde(default, skip_serializing_if = "Option::is_none")]
200 pub closed: Option<String>,
201}
202
203impl Record {
204 pub fn pending(&self, now: DateTime<Utc>) -> bool {
205 self.resolution.is_none() && self.closed.is_none() && now < self.request.proposal.deadline
206 }
207
208 pub fn validate_answer(
209 &self,
210 binding_digest: &str,
211 allow: bool,
212 now: DateTime<Utc>,
213 ) -> Result<()> {
214 if !self.pending(now)
215 || self.request.binding_digest != binding_digest
216 || (allow
217 && (self.request.ambiguous_display()
218 || self.request.proposal.prohibition.is_some()
219 || self.request.proposal.option(true).is_none()))
220 {
221 return Err(EngineError::InvalidState(
222 "permission is stale, expired, already resolved or prohibited".into(),
223 ));
224 }
225 Ok(())
226 }
227}
228
229pub(crate) fn fold(
230 state: &mut crate::types::MissionState,
231 event: &crate::events::Event,
232) -> Result<()> {
233 use crate::events::EventKind;
234 let invalid =
235 || EngineError::InvalidState("invalid or stale one-call permission transition".into());
236 match &event.kind {
237 EventKind::PermissionRequested { request } => {
238 request.validate()?;
239 let run = state
240 .runs
241 .get(&request.binding.run_id)
242 .ok_or_else(invalid)?;
243 if request.binding.mission_id != state.mission.id
244 || request.proposal.engine_session_id != run.sdk_session_id
245 || run.role != crate::types::Role::Worker
246 || run.ended_at.is_some()
247 || event.ts < request.proposal.observed_at
248 || event.ts >= request.proposal.deadline
249 || state.permissions.contains_key(&request.proposal.id)
250 || state.permissions.len() >= MAX_MISSION_REQUESTS
251 || state
252 .permissions
253 .values()
254 .filter(|r| r.closed.is_none() && r.delivery.is_none())
255 .count()
256 >= MAX_PENDING
257 {
258 return Err(invalid());
259 }
260 state.permissions.insert(
261 request.proposal.id.clone(),
262 Record {
263 request: request.clone(),
264 resolution: None,
265 resolved_at: None,
266 delivery: None,
267 responded_at: None,
268 closed: None,
269 },
270 );
271 }
272 EventKind::PermissionResolved { resolution } => {
273 let record = state
274 .permissions
275 .get_mut(&resolution.request_id)
276 .ok_or_else(invalid)?;
277 if !record.pending(event.ts)
278 || resolution.binding_digest != record.request.binding_digest
279 || (resolution.allow
280 && (record.request.proposal.prohibition.is_some()
281 || record.request.proposal.option(true).is_none()
282 || resolution.actor == Actor::Policy))
283 || resolution.reason.trim().is_empty()
284 || matches!(&resolution.actor, Actor::SlackUser(id) if id.trim().is_empty())
285 {
286 return Err(invalid());
287 }
288 record.resolution = Some(resolution.clone());
289 record.resolved_at = Some(event.ts);
290 }
291 EventKind::PermissionResponseRecorded {
292 request_id,
293 delivery,
294 } => {
295 let record = state.permissions.get_mut(request_id).ok_or_else(invalid)?;
296 if record.resolution.is_none() || record.delivery.is_some() || record.closed.is_some() {
297 return Err(invalid());
298 }
299 record.delivery = Some(delivery.clone());
300 record.responded_at = Some(event.ts);
301 }
302 EventKind::PermissionClosed { request_id, reason } => {
303 let record = state.permissions.get_mut(request_id).ok_or_else(invalid)?;
304 if record.closed.is_some() || reason.trim().is_empty() {
305 return Err(invalid());
306 }
307 record.closed = Some(reason.clone());
308 }
309 _ => return Err(invalid()),
310 }
311 Ok(())
312}
313
314#[derive(Debug)]
315pub(crate) struct Answer {
316 pub proposal: Proposal,
317 pub allow: bool,
318}
319
320#[derive(Debug, Clone)]
321pub struct PermissionResponder(mpsc::Sender<Answer>);
322
323impl PermissionResponder {
324 pub fn respond(&self, proposal: &Proposal, allow: bool) -> Result<()> {
327 if allow && (proposal.ambiguous_display() || proposal.prohibition.is_some()) {
328 return Err(EngineError::InvalidState(
329 "ambiguous or prohibited permission cannot be allowed".into(),
330 ));
331 }
332 self.0
333 .try_send(Answer {
334 proposal: proposal.clone(),
335 allow,
336 })
337 .map_err(|_| {
338 EngineError::Backend("permission response channel unavailable or full".into())
339 })
340 }
341
342 pub(crate) fn channel() -> (Self, mpsc::Receiver<Answer>) {
343 let (tx, rx) = mpsc::channel(MAX_PENDING);
344 (Self(tx), rx)
345 }
346}