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    /// Check source values, including JSON keys and adapter option labels.
64    /// Kept out of historical validation so old event logs still replay.
65    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    /// IDs are opaque. Only a unique, offered, one-time kind supplies semantics.
72    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/// Engine-authored identity of the exact policy and workspace that spawned a run.
91#[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/// Capability attribution is deliberately distinct from a verified Slack user.
161#[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    /// The caller must durably record the matching resolution before queuing.
325    /// Queued is not sent: the session emits a separate delivery receipt.
326    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}