Skip to main content

kranz_engine/
live_permission.rs

1//! One-call ACP consent. Persisted records are evidence, never response handles.
2//!
3//! The backend owns the live invocation; the engine owns approval authority.
4//! A bounded channel connects them without borrowing the session's output pump.
5
6use 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/// Adapter observations, before the engine attaches mission authority.
26#[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    /// A prohibition is not overridable by the consent UI.
41    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    /// IDs are opaque. Only a unique, offered, one-time kind supplies semantics.
64    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/// Engine-authored identity of the exact policy and workspace that spawned a run.
83#[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/// Capability attribution is deliberately distinct from a verified Slack user.
146#[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    /// The caller must durably record the matching resolution before queuing.
309    /// Queued is not sent: the session emits a separate delivery receipt.
310    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}