Skip to main content

atman_runtime/
permission.rs

1use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque};
2use std::path::PathBuf;
3use std::sync::{Arc, Mutex};
4
5use chrono::{DateTime, Utc};
6use serde::{Deserialize, Serialize};
7use tokio::sync::{oneshot, watch};
8use uuid::Uuid;
9
10use crate::event::FlowRunId;
11use crate::flow_authority::{FlowExecutionState, FlowIdentity};
12use crate::tool::{PathOrigin, Tier};
13use crate::tools::agent_ctrl::FlowRegistry;
14use crate::trust::{ExecutionPolicy, PolicyAction, RiskKind, TrustConfig};
15
16pub const ANCESTOR_OFFER_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
17const RECENT_TERMINAL_REQUEST_LIMIT: usize = 256;
18
19macro_rules! permission_id {
20    ($name:ident) => {
21        #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
22        #[serde(transparent)]
23        pub struct $name(pub Uuid);
24
25        impl $name {
26            pub fn now() -> Self {
27                Self(Uuid::now_v7())
28            }
29        }
30
31        impl std::fmt::Display for $name {
32            fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
33                self.0.fmt(f)
34            }
35        }
36    };
37}
38
39permission_id!(PermissionRequestId);
40permission_id!(PermissionDecisionId);
41permission_id!(PermissionGrantId);
42permission_id!(PermissionGroupId);
43
44#[derive(Debug, Clone, PartialEq, Eq)]
45pub enum GroupOwner {
46    Flow(FlowRunId),
47    User,
48    System,
49}
50
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct PermissionGroup {
53    pub group_id: PermissionGroupId,
54    pub owner: GroupOwner,
55    pub label: String,
56    pub request_ids: BTreeSet<PermissionRequestId>,
57    pub member_revisions: BTreeMap<PermissionRequestId, u64>,
58    pub created_at: DateTime<Utc>,
59    pub revision: u64,
60}
61
62#[derive(Debug, Clone, Default, PartialEq, Eq)]
63pub struct ResourceProvenance {
64    pub cwd: Option<std::path::PathBuf>,
65    pub path: Option<std::path::PathBuf>,
66    pub path_origin: Option<PathOrigin>,
67    pub workspace_id: Option<String>,
68    pub workspace_root: Option<std::path::PathBuf>,
69    pub repository_root: Option<std::path::PathBuf>,
70    pub network: bool,
71    pub risks: BTreeSet<RiskKind>,
72    extra_targets: Vec<std::path::PathBuf>,
73}
74
75impl ResourceProvenance {
76    pub fn none() -> Self {
77        Self::default()
78    }
79
80    pub fn for_ctx(ctx: &crate::tool::ToolCtx) -> Self {
81        Self {
82            workspace_id: ctx.workspace.as_ref().map(|w| w.workspace_id.clone()),
83            workspace_root: ctx.workspace.as_ref().map(|w| w.path.clone()),
84            repository_root: ctx.workspace.as_ref().map(|w| w.repository_root.clone()),
85            ..Self::default()
86        }
87    }
88
89    pub fn with_cwd(
90        mut self,
91        ctx: &crate::tool::ToolCtx,
92        explicit: Option<&std::path::Path>,
93    ) -> Result<Self, crate::error::RuntimeError> {
94        let resolved = ctx.resolve_cwd_with_origin(explicit)?;
95        self.cwd = Some(resolved.path);
96        self.path_origin = Some(resolved.origin);
97        Ok(self)
98    }
99
100    pub fn with_path(
101        mut self,
102        ctx: &crate::tool::ToolCtx,
103        path: &std::path::Path,
104    ) -> Result<Self, crate::error::RuntimeError> {
105        let resolved = ctx.resolve_path_with_origin(path)?;
106        self.path = Some(resolved.path);
107        self.path_origin = Some(resolved.origin);
108        Ok(self)
109    }
110
111    pub fn with_network(mut self) -> Self {
112        self.network = true;
113        self
114    }
115
116    pub fn with_risk(mut self, risk: RiskKind) -> Self {
117        self.risks.insert(risk);
118        self
119    }
120
121    pub fn with_extra_target(
122        mut self,
123        ctx: &crate::tool::ToolCtx,
124        path: &std::path::Path,
125    ) -> Result<Self, crate::error::RuntimeError> {
126        let resolved = ctx.resolve_path_with_origin(path)?;
127        // Any external target makes the whole invocation external.
128        if matches!(resolved.origin, PathOrigin::ExplicitExternal) {
129            self.path_origin = Some(resolved.origin);
130        }
131        self.extra_targets.push(resolved.path);
132        Ok(self)
133    }
134
135    /// Every resource this provenance authorizes, primary and extra alike. The
136    /// permit check compares against this set, so a target the tool never
137    /// declared can never be written under another target's approval.
138    pub fn authorized_targets(&self) -> impl Iterator<Item = &std::path::Path> {
139        self.path
140            .as_deref()
141            .into_iter()
142            .chain(self.cwd.as_deref())
143            .chain(self.extra_targets.iter().map(|p| p.as_path()))
144    }
145
146    pub fn is_external(&self) -> bool {
147        matches!(self.path_origin, Some(PathOrigin::ExplicitExternal))
148    }
149
150    pub fn is_unbound(&self) -> bool {
151        matches!(self.path_origin, Some(PathOrigin::Unbound))
152    }
153
154    pub fn workspace_relative_path(&self) -> Option<String> {
155        if self.is_external() || self.is_unbound() || self.authorized_targets().count() != 1 {
156            return None;
157        }
158        let root = self.workspace_root.as_deref()?;
159        let relative = self.authorized_targets().next()?.strip_prefix(root).ok()?;
160        if relative.as_os_str().is_empty()
161            || relative.components().any(|component| {
162                matches!(
163                    component,
164                    std::path::Component::ParentDir | std::path::Component::RootDir
165                )
166            })
167        {
168            return None;
169        }
170        relative.to_str().map(str::to_owned)
171    }
172}
173
174#[derive(Debug, Clone, PartialEq, Eq)]
175pub struct PermissionIntent {
176    pub tool_use_id: String,
177    pub tool_name: String,
178    pub call_intent: Option<crate::message::ToolCallIntent>,
179    pub tier: Tier,
180    pub risks: BTreeSet<RiskKind>,
181    pub args_digest: String,
182    pub preview: Option<String>,
183    pub provenance: ResourceProvenance,
184}
185
186impl PermissionIntent {
187    pub fn minimal(
188        tool_use_id: impl Into<String>,
189        tool_name: impl Into<String>,
190        tier: Tier,
191    ) -> Self {
192        Self {
193            tool_use_id: tool_use_id.into(),
194            tool_name: tool_name.into(),
195            call_intent: None,
196            tier,
197            risks: BTreeSet::new(),
198            args_digest: String::new(),
199            preview: None,
200            provenance: ResourceProvenance::none(),
201        }
202    }
203}
204
205#[derive(Debug, Clone, PartialEq, Eq)]
206pub struct AuthorityRequirement {
207    pub tier: Tier,
208    pub risks: BTreeSet<RiskKind>,
209    pub shell: bool,
210}
211
212impl AuthorityRequirement {
213    pub fn from_intent(intent: &PermissionIntent, shell: bool) -> Self {
214        Self {
215            tier: intent.tier,
216            risks: intent.risks.clone(),
217            shell,
218        }
219    }
220}
221
222#[derive(Debug, Clone, PartialEq, Eq)]
223pub enum ApprovalTarget {
224    Flow(FlowRunId),
225    User,
226}
227
228#[derive(Debug, Clone, PartialEq, Eq)]
229pub enum PermissionRequestState {
230    Evaluating,
231    Pending { target: ApprovalTarget },
232    Approved { decision_id: PermissionDecisionId },
233    Denied { decision_id: PermissionDecisionId },
234    Cancelled { reason: String },
235}
236
237impl PermissionRequestState {
238    fn is_terminal(&self) -> bool {
239        matches!(
240            self,
241            Self::Approved { .. } | Self::Denied { .. } | Self::Cancelled { .. }
242        )
243    }
244}
245
246#[derive(Debug, Clone, PartialEq, Eq)]
247pub enum DecisionActor {
248    Policy {
249        policy_version: String,
250        rule_id: String,
251    },
252    Flow {
253        session_id: String,
254        run_id: FlowRunId,
255    },
256    User {
257        session_id: String,
258        principal_id: Option<String>,
259    },
260    System {
261        component: String,
262    },
263}
264
265#[derive(Debug, Clone, Copy, PartialEq, Eq)]
266pub enum PermissionAction {
267    Approve,
268    Deny,
269    Defer,
270}
271
272#[derive(Debug, Clone, PartialEq, Eq)]
273pub enum GrantScope {
274    CurrentCall,
275    ChildRunSameTool {
276        run_id: FlowRunId,
277        tool_name: String,
278    },
279    ChildRunSamePathRule {
280        run_id: FlowRunId,
281        tool_name: String,
282        workspace_relative_path: String,
283    },
284}
285
286#[derive(Debug, Clone, PartialEq, Eq)]
287pub struct EscalationHop {
288    pub target: ApprovalTarget,
289    pub actor: Option<DecisionActor>,
290    pub action: Option<PermissionAction>,
291    pub reason: Option<String>,
292    pub at: DateTime<Utc>,
293}
294
295#[derive(Debug, Clone, PartialEq, Eq)]
296pub struct PermissionRequest {
297    pub request_id: PermissionRequestId,
298    pub session_id: String,
299    pub requesting_run_id: FlowRunId,
300    pub parent_run_id: Option<FlowRunId>,
301    pub root_run_id: FlowRunId,
302    pub intent: PermissionIntent,
303    pub requirement: AuthorityRequirement,
304    pub policy_reference: crate::permission_audit::PermissionPolicyReference,
305    pub state: PermissionRequestState,
306    /// Optimistic concurrency token for client decisions. Incremented for every
307    /// state or target transition while the stable request ID is unchanged.
308    pub revision: u64,
309    pub escalation_path: Vec<EscalationHop>,
310    pub requested_at: DateTime<Utc>,
311}
312
313#[derive(Debug, Clone, PartialEq, Eq)]
314pub struct PermissionDecision {
315    pub decision_id: PermissionDecisionId,
316    pub request_id: PermissionRequestId,
317    pub actor: DecisionActor,
318    pub action: PermissionAction,
319    pub execution_boundary: ExecutionBoundary,
320    pub grant_scope: Option<GrantScope>,
321    pub reason: Option<String>,
322    pub decided_at: DateTime<Utc>,
323}
324
325#[derive(Debug, Clone, PartialEq, Eq)]
326pub struct PermissionGrant {
327    pub grant_id: PermissionGrantId,
328    pub request_id: PermissionRequestId,
329    pub session_id: String,
330    pub requesting_run_id: FlowRunId,
331    pub requirement: AuthorityRequirement,
332    pub execution_boundary: ExecutionBoundary,
333    pub scope: GrantScope,
334    pub actor: DecisionActor,
335    pub granted_at: DateTime<Utc>,
336}
337
338#[derive(
339    Debug,
340    Clone,
341    Copy,
342    Default,
343    PartialEq,
344    Eq,
345    PartialOrd,
346    Ord,
347    serde::Serialize,
348    serde::Deserialize,
349)]
350#[serde(rename_all = "snake_case")]
351pub enum ExecutionBoundary {
352    #[default]
353    Sandboxed,
354    Direct,
355}
356
357#[derive(Debug, Clone, PartialEq, Eq)]
358pub enum ImmediateAuthorization {
359    Unrestricted,
360    Auto {
361        execution_boundary: ExecutionBoundary,
362    },
363    Denied {
364        reason: String,
365    },
366    Granted {
367        grant: Box<PermissionGrant>,
368    },
369}
370
371#[derive(Debug)]
372pub struct ImmediateSubmission {
373    pub request: PermissionRequest,
374    pub authorization: ImmediateAuthorization,
375}
376
377#[derive(Debug)]
378pub enum SubmissionOutcome {
379    Immediate(Box<ImmediateSubmission>),
380    Pending(Box<PendingPermission>),
381}
382
383/// Unforgeable authorization for one invocation and its exact resources.
384#[derive(Debug, Clone, PartialEq, Eq)]
385pub struct InvocationAuthorization {
386    request_id: PermissionRequestId,
387    tool_use_id: String,
388    tool_name: String,
389    provenance: ResourceProvenance,
390    execution_boundary: ExecutionBoundary,
391}
392
393impl InvocationAuthorization {
394    pub(crate) fn new(
395        request_id: PermissionRequestId,
396        tool_use_id: impl Into<String>,
397        tool_name: impl Into<String>,
398        provenance: ResourceProvenance,
399        execution_boundary: ExecutionBoundary,
400    ) -> Self {
401        Self {
402            request_id,
403            tool_use_id: tool_use_id.into(),
404            tool_name: tool_name.into(),
405            provenance,
406            execution_boundary,
407        }
408    }
409
410    pub fn request_id(&self) -> &PermissionRequestId {
411        &self.request_id
412    }
413
414    pub fn tool_name(&self) -> &str {
415        &self.tool_name
416    }
417
418    pub fn tool_use_id(&self) -> &str {
419        &self.tool_use_id
420    }
421
422    pub fn execution_boundary(&self) -> ExecutionBoundary {
423        self.execution_boundary
424    }
425
426    pub(crate) fn provenance(&self) -> &ResourceProvenance {
427        &self.provenance
428    }
429
430    // Prefix matching would turn a directory permit into blanket subtree access.
431    pub fn covers(&self, operation: &str, target: &std::path::Path) -> bool {
432        if self.tool_name != operation {
433            return false;
434        }
435        let canonical = crate::fs_access::canonicalize_stable(target);
436        self.provenance
437            .authorized_targets()
438            .any(|authorized| canonical == authorized)
439    }
440
441    pub fn is_for_call(&self, tool_use_id: &str, tool_name: &str) -> bool {
442        self.tool_use_id == tool_use_id && self.tool_name == tool_name
443    }
444}
445
446#[derive(Debug, Clone, PartialEq, Eq)]
447pub enum PermissionResolution {
448    Decision(PermissionDecision),
449    Cancelled { reason: String },
450}
451
452#[derive(Debug)]
453pub struct PendingPermission {
454    pub request: PermissionRequest,
455    pub resolution: oneshot::Receiver<PermissionResolution>,
456    pub target_changes: watch::Receiver<ApprovalTarget>,
457}
458
459#[derive(Debug, Clone, PartialEq, Eq)]
460pub enum ResolveOutcome {
461    Resolved(PermissionDecision),
462    Deferred(PermissionDecision),
463}
464
465/// Stable projections used by batch permission controls. Selectors never carry
466/// authority; they are expanded and checked against the authenticated actor.
467#[derive(Debug, Clone, PartialEq, Eq)]
468pub enum PermissionSelector {
469    RequestIds(Vec<PermissionRequestId>),
470    Group(PermissionGroupId),
471    DescendantRun(FlowRunId),
472    ChildRun,
473    Tool(String),
474    Tier(Tier),
475    Risk(RiskKind),
476    PathPrefix(PathBuf),
477    Target(ApprovalTarget),
478}
479
480#[derive(Debug, Clone, Copy, PartialEq, Eq)]
481pub enum BatchMode {
482    BestEffort,
483    Atomic,
484}
485
486#[derive(Debug, Clone, PartialEq, Eq)]
487pub enum BatchRequestOutcome {
488    Approved(PermissionDecision),
489    Denied(PermissionDecision),
490    Deferred(PermissionDecision),
491    SkippedAlreadyResolved,
492    RejectedNotAncestor,
493    RejectedOverAuthority,
494    RejectedStale,
495    RejectedNotFound,
496    RejectedNotRunning,
497    RejectedPermissionManagementRequired,
498    RejectedUnsupportedGrantScope,
499    RejectedNoEscalationTarget,
500    Rejected(PermissionError),
501}
502
503#[derive(Debug, Clone, PartialEq, Eq)]
504pub struct BatchResolution {
505    pub request_id: PermissionRequestId,
506    pub outcome: BatchRequestOutcome,
507}
508
509struct PreparedResolution {
510    request: PermissionRequest,
511    actor: DecisionActor,
512    target: ApprovalTarget,
513    next_target: Option<ApprovalTarget>,
514}
515
516#[derive(Debug, Clone, PartialEq, Eq)]
517pub enum PermissionError {
518    MissingIdentity,
519    IdentityMismatch,
520    RequestNotFound,
521    AlreadyResolved,
522    ActorNotAuthorized,
523    ActorNotRunning,
524    PermissionManagementRequired,
525    GrantExceedsAuthority,
526    UnsupportedGrantScope,
527    NoEscalationTarget,
528    GroupNotFound,
529    GroupNotEmpty,
530    GroupNotOwner,
531    GroupRevisionConflict,
532    EmptyGroup,
533}
534
535impl std::fmt::Display for PermissionError {
536    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
537        let message = match self {
538            Self::MissingIdentity => "permission identity is missing",
539            Self::IdentityMismatch => "permission identity does not match the authority graph",
540            Self::RequestNotFound => "permission request was not found",
541            Self::AlreadyResolved => "permission request is already resolved",
542            Self::ActorNotAuthorized => "decision actor is not authorized for this request",
543            Self::ActorNotRunning => "decision actor is not running",
544            Self::PermissionManagementRequired => {
545                "decision actor lacks permission-management authority"
546            }
547            Self::GrantExceedsAuthority => "requested grant exceeds actor authority",
548            Self::UnsupportedGrantScope => "path-rule grants require structured provenance",
549            Self::NoEscalationTarget => "permission request has no remaining escalation target",
550            Self::GroupNotFound => "permission group was not found",
551            Self::GroupNotEmpty => "permission group must be empty before deletion",
552            Self::GroupNotOwner => "permission group is not owned by this actor",
553            Self::GroupRevisionConflict => "permission group revision is stale",
554            Self::EmptyGroup => "permission group cannot be empty",
555        };
556        f.write_str(message)
557    }
558}
559
560impl std::error::Error for PermissionError {}
561
562struct RequestEntry {
563    request: PermissionRequest,
564    responder: Option<oneshot::Sender<PermissionResolution>>,
565    target_tx: Option<watch::Sender<ApprovalTarget>>,
566}
567
568struct SubmissionContext {
569    target_authority: Option<ApprovalAuthority>,
570}
571
572#[derive(Default)]
573struct BrokerState {
574    requests: HashMap<PermissionRequestId, RequestEntry>,
575    recent_terminal_requests: VecDeque<PermissionRequest>,
576    grants: Vec<PermissionGrant>,
577    groups: HashMap<PermissionGroupId, PermissionGroup>,
578    resolved_group_audits: BTreeSet<PermissionGroupId>,
579    next_audit_bundle: u64,
580    pending_audit_bundles: BTreeMap<u64, Vec<crate::permission_audit::PermissionAuditRecord>>,
581}
582
583impl BrokerState {
584    fn request(&self, request_id: &PermissionRequestId) -> Option<&PermissionRequest> {
585        self.requests
586            .get(request_id)
587            .map(|entry| &entry.request)
588            .or_else(|| {
589                self.recent_terminal_requests
590                    .iter()
591                    .find(|request| request.request_id == *request_id)
592            })
593    }
594
595    fn is_recent_terminal(&self, request_id: &PermissionRequestId) -> bool {
596        self.recent_terminal_requests
597            .iter()
598            .any(|request| request.request_id == *request_id)
599    }
600
601    fn remember_terminal(&mut self, request: PermissionRequest) {
602        debug_assert!(request.state.is_terminal());
603        if let Some(position) = self
604            .recent_terminal_requests
605            .iter()
606            .position(|existing| existing.request_id == request.request_id)
607        {
608            self.recent_terminal_requests.remove(position);
609        }
610        if self.recent_terminal_requests.len() == RECENT_TERMINAL_REQUEST_LIMIT {
611            self.recent_terminal_requests.pop_front();
612        }
613        self.recent_terminal_requests.push_back(request);
614    }
615
616    fn archive_terminal(&mut self, request_id: &PermissionRequestId) {
617        let Some(entry) = self.requests.get(request_id) else {
618            return;
619        };
620        if !entry.request.state.is_terminal() {
621            return;
622        }
623        let entry = self
624            .requests
625            .remove(request_id)
626            .expect("terminal request exists");
627        self.remember_terminal(entry.request);
628    }
629
630    fn archive_terminals<'a>(
631        &mut self,
632        request_ids: impl IntoIterator<Item = &'a PermissionRequestId>,
633    ) {
634        for request_id in request_ids {
635            self.archive_terminal(request_id);
636        }
637    }
638}
639
640#[derive(Default)]
641struct AuditDispatcher {
642    next_bundle: u64,
643}
644
645#[cfg(test)]
646type AuditDispatchHook = Arc<dyn Fn(u64) + Send + Sync>;
647
648#[derive(Debug, Clone)]
649pub struct FlowDecisionAuthority {
650    identity: Arc<FlowIdentity>,
651}
652
653#[derive(Debug, Clone)]
654pub struct UserDecisionAuthority {
655    broker_id: Uuid,
656    session_id: String,
657    principal_id: Option<String>,
658}
659
660#[derive(Debug, Clone)]
661pub enum DecisionAuthority {
662    Flow(FlowDecisionAuthority),
663    User(UserDecisionAuthority),
664}
665
666#[derive(Debug, Clone)]
667pub enum ApprovalAuthority {
668    Flow(FlowDecisionAuthority),
669    User(UserDecisionAuthority),
670}
671
672pub struct PermissionBroker {
673    broker_id: Uuid,
674    flows: Arc<FlowRegistry>,
675    state: Mutex<BrokerState>,
676    permission_clients: std::sync::atomic::AtomicUsize,
677    audit: Mutex<Option<crate::permission_audit::PermissionAuditProjector>>,
678    audit_dispatcher: Mutex<AuditDispatcher>,
679    #[cfg(test)]
680    audit_dispatch_hook: Mutex<Option<AuditDispatchHook>>,
681    /// Keeps the registry-observed cleanup hook alive for exactly this broker's lifetime.
682    observer: Mutex<Option<Arc<dyn crate::tools::agent_ctrl::FlowTerminalObserver>>>,
683}
684
685/// Bridges registry terminal transitions into the broker without keeping the broker
686/// alive: the registry stores only a `Weak`, so a dropped broker deregisters itself.
687struct TerminalCleanup {
688    broker: std::sync::Weak<PermissionBroker>,
689}
690
691pub struct PermissionClientGuard {
692    broker: std::sync::Weak<PermissionBroker>,
693}
694
695impl Drop for PermissionClientGuard {
696    fn drop(&mut self) {
697        if let Some(broker) = self.broker.upgrade()
698            && broker
699                .permission_clients
700                .fetch_sub(1, std::sync::atomic::Ordering::AcqRel)
701                == 1
702        {
703            broker.cancel_user_pending("permission client disconnected");
704        }
705    }
706}
707
708impl crate::tools::agent_ctrl::FlowTerminalObserver for TerminalCleanup {
709    fn flow_became_terminal(
710        &self,
711        session_id: &str,
712        run_id: &FlowRunId,
713    ) -> Option<Box<dyn FnOnce() + Send>> {
714        let broker = self.broker.upgrade()?;
715        broker.handle_terminal_locked(session_id, run_id);
716        Some(Box::new(move || broker.dispatch_audit_bundles()))
717    }
718}
719
720impl PermissionBroker {
721    /// Creates a broker without terminal-driven cleanup. Deliberately private: only
722    /// [`PermissionBroker::shared`] returns a broker in production, because a
723    /// non-observing broker never wakes pending awaiters on terminal transitions
724    /// unless some caller drives an `expire_terminal` sweep.
725    fn new(flows: Arc<FlowRegistry>) -> Self {
726        Self {
727            broker_id: Uuid::now_v7(),
728            flows,
729            state: Mutex::new(BrokerState::default()),
730            permission_clients: std::sync::atomic::AtomicUsize::new(0),
731            audit: Mutex::new(None),
732            audit_dispatcher: Mutex::new(AuditDispatcher::default()),
733            #[cfg(test)]
734            audit_dispatch_hook: Mutex::new(None),
735            observer: Mutex::new(None),
736        }
737    }
738
739    pub fn set_audit_projector(
740        &self,
741        projector: crate::permission_audit::PermissionAuditProjector,
742    ) {
743        self.audit.lock().unwrap().replace(projector);
744    }
745
746    pub fn register_client(self: &Arc<Self>) -> PermissionClientGuard {
747        self.permission_clients
748            .fetch_add(1, std::sync::atomic::Ordering::AcqRel);
749        PermissionClientGuard {
750            broker: Arc::downgrade(self),
751        }
752    }
753
754    pub fn has_clients(&self) -> bool {
755        self.permission_clients
756            .load(std::sync::atomic::Ordering::Acquire)
757            > 0
758    }
759
760    fn group_is_resolved(state: &BrokerState, group: &PermissionGroup) -> bool {
761        group.request_ids.iter().all(|request_id| {
762            state.requests.get(request_id).is_none_or(|entry| {
763                !matches!(entry.request.state, PermissionRequestState::Pending { .. })
764            })
765        })
766    }
767
768    fn unresolved_groups_for_requests(
769        state: &BrokerState,
770        request_ids: &BTreeSet<PermissionRequestId>,
771    ) -> BTreeSet<PermissionGroupId> {
772        state
773            .groups
774            .values()
775            .filter(|group| {
776                !group.request_ids.is_disjoint(request_ids)
777                    && !Self::group_is_resolved(state, group)
778            })
779            .map(|group| group.group_id.clone())
780            .collect()
781    }
782
783    fn append_resolved_group_audits(
784        state: &mut BrokerState,
785        candidate_group_ids: &BTreeSet<PermissionGroupId>,
786        records: &mut Vec<crate::permission_audit::PermissionAuditRecord>,
787        at: DateTime<Utc>,
788    ) {
789        for group_id in candidate_group_ids {
790            if state.resolved_group_audits.contains(group_id) {
791                continue;
792            }
793            let Some(group) = state.groups.get(group_id) else {
794                continue;
795            };
796            if !Self::group_is_resolved(state, group) {
797                continue;
798            }
799            let session_id = group
800                .request_ids
801                .iter()
802                .find_map(|request_id| {
803                    state
804                        .request(request_id)
805                        .map(|request| request.session_id.as_str())
806                })
807                .unwrap_or("unknown");
808            let audit =
809                crate::permission_audit::PermissionGroupAudit::from_group(group, session_id, at);
810            records.push(crate::permission_audit::PermissionAuditRecord::GroupResolved(audit));
811            state.resolved_group_audits.insert(group_id.clone());
812        }
813    }
814
815    fn enqueue_audit_bundle(
816        state: &mut BrokerState,
817        records: Vec<crate::permission_audit::PermissionAuditRecord>,
818    ) -> u64 {
819        let sequence = state.next_audit_bundle;
820        state.next_audit_bundle += 1;
821        assert!(
822            state
823                .pending_audit_bundles
824                .insert(sequence, records)
825                .is_none(),
826            "audit bundle sequence is unique"
827        );
828        sequence
829    }
830
831    fn dispatch_audit_bundles(&self) {
832        let mut dispatcher = self.audit_dispatcher.lock().unwrap();
833        loop {
834            let records = self
835                .state
836                .lock()
837                .unwrap()
838                .pending_audit_bundles
839                .remove(&dispatcher.next_bundle);
840            let Some(records) = records else {
841                break;
842            };
843            dispatcher.next_bundle += 1;
844            let projector = self.audit.lock().unwrap().clone();
845            if let Some(projector) = projector {
846                for record in records {
847                    projector.emit(record);
848                }
849            }
850        }
851    }
852
853    fn dispatch_audit_bundle(&self, _sequence: u64) {
854        #[cfg(test)]
855        let hook = self.audit_dispatch_hook.lock().unwrap().clone();
856        #[cfg(test)]
857        if let Some(hook) = hook {
858            hook(_sequence);
859        }
860        self.dispatch_audit_bundles();
861    }
862
863    #[cfg(test)]
864    fn set_audit_dispatch_hook(&self, hook: AuditDispatchHook) {
865        self.audit_dispatch_hook.lock().unwrap().replace(hook);
866    }
867
868    /// Wires terminal-driven cleanup. Only shared brokers can observe the registry,
869    /// because the observer must not keep the broker alive.
870    pub fn shared(flows: Arc<FlowRegistry>) -> Arc<Self> {
871        let broker = Arc::new(Self::new(Arc::clone(&flows)));
872        let observer: Arc<dyn crate::tools::agent_ctrl::FlowTerminalObserver> =
873            Arc::new(TerminalCleanup {
874                broker: Arc::downgrade(&broker),
875            });
876        // The registry holds a Weak to the trait object, so the Arc must live in the
877        // broker itself to stay upgradable for the broker's lifetime.
878        flows.register_terminal_observer(Arc::downgrade(&observer));
879        broker.observer.lock().unwrap().replace(observer);
880        broker
881    }
882
883    /// True when this broker arbitrates `flows`. Callers that already hold a
884    /// broker must confirm it is bound to the registry they are about to submit
885    /// against, otherwise identity authentication would run against a different
886    /// authority graph.
887    pub fn is_for_registry(&self, flows: &Arc<FlowRegistry>) -> bool {
888        Arc::ptr_eq(&self.flows, flows)
889    }
890
891    pub fn flow_authority(
892        &self,
893        identity: Arc<FlowIdentity>,
894    ) -> Result<FlowDecisionAuthority, PermissionError> {
895        self.authenticate_registered_identity(&identity)?;
896        Ok(FlowDecisionAuthority { identity })
897    }
898
899    pub fn user_authority(
900        &self,
901        session_id: impl Into<String>,
902        principal_id: Option<String>,
903    ) -> UserDecisionAuthority {
904        UserDecisionAuthority {
905            broker_id: self.broker_id,
906            session_id: session_id.into(),
907            principal_id,
908        }
909    }
910
911    pub fn submit(
912        &self,
913        session_id: Option<&str>,
914        run_id: Option<&FlowRunId>,
915        intent: PermissionIntent,
916        shell: bool,
917        policy: &TrustConfig,
918    ) -> Result<SubmissionOutcome, PermissionError> {
919        self.submit_to_with_context(
920            session_id,
921            run_id,
922            intent,
923            shell,
924            SubmissionContext {
925                target_authority: None,
926            },
927            policy,
928        )
929    }
930
931    pub fn submit_to(
932        &self,
933        session_id: Option<&str>,
934        run_id: Option<&FlowRunId>,
935        intent: PermissionIntent,
936        shell: bool,
937        target_authority: ApprovalAuthority,
938        policy: &TrustConfig,
939    ) -> Result<SubmissionOutcome, PermissionError> {
940        self.submit_to_with_context(
941            session_id,
942            run_id,
943            intent,
944            shell,
945            SubmissionContext {
946                target_authority: Some(target_authority),
947            },
948            policy,
949        )
950    }
951
952    fn submit_to_with_context(
953        &self,
954        session_id: Option<&str>,
955        run_id: Option<&FlowRunId>,
956        intent: PermissionIntent,
957        shell: bool,
958        context: SubmissionContext,
959        policy: &TrustConfig,
960    ) -> Result<SubmissionOutcome, PermissionError> {
961        // Liveness checks and the pending insert must observe the same lifecycle
962        // snapshot, otherwise a run could go terminal between them and leave an
963        // orphan pending request that terminal cleanup already walked past.
964        let (outcome, sequence) = self.flows.with_lifecycle_arbitration(|| {
965            self.submit_to_locked(session_id, run_id, intent, shell, context, policy)
966        })?;
967        self.dispatch_audit_bundle(sequence);
968        Ok(outcome)
969    }
970
971    fn submit_to_locked(
972        &self,
973        session_id: Option<&str>,
974        run_id: Option<&FlowRunId>,
975        intent: PermissionIntent,
976        shell: bool,
977        context: SubmissionContext,
978        policy: &TrustConfig,
979    ) -> Result<(SubmissionOutcome, u64), PermissionError> {
980        let requirement = AuthorityRequirement::from_intent(&intent, shell);
981        let identity =
982            self.authenticate_requester(session_id, run_id, &requirement, &intent.provenance)?;
983        let explicit_target = context
984            .target_authority
985            .map(|target| {
986                self.authenticate_target(&identity, &requirement, &intent.provenance, target)
987            })
988            .transpose()?;
989        let (execution_policy, resolution) = identity
990            .effective_authority
991            .constrain_policy_resolution(policy, intent.tier, intent.risks.iter().copied());
992        let mut state = self.state.lock().unwrap();
993        let request_id = PermissionRequestId::now();
994        let now = Utc::now();
995        let immediate = if execution_policy == ExecutionPolicy::Unrestricted {
996            Some(ImmediateAuthorization::Unrestricted)
997        } else {
998            match resolution.action {
999                PolicyAction::Auto => Some(ImmediateAuthorization::Auto {
1000                    execution_boundary: if intent.risks.contains(&RiskKind::ProcessSpawn)
1001                        && resolution.escalation == crate::trust::PolicyEscalation::Allowed
1002                    {
1003                        ExecutionBoundary::Direct
1004                    } else {
1005                        ExecutionBoundary::Sandboxed
1006                    },
1007                }),
1008                PolicyAction::Deny => Some(ImmediateAuthorization::Denied {
1009                    reason: "permission policy denied the invocation".into(),
1010                }),
1011                PolicyAction::Ask => {
1012                    matching_grant(&state.grants, &identity, &intent, &requirement).map(|grant| {
1013                        ImmediateAuthorization::Granted {
1014                            grant: Box::new(grant),
1015                        }
1016                    })
1017                }
1018            }
1019        };
1020        let target = explicit_target.unwrap_or_else(|| {
1021            if immediate.is_none() {
1022                self.next_eligible_target(&identity, &requirement, &intent.provenance, None)
1023            } else {
1024                ApprovalTarget::User
1025            }
1026        });
1027        let request_state = match &immediate {
1028            Some(ImmediateAuthorization::Denied { .. }) => PermissionRequestState::Denied {
1029                decision_id: PermissionDecisionId::now(),
1030            },
1031            Some(_) => PermissionRequestState::Approved {
1032                decision_id: PermissionDecisionId::now(),
1033            },
1034            None => PermissionRequestState::Pending {
1035                target: target.clone(),
1036            },
1037        };
1038        let request = PermissionRequest {
1039            request_id: request_id.clone(),
1040            session_id: identity.session_id.clone(),
1041            requesting_run_id: identity.run_id.clone(),
1042            parent_run_id: identity.parent_run_id.clone(),
1043            root_run_id: identity.root_run_id.clone(),
1044            policy_reference: crate::permission_audit::PermissionPolicyReference::capture(
1045                policy,
1046                intent.tier,
1047                &intent.risks,
1048                execution_policy,
1049                resolution,
1050            ),
1051            intent,
1052            requirement,
1053            state: request_state,
1054            revision: 0,
1055            escalation_path: vec![EscalationHop {
1056                target,
1057                actor: None,
1058                action: None,
1059                reason: None,
1060                at: now,
1061            }],
1062            requested_at: now,
1063        };
1064        let audits = submission_audits(&request, immediate.as_ref(), now);
1065        if let Some(authorization) = immediate {
1066            state.remember_terminal(request.clone());
1067            let sequence = Self::enqueue_audit_bundle(&mut state, audits);
1068            return Ok((
1069                SubmissionOutcome::Immediate(Box::new(ImmediateSubmission {
1070                    request,
1071                    authorization,
1072                })),
1073                sequence,
1074            ));
1075        }
1076
1077        let (responder, resolution) = oneshot::channel();
1078        let (target_tx, target_changes) = watch::channel(match &request.state {
1079            PermissionRequestState::Pending { target } => target.clone(),
1080            _ => unreachable!("pending submission must have a target"),
1081        });
1082        state.requests.insert(
1083            request_id,
1084            RequestEntry {
1085                request: request.clone(),
1086                responder: Some(responder),
1087                target_tx: Some(target_tx),
1088            },
1089        );
1090        let sequence = Self::enqueue_audit_bundle(&mut state, audits);
1091        Ok((
1092            SubmissionOutcome::Pending(Box::new(PendingPermission {
1093                request,
1094                resolution,
1095                target_changes,
1096            })),
1097            sequence,
1098        ))
1099    }
1100
1101    pub fn resolve(
1102        &self,
1103        request_id: &PermissionRequestId,
1104        authority: &DecisionAuthority,
1105        action: PermissionAction,
1106        grant_scope: Option<GrantScope>,
1107        reason: Option<String>,
1108    ) -> Result<ResolveOutcome, PermissionError> {
1109        // Holding the lifecycle arbitration across the liveness check and the decision
1110        // commit makes terminal-vs-approve a single linearization point: either this
1111        // decision commits before the run is terminal, or terminal cleanup runs first
1112        // and this call observes a cancelled request.
1113        let (outcome, sequence) = self.flows.with_lifecycle_arbitration(|| {
1114            self.resolve_locked(request_id, authority, action, grant_scope, reason)
1115        })?;
1116        self.dispatch_audit_bundle(sequence);
1117        Ok(outcome)
1118    }
1119
1120    fn resolve_locked(
1121        &self,
1122        request_id: &PermissionRequestId,
1123        authority: &DecisionAuthority,
1124        action: PermissionAction,
1125        grant_scope: Option<GrantScope>,
1126        reason: Option<String>,
1127    ) -> Result<(ResolveOutcome, u64), PermissionError> {
1128        let mut state = self.state.lock().unwrap();
1129        let request_ids = BTreeSet::from([request_id.clone()]);
1130        let candidate_groups = Self::unresolved_groups_for_requests(&state, &request_ids);
1131        let outcome = self.resolve_locked_state(
1132            &mut state,
1133            request_id,
1134            authority,
1135            action,
1136            grant_scope,
1137            reason,
1138        )?;
1139        let decision = match &outcome {
1140            ResolveOutcome::Resolved(decision) | ResolveOutcome::Deferred(decision) => decision,
1141        };
1142        let mut audits = resolution_audits_from_state(&state, request_id, decision);
1143        Self::append_resolved_group_audits(
1144            &mut state,
1145            &candidate_groups,
1146            &mut audits,
1147            decision.decided_at,
1148        );
1149        state.archive_terminal(request_id);
1150        let sequence = Self::enqueue_audit_bundle(&mut state, audits);
1151        Ok((outcome, sequence))
1152    }
1153
1154    fn prepare_resolution(
1155        &self,
1156        state: &BrokerState,
1157        request_id: &PermissionRequestId,
1158        authority: &DecisionAuthority,
1159        action: PermissionAction,
1160        grant_scope: Option<&GrantScope>,
1161    ) -> Result<PreparedResolution, PermissionError> {
1162        let request = match state.requests.get(request_id) {
1163            Some(entry) => entry.request.clone(),
1164            None if state.is_recent_terminal(request_id) => {
1165                return Err(PermissionError::AlreadyResolved);
1166            }
1167            None => return Err(PermissionError::RequestNotFound),
1168        };
1169        if request.state.is_terminal() {
1170            return Err(PermissionError::AlreadyResolved);
1171        }
1172        let requester = self
1173            .flows
1174            .lookup_run(&request.requesting_run_id)
1175            .filter(|identity| !matches!(identity.execution_state(), FlowExecutionState::Terminal))
1176            .ok_or(PermissionError::ActorNotRunning)?;
1177        let actor = self.validate_authority(&request, authority, action, grant_scope)?;
1178        let target = match &request.state {
1179            PermissionRequestState::Pending { target } => target.clone(),
1180            PermissionRequestState::Evaluating => {
1181                return Err(PermissionError::ActorNotAuthorized);
1182            }
1183            _ => return Err(PermissionError::AlreadyResolved),
1184        };
1185        if !actor_matches_target(&actor, &target) {
1186            return Err(PermissionError::ActorNotAuthorized);
1187        }
1188        let next_target = (action == PermissionAction::Defer
1189            && !matches!(target, ApprovalTarget::User))
1190        .then(|| {
1191            self.next_eligible_target(
1192                &requester,
1193                &request.requirement,
1194                &request.intent.provenance,
1195                match &target {
1196                    ApprovalTarget::Flow(run_id) => Some(run_id),
1197                    ApprovalTarget::User => None,
1198                },
1199            )
1200        });
1201        Ok(PreparedResolution {
1202            request,
1203            actor,
1204            target,
1205            next_target,
1206        })
1207    }
1208
1209    fn commit_prepared_resolution(
1210        &self,
1211        state: &mut BrokerState,
1212        prepared: PreparedResolution,
1213        action: PermissionAction,
1214        grant_scope: Option<GrantScope>,
1215        reason: Option<String>,
1216    ) -> ResolveOutcome {
1217        let request_id = prepared.request.request_id.clone();
1218        let decision = PermissionDecision {
1219            decision_id: PermissionDecisionId::now(),
1220            request_id: request_id.clone(),
1221            actor: prepared.actor.clone(),
1222            action,
1223            execution_boundary: decision_execution_boundary(
1224                &prepared.request,
1225                &prepared.actor,
1226                action,
1227            ),
1228            grant_scope: grant_scope.clone(),
1229            reason: reason.clone(),
1230            decided_at: Utc::now(),
1231        };
1232        let entry = state
1233            .requests
1234            .get_mut(&request_id)
1235            .expect("prepared request exists");
1236        entry.request.escalation_path.push(EscalationHop {
1237            target: prepared.target.clone(),
1238            actor: Some(prepared.actor.clone()),
1239            action: Some(action),
1240            reason: reason.clone(),
1241            at: decision.decided_at,
1242        });
1243        if action == PermissionAction::Defer {
1244            if matches!(prepared.target, ApprovalTarget::User) {
1245                cancel_entry(
1246                    entry,
1247                    "permission.user_defer",
1248                    reason.unwrap_or_else(|| "user deferred without a fallback".into()),
1249                    decision.decided_at,
1250                );
1251            } else {
1252                let next_target = prepared
1253                    .next_target
1254                    .expect("flow-target defer has a prepared next target");
1255                entry.request.state = PermissionRequestState::Pending {
1256                    target: next_target.clone(),
1257                };
1258                entry.request.revision += 1;
1259                if let Some(target_tx) = &entry.target_tx {
1260                    let _ = target_tx.send(next_target.clone());
1261                }
1262                entry.request.escalation_path.push(EscalationHop {
1263                    target: next_target,
1264                    actor: None,
1265                    action: None,
1266                    reason: None,
1267                    at: Utc::now(),
1268                });
1269            }
1270            return ResolveOutcome::Deferred(decision);
1271        }
1272        let persistent_grant = if action == PermissionAction::Approve {
1273            grant_scope
1274                .clone()
1275                .filter(|scope| !matches!(scope, GrantScope::CurrentCall))
1276                .map(|scope| PermissionGrant {
1277                    grant_id: PermissionGrantId::now(),
1278                    request_id: request_id.clone(),
1279                    session_id: prepared.request.session_id,
1280                    requesting_run_id: prepared.request.requesting_run_id,
1281                    requirement: prepared.request.requirement,
1282                    execution_boundary: decision.execution_boundary,
1283                    scope,
1284                    actor: prepared.actor,
1285                    granted_at: decision.decided_at,
1286                })
1287        } else {
1288            None
1289        };
1290        let entry = state
1291            .requests
1292            .get_mut(&request_id)
1293            .expect("prepared request exists");
1294        entry.request.state = if action == PermissionAction::Approve {
1295            PermissionRequestState::Approved {
1296                decision_id: decision.decision_id.clone(),
1297            }
1298        } else {
1299            PermissionRequestState::Denied {
1300                decision_id: decision.decision_id.clone(),
1301            }
1302        };
1303        entry.request.revision += 1;
1304        let responder = entry.responder.take();
1305        if let Some(grant) = persistent_grant
1306            && !state
1307                .grants
1308                .iter()
1309                .any(|existing| equivalent_grant(existing, &grant))
1310        {
1311            state.grants.push(grant);
1312        }
1313        if let Some(responder) = responder {
1314            let _ = responder.send(PermissionResolution::Decision(decision.clone()));
1315        }
1316        ResolveOutcome::Resolved(decision)
1317    }
1318
1319    fn resolve_locked_state(
1320        &self,
1321        state: &mut BrokerState,
1322        request_id: &PermissionRequestId,
1323        authority: &DecisionAuthority,
1324        action: PermissionAction,
1325        grant_scope: Option<GrantScope>,
1326        reason: Option<String>,
1327    ) -> Result<ResolveOutcome, PermissionError> {
1328        let request = match state.requests.get(request_id) {
1329            Some(entry) => entry.request.clone(),
1330            None if state.is_recent_terminal(request_id) => {
1331                return Err(PermissionError::AlreadyResolved);
1332            }
1333            None => return Err(PermissionError::RequestNotFound),
1334        };
1335        if request.state.is_terminal() {
1336            return Err(PermissionError::AlreadyResolved);
1337        }
1338        if matches!(
1339            self.flows.execution_state(&request.requesting_run_id),
1340            None | Some(FlowExecutionState::Terminal)
1341        ) {
1342            let reason = "requesting flow is terminal".to_owned();
1343            cancel_entry(
1344                state.requests.get_mut(request_id).unwrap(),
1345                "permission.requester_terminal",
1346                reason,
1347                Utc::now(),
1348            );
1349            state.grants.retain(|grant| {
1350                grant.session_id != request.session_id
1351                    || grant.requesting_run_id != request.requesting_run_id
1352            });
1353            state.archive_terminal(request_id);
1354            return Err(PermissionError::ActorNotRunning);
1355        }
1356
1357        let actor = self.validate_authority(&request, authority, action, grant_scope.as_ref())?;
1358        let target = match &request.state {
1359            PermissionRequestState::Pending { target } => target.clone(),
1360            PermissionRequestState::Evaluating => return Err(PermissionError::ActorNotAuthorized),
1361            _ => return Err(PermissionError::AlreadyResolved),
1362        };
1363        if !actor_matches_target(&actor, &target) {
1364            return Err(PermissionError::ActorNotAuthorized);
1365        }
1366
1367        let decision = PermissionDecision {
1368            decision_id: PermissionDecisionId::now(),
1369            request_id: request_id.clone(),
1370            actor: actor.clone(),
1371            action,
1372            execution_boundary: decision_execution_boundary(&request, &actor, action),
1373            grant_scope: grant_scope.clone(),
1374            reason: reason.clone(),
1375            decided_at: Utc::now(),
1376        };
1377        let entry = state.requests.get_mut(request_id).unwrap();
1378        entry.request.escalation_path.push(EscalationHop {
1379            target: target.clone(),
1380            actor: Some(actor.clone()),
1381            action: Some(action),
1382            reason: reason.clone(),
1383            at: decision.decided_at,
1384        });
1385
1386        if action == PermissionAction::Defer {
1387            if matches!(
1388                entry.request.state,
1389                PermissionRequestState::Pending {
1390                    target: ApprovalTarget::User
1391                }
1392            ) {
1393                let reason = reason.unwrap_or_else(|| "user deferred without a fallback".into());
1394                cancel_entry(entry, "permission.user_defer", reason, decision.decided_at);
1395                return Ok(ResolveOutcome::Deferred(decision));
1396            }
1397            let requester = self
1398                .flows
1399                .lookup_run(&entry.request.requesting_run_id)
1400                .ok_or(PermissionError::RequestNotFound)?;
1401            let next_target = self.next_eligible_target(
1402                &requester,
1403                &entry.request.requirement,
1404                &entry.request.intent.provenance,
1405                match &target {
1406                    ApprovalTarget::Flow(run_id) => Some(run_id),
1407                    ApprovalTarget::User => None,
1408                },
1409            );
1410            entry.request.state = PermissionRequestState::Pending {
1411                target: next_target.clone(),
1412            };
1413            if let Some(target_tx) = &entry.target_tx {
1414                let _ = target_tx.send(next_target.clone());
1415            }
1416            entry.request.escalation_path.push(EscalationHop {
1417                target: next_target,
1418                actor: None,
1419                action: None,
1420                reason: None,
1421                at: Utc::now(),
1422            });
1423            return Ok(ResolveOutcome::Deferred(decision));
1424        }
1425
1426        let persistent_grant = if action == PermissionAction::Approve {
1427            grant_scope
1428                .clone()
1429                .filter(|scope| !matches!(scope, GrantScope::CurrentCall))
1430                .map(|scope| PermissionGrant {
1431                    grant_id: PermissionGrantId::now(),
1432                    request_id: request_id.clone(),
1433                    session_id: request.session_id.clone(),
1434                    requesting_run_id: request.requesting_run_id.clone(),
1435                    requirement: request.requirement.clone(),
1436                    execution_boundary: decision.execution_boundary,
1437                    scope,
1438                    actor,
1439                    granted_at: decision.decided_at,
1440                })
1441        } else {
1442            None
1443        };
1444        let entry = state.requests.get_mut(request_id).unwrap();
1445        entry.request.state = if action == PermissionAction::Approve {
1446            PermissionRequestState::Approved {
1447                decision_id: decision.decision_id.clone(),
1448            }
1449        } else {
1450            PermissionRequestState::Denied {
1451                decision_id: decision.decision_id.clone(),
1452            }
1453        };
1454        entry.request.revision += 1;
1455        let responder = entry.responder.take();
1456        if let Some(grant) = persistent_grant
1457            && !state
1458                .grants
1459                .iter()
1460                .any(|existing| equivalent_grant(existing, &grant))
1461        {
1462            state.grants.push(grant);
1463        }
1464        if let Some(responder) = responder {
1465            let _ = responder.send(PermissionResolution::Decision(decision.clone()));
1466        }
1467        Ok(ResolveOutcome::Resolved(decision))
1468    }
1469
1470    pub fn defer_timed_out_target(
1471        &self,
1472        request_id: &PermissionRequestId,
1473        expected_target: &FlowRunId,
1474    ) -> Result<bool, PermissionError> {
1475        let (deferred, sequence) = self.flows.with_lifecycle_arbitration(|| {
1476            let mut state = self.state.lock().unwrap();
1477            let (deferred, audits) = self.defer_unavailable_target_state(
1478                &mut state,
1479                request_id,
1480                expected_target,
1481                "permission.parent_timeout",
1482                "ancestor offer timed out",
1483            )?;
1484            Ok((deferred, Self::enqueue_audit_bundle(&mut state, audits)))
1485        })?;
1486        self.dispatch_audit_bundle(sequence);
1487        Ok(deferred)
1488    }
1489
1490    fn defer_unavailable_target_state(
1491        &self,
1492        state: &mut BrokerState,
1493        request_id: &PermissionRequestId,
1494        expected_target: &FlowRunId,
1495        component: &str,
1496        reason: &str,
1497    ) -> Result<(bool, Vec<crate::permission_audit::PermissionAuditRecord>), PermissionError> {
1498        if !state.requests.contains_key(request_id) {
1499            return if state.is_recent_terminal(request_id) {
1500                Ok((false, Vec::new()))
1501            } else {
1502                Err(PermissionError::RequestNotFound)
1503            };
1504        }
1505        let entry = state
1506            .requests
1507            .get_mut(request_id)
1508            .expect("active request exists");
1509        if !matches!(
1510            &entry.request.state,
1511            PermissionRequestState::Pending { target: ApprovalTarget::Flow(run_id) }
1512                if run_id == expected_target
1513        ) {
1514            return Ok((false, Vec::new()));
1515        }
1516        let requester = self
1517            .flows
1518            .lookup_run(&entry.request.requesting_run_id)
1519            .ok_or(PermissionError::RequestNotFound)?;
1520        let now = Utc::now();
1521        entry.request.escalation_path.push(EscalationHop {
1522            target: ApprovalTarget::Flow(expected_target.clone()),
1523            actor: Some(DecisionActor::System {
1524                component: component.into(),
1525            }),
1526            action: Some(PermissionAction::Defer),
1527            reason: Some(reason.into()),
1528            at: now,
1529        });
1530        let next_target = self.next_eligible_target(
1531            &requester,
1532            &entry.request.requirement,
1533            &entry.request.intent.provenance,
1534            Some(expected_target),
1535        );
1536        entry.request.state = PermissionRequestState::Pending {
1537            target: next_target.clone(),
1538        };
1539        entry.request.revision += 1;
1540        entry.request.escalation_path.push(EscalationHop {
1541            target: next_target.clone(),
1542            actor: None,
1543            action: None,
1544            reason: None,
1545            at: now,
1546        });
1547        if let Some(target_tx) = &entry.target_tx {
1548            let _ = target_tx.send(next_target);
1549        }
1550        let audit = request_audit_from_state(state, request_id, None, now)
1551            .expect("deferred request exists");
1552        Ok((
1553            true,
1554            vec![
1555                crate::permission_audit::PermissionAuditRecord::RequestDeferred(audit.clone()),
1556                crate::permission_audit::PermissionAuditRecord::RequestTargeted(audit),
1557            ],
1558        ))
1559    }
1560
1561    pub fn cancel(
1562        &self,
1563        request_id: &PermissionRequestId,
1564        reason: impl Into<String>,
1565    ) -> Result<(), PermissionError> {
1566        let sequence = {
1567            let mut state = self.state.lock().unwrap();
1568            let at = Utc::now();
1569            let request_ids = BTreeSet::from([request_id.clone()]);
1570            let candidate_groups = Self::unresolved_groups_for_requests(&state, &request_ids);
1571            if !state.requests.contains_key(request_id) {
1572                return if state.is_recent_terminal(request_id) {
1573                    Err(PermissionError::AlreadyResolved)
1574                } else {
1575                    Err(PermissionError::RequestNotFound)
1576                };
1577            }
1578            let entry = state
1579                .requests
1580                .get_mut(request_id)
1581                .expect("active request exists");
1582            if entry.request.state.is_terminal() || entry.responder.is_none() {
1583                return Err(PermissionError::AlreadyResolved);
1584            }
1585            cancel_entry(entry, "permission.user_cancel", reason.into(), at);
1586            let audit = request_audit_from_state(&state, request_id, None, at)
1587                .expect("cancelled request exists");
1588            let mut audits =
1589                vec![crate::permission_audit::PermissionAuditRecord::RequestCancelled(audit)];
1590            Self::append_resolved_group_audits(&mut state, &candidate_groups, &mut audits, at);
1591            state.archive_terminal(request_id);
1592            Self::enqueue_audit_bundle(&mut state, audits)
1593        };
1594        self.dispatch_audit_bundle(sequence);
1595        Ok(())
1596    }
1597
1598    fn cancel_user_pending(&self, reason: &str) {
1599        let sequence = {
1600            let mut state = self.state.lock().unwrap();
1601            let at = Utc::now();
1602            let request_ids: BTreeSet<_> = state
1603                .requests
1604                .iter()
1605                .filter(|(_, entry)| {
1606                    matches!(
1607                        entry.request.state,
1608                        PermissionRequestState::Pending {
1609                            target: ApprovalTarget::User
1610                        }
1611                    ) && entry.responder.is_some()
1612                })
1613                .map(|(request_id, _)| request_id.clone())
1614                .collect();
1615            if request_ids.is_empty() {
1616                return;
1617            }
1618            let candidate_groups = Self::unresolved_groups_for_requests(&state, &request_ids);
1619            for request_id in &request_ids {
1620                let entry = state
1621                    .requests
1622                    .get_mut(request_id)
1623                    .expect("pending permission exists");
1624                cancel_entry(
1625                    entry,
1626                    "permission.client_disconnect",
1627                    reason.to_string(),
1628                    at,
1629                );
1630            }
1631            let mut audits = request_ids
1632                .iter()
1633                .filter_map(|request_id| request_audit_from_state(&state, request_id, None, at))
1634                .map(crate::permission_audit::PermissionAuditRecord::RequestCancelled)
1635                .collect::<Vec<_>>();
1636            Self::append_resolved_group_audits(&mut state, &candidate_groups, &mut audits, at);
1637            state.archive_terminals(&request_ids);
1638            Self::enqueue_audit_bundle(&mut state, audits)
1639        };
1640        self.dispatch_audit_bundle(sequence);
1641    }
1642
1643    pub fn cancel_for_run(&self, session_id: &str, run_id: &FlowRunId, reason: &str) -> usize {
1644        let (cancelled, sequence) = self.flows.with_lifecycle_arbitration(|| {
1645            let mut state = self.state.lock().unwrap();
1646            let (cancelled, audits) =
1647                self.cancel_for_run_state(&mut state, session_id, run_id, reason);
1648            (cancelled, Self::enqueue_audit_bundle(&mut state, audits))
1649        });
1650        self.dispatch_audit_bundle(sequence);
1651        cancelled
1652    }
1653
1654    fn handle_terminal_locked(&self, session_id: &str, run_id: &FlowRunId) -> u64 {
1655        let mut state = self.state.lock().unwrap();
1656        let (_, mut audits) = self.cancel_for_run_state(
1657            &mut state,
1658            session_id,
1659            run_id,
1660            "requesting flow is terminal",
1661        );
1662        let request_ids: Vec<_> = state
1663            .requests
1664            .iter()
1665            .filter(|(_, entry)| {
1666                matches!(
1667                    &entry.request.state,
1668                    PermissionRequestState::Pending { target: ApprovalTarget::Flow(target) }
1669                        if target == run_id
1670                )
1671            })
1672            .map(|(request_id, _)| request_id.clone())
1673            .collect();
1674        for request_id in request_ids {
1675            if let Ok((_, deferred_audits)) = self.defer_unavailable_target_state(
1676                &mut state,
1677                &request_id,
1678                run_id,
1679                "permission.target_terminal",
1680                "target flow became terminal",
1681            ) {
1682                audits.extend(deferred_audits);
1683            }
1684        }
1685        Self::enqueue_audit_bundle(&mut state, audits)
1686    }
1687
1688    fn cancel_for_run_state(
1689        &self,
1690        state: &mut BrokerState,
1691        session_id: &str,
1692        run_id: &FlowRunId,
1693        reason: &str,
1694    ) -> (usize, Vec<crate::permission_audit::PermissionAuditRecord>) {
1695        let now = Utc::now();
1696        let request_ids: Vec<_> = state
1697            .requests
1698            .iter()
1699            .filter(|(_, entry)| {
1700                entry.request.session_id == session_id
1701                    && entry.request.requesting_run_id == *run_id
1702                    && !entry.request.state.is_terminal()
1703                    && entry.responder.is_some()
1704            })
1705            .map(|(request_id, _)| request_id.clone())
1706            .collect();
1707        let request_id_set = request_ids.iter().cloned().collect();
1708        let candidate_groups = Self::unresolved_groups_for_requests(state, &request_id_set);
1709        for request_id in &request_ids {
1710            cancel_entry(
1711                state.requests.get_mut(request_id).expect("request exists"),
1712                "permission.run_cleanup",
1713                reason.to_owned(),
1714                now,
1715            );
1716        }
1717        let audits = request_ids
1718            .iter()
1719            .map(|request_id| {
1720                request_audit_from_state(state, request_id, None, now)
1721                    .expect("cancelled request exists")
1722            })
1723            .collect::<Vec<_>>();
1724        let mut expired_grants = Vec::new();
1725        state.grants.retain(|grant| {
1726            let remove = grant.session_id == session_id && grant.requesting_run_id == *run_id;
1727            if remove {
1728                expired_grants.push(grant.clone());
1729            }
1730            !remove
1731        });
1732        let cancelled = audits.len();
1733        let mut records = audits
1734            .into_iter()
1735            .map(crate::permission_audit::PermissionAuditRecord::RequestCancelled)
1736            .collect::<Vec<_>>();
1737        records.extend(expired_grants.into_iter().map(|grant| {
1738            crate::permission_audit::PermissionAuditRecord::GrantExpired(
1739                crate::permission_audit::PermissionGrantAudit::from_grant_with_actor(
1740                    &grant,
1741                    crate::permission_audit::PermissionProjectionActor::System {
1742                        component: "permission.run_cleanup".into(),
1743                    },
1744                    Some(reason.to_owned()),
1745                    now,
1746                ),
1747            )
1748        }));
1749        Self::append_resolved_group_audits(state, &candidate_groups, &mut records, now);
1750        state.archive_terminals(&request_ids);
1751        (cancelled, records)
1752    }
1753
1754    pub fn expire_terminal(&self, reason: &str) -> usize {
1755        let (expired, sequence) = self.flows.with_lifecycle_arbitration(|| {
1756            let mut state = self.state.lock().unwrap();
1757            let terminal_runs: HashSet<_> = state
1758                .requests
1759                .values()
1760                .filter(|entry| {
1761                    !entry.request.state.is_terminal()
1762                        && matches!(
1763                            self.flows.execution_state(&entry.request.requesting_run_id),
1764                            None | Some(FlowExecutionState::Terminal)
1765                        )
1766                })
1767                .map(|entry| {
1768                    (
1769                        entry.request.session_id.clone(),
1770                        entry.request.requesting_run_id.clone(),
1771                    )
1772                })
1773                .collect();
1774            let mut expired = 0;
1775            let mut audits = Vec::new();
1776            for (session_id, run_id) in &terminal_runs {
1777                let (run_expired, run_audits) =
1778                    self.cancel_for_run_state(&mut state, session_id, run_id, reason);
1779                expired += run_expired;
1780                audits.extend(run_audits);
1781            }
1782            (expired, Self::enqueue_audit_bundle(&mut state, audits))
1783        });
1784        self.dispatch_audit_bundle(sequence);
1785        expired
1786    }
1787
1788    pub fn expire_before(&self, deadline: DateTime<Utc>, reason: &str) -> usize {
1789        let (expired, sequence) = {
1790            let mut state = self.state.lock().unwrap();
1791            let now = Utc::now();
1792            let request_ids: Vec<_> = state
1793                .requests
1794                .iter()
1795                .filter(|(_, entry)| {
1796                    entry.request.requested_at <= deadline
1797                        && !entry.request.state.is_terminal()
1798                        && entry.responder.is_some()
1799                })
1800                .map(|(request_id, _)| request_id.clone())
1801                .collect();
1802            let request_id_set = request_ids.iter().cloned().collect();
1803            let candidate_groups = Self::unresolved_groups_for_requests(&state, &request_id_set);
1804            for request_id in &request_ids {
1805                cancel_entry(
1806                    state.requests.get_mut(request_id).expect("request exists"),
1807                    "permission.expiry",
1808                    reason.to_owned(),
1809                    now,
1810                );
1811            }
1812            let mut audits = request_ids
1813                .iter()
1814                .filter_map(|request_id| request_audit_from_state(&state, request_id, None, now))
1815                .map(crate::permission_audit::PermissionAuditRecord::RequestCancelled)
1816                .collect::<Vec<_>>();
1817            let expired = audits.len();
1818            Self::append_resolved_group_audits(&mut state, &candidate_groups, &mut audits, now);
1819            state.archive_terminals(&request_ids);
1820            (expired, Self::enqueue_audit_bundle(&mut state, audits))
1821        };
1822        self.dispatch_audit_bundle(sequence);
1823        expired
1824    }
1825
1826    pub fn get(&self, request_id: &PermissionRequestId) -> Option<PermissionRequest> {
1827        self.state.lock().unwrap().request(request_id).cloned()
1828    }
1829
1830    pub fn user_list(&self, session_id: &str) -> (Vec<PermissionRequest>, Vec<PermissionGroup>) {
1831        let state = self.state.lock().unwrap();
1832        let requests: Vec<_> = state
1833            .requests
1834            .values()
1835            .filter(|entry| Self::user_visible_request(&entry.request, session_id))
1836            .map(|entry| entry.request.clone())
1837            .collect();
1838        let groups: Vec<_> = state
1839            .groups
1840            .values()
1841            .filter(|group| group.owner == GroupOwner::User)
1842            .filter(|group| {
1843                !group.request_ids.is_empty()
1844                    && group.request_ids.iter().all(|request_id| {
1845                        state.requests.get(request_id).is_some_and(|entry| {
1846                            Self::user_visible_request(&entry.request, session_id)
1847                        })
1848                    })
1849            })
1850            .cloned()
1851            .collect();
1852        (requests, groups)
1853    }
1854
1855    pub fn user_create_group(
1856        &self,
1857        session_id: &str,
1858        request_ids: BTreeSet<PermissionRequestId>,
1859        label: String,
1860        expected_revisions: &HashMap<PermissionRequestId, u64>,
1861    ) -> Result<PermissionGroup, PermissionError> {
1862        if request_ids.is_empty() {
1863            return Err(PermissionError::EmptyGroup);
1864        }
1865        let mut state = self.state.lock().unwrap();
1866        if request_ids.iter().any(|id| {
1867            state.requests.get(id).is_none_or(|entry| {
1868                !Self::user_visible_request(&entry.request, session_id)
1869                    || expected_revisions.get(id) != Some(&entry.request.revision)
1870            })
1871        }) {
1872            return Err(PermissionError::GroupRevisionConflict);
1873        }
1874        let member_revisions = request_ids
1875            .iter()
1876            .map(|id| {
1877                (
1878                    id.clone(),
1879                    state
1880                        .requests
1881                        .get(id)
1882                        .expect("validated group member")
1883                        .request
1884                        .revision,
1885                )
1886            })
1887            .collect();
1888        let group = PermissionGroup {
1889            group_id: PermissionGroupId::now(),
1890            owner: GroupOwner::User,
1891            label,
1892            request_ids,
1893            member_revisions,
1894            created_at: Utc::now(),
1895            revision: 0,
1896        };
1897        state.groups.insert(group.group_id.clone(), group.clone());
1898        let records = vec![
1899            crate::permission_audit::PermissionAuditRecord::GroupCreated(
1900                crate::permission_audit::PermissionGroupAudit::from_group(
1901                    &group,
1902                    session_id,
1903                    group.created_at,
1904                ),
1905            ),
1906        ];
1907        let sequence = Self::enqueue_audit_bundle(&mut state, records);
1908        drop(state);
1909        self.dispatch_audit_bundle(sequence);
1910        Ok(group)
1911    }
1912
1913    #[allow(clippy::too_many_arguments)]
1914    pub fn user_resolve(
1915        &self,
1916        session_id: &str,
1917        principal_id: Option<String>,
1918        request_ids: Vec<PermissionRequestId>,
1919        expected_revisions: &HashMap<PermissionRequestId, u64>,
1920        group_id: Option<(PermissionGroupId, u64)>,
1921        action: PermissionAction,
1922        grant_scope: Option<GrantScope>,
1923        reason: Option<String>,
1924    ) -> Result<Vec<BatchResolution>, PermissionError> {
1925        let (results, sequence) = self.flows.with_lifecycle_arbitration(|| {
1926            let mut state = self.state.lock().unwrap();
1927            let is_group_selector = group_id.is_some();
1928            let ids = if let Some((group_id, expected_revision)) = group_id {
1929                let group = state
1930                    .groups
1931                    .get(&group_id)
1932                    .ok_or(PermissionError::GroupNotFound)?;
1933                if group.owner != GroupOwner::User || group.revision != expected_revision {
1934                    return Err(PermissionError::GroupRevisionConflict);
1935                }
1936                if group.request_ids.iter().any(|id| {
1937                    state.requests.get(id).is_none_or(|entry| {
1938                        group.member_revisions.get(id) != Some(&entry.request.revision)
1939                    })
1940                }) {
1941                    return Err(PermissionError::GroupRevisionConflict);
1942                }
1943                group.request_ids.iter().cloned().collect()
1944            } else {
1945                request_ids
1946            };
1947            if ids.iter().any(|id| {
1948                state.requests.get(id).is_none_or(|entry| {
1949                    !Self::user_visible_request(&entry.request, session_id)
1950                        || (!is_group_selector
1951                            && expected_revisions.get(id) != Some(&entry.request.revision))
1952                })
1953            }) {
1954                return Err(PermissionError::GroupRevisionConflict);
1955            }
1956            let authority =
1957                DecisionAuthority::User(self.user_authority(session_id.to_string(), principal_id));
1958            let request_set = ids.iter().cloned().collect();
1959            let candidate_groups = Self::unresolved_groups_for_requests(&state, &request_set);
1960            let prepared: Vec<_> = ids
1961                .iter()
1962                .map(|id| {
1963                    self.prepare_resolution(&state, id, &authority, action, grant_scope.as_ref())
1964                })
1965                .collect::<Result<_, _>>()?;
1966            let mut results = Vec::with_capacity(prepared.len());
1967            let mut audits = Vec::new();
1968            for prepared in prepared {
1969                let request_id = prepared.request.request_id.clone();
1970                let outcome = self.commit_prepared_resolution(
1971                    &mut state,
1972                    prepared,
1973                    action,
1974                    grant_scope.clone(),
1975                    reason.clone(),
1976                );
1977                let decision = match &outcome {
1978                    ResolveOutcome::Resolved(decision) | ResolveOutcome::Deferred(decision) => {
1979                        decision
1980                    }
1981                };
1982                audits.extend(resolution_audits_from_state(&state, &request_id, decision));
1983                results.push(BatchResolution {
1984                    request_id,
1985                    outcome: batch_outcome(action, outcome),
1986                });
1987            }
1988            Self::append_resolved_group_audits(
1989                &mut state,
1990                &candidate_groups,
1991                &mut audits,
1992                Utc::now(),
1993            );
1994            let terminal_ids = results
1995                .iter()
1996                .map(|result| &result.request_id)
1997                .collect::<Vec<_>>();
1998            state.archive_terminals(terminal_ids);
1999            Ok((results, Self::enqueue_audit_bundle(&mut state, audits)))
2000        })?;
2001        self.dispatch_audit_bundle(sequence);
2002        Ok(results)
2003    }
2004
2005    pub fn list(&self) -> Vec<PermissionRequest> {
2006        let state = self.state.lock().unwrap();
2007        let mut requests: Vec<_> = state
2008            .requests
2009            .values()
2010            .map(|entry| entry.request.clone())
2011            .chain(state.recent_terminal_requests.iter().cloned())
2012            .collect();
2013        requests.sort_by_key(|request| request.requested_at);
2014        requests
2015    }
2016
2017    pub fn visible_get(
2018        &self,
2019        actor: &Arc<FlowIdentity>,
2020        request_id: &PermissionRequestId,
2021    ) -> Result<Option<PermissionRequest>, PermissionError> {
2022        self.authenticate_visibility_actor(actor)?;
2023        let request = self
2024            .state
2025            .lock()
2026            .unwrap()
2027            .requests
2028            .get(request_id)
2029            .map(|entry| entry.request.clone());
2030        Ok(request.filter(|request| self.visible_to(actor, request)))
2031    }
2032
2033    pub fn visible_list(
2034        &self,
2035        actor: &Arc<FlowIdentity>,
2036    ) -> Result<Vec<PermissionRequest>, PermissionError> {
2037        self.authenticate_visibility_actor(actor)?;
2038        let mut requests: Vec<_> = self
2039            .state
2040            .lock()
2041            .unwrap()
2042            .requests
2043            .values()
2044            .map(|entry| entry.request.clone())
2045            .filter(|request| self.visible_to(actor, request))
2046            .collect();
2047        requests.sort_by_key(|request| request.requested_at);
2048        Ok(requests)
2049    }
2050
2051    pub fn authenticate_control_actor(
2052        &self,
2053        actor: &Arc<FlowIdentity>,
2054    ) -> Result<(), PermissionError> {
2055        self.authenticate_visibility_actor(actor)
2056    }
2057
2058    pub fn create_group(
2059        &self,
2060        actor: &Arc<FlowIdentity>,
2061        request_ids: BTreeSet<PermissionRequestId>,
2062        label: String,
2063    ) -> Result<PermissionGroup, PermissionError> {
2064        self.authenticate_visibility_actor(actor)?;
2065        if request_ids.is_empty() {
2066            return Err(PermissionError::EmptyGroup);
2067        }
2068        let mut state = self.state.lock().unwrap();
2069        if request_ids.iter().any(|request_id| {
2070            state.request(request_id).is_none_or(|request| {
2071                !(self.visible_to(actor, request)
2072                    || request.state.is_terminal()
2073                        && request.session_id == actor.session_id
2074                        && self
2075                            .flows
2076                            .is_strict_ancestor(&actor.run_id, &request.requesting_run_id))
2077            })
2078        }) {
2079            return Err(PermissionError::ActorNotAuthorized);
2080        }
2081        let member_revisions = request_ids
2082            .iter()
2083            .filter_map(|id| {
2084                state
2085                    .request(id)
2086                    .map(|request| (id.clone(), request.revision))
2087            })
2088            .collect();
2089        let group = PermissionGroup {
2090            group_id: PermissionGroupId::now(),
2091            owner: GroupOwner::Flow(actor.run_id.clone()),
2092            label,
2093            request_ids,
2094            member_revisions,
2095            created_at: Utc::now(),
2096            revision: 0,
2097        };
2098        state.groups.insert(group.group_id.clone(), group.clone());
2099        let mut records = vec![
2100            crate::permission_audit::PermissionAuditRecord::GroupCreated(
2101                crate::permission_audit::PermissionGroupAudit::from_group(
2102                    &group,
2103                    &actor.session_id,
2104                    group.created_at,
2105                ),
2106            ),
2107        ];
2108        Self::append_resolved_group_audits(
2109            &mut state,
2110            &BTreeSet::from([group.group_id.clone()]),
2111            &mut records,
2112            group.created_at,
2113        );
2114        let sequence = Self::enqueue_audit_bundle(&mut state, records);
2115        drop(state);
2116        self.dispatch_audit_bundle(sequence);
2117        Ok(group)
2118    }
2119
2120    pub fn visible_group_get(
2121        &self,
2122        actor: &Arc<FlowIdentity>,
2123        group_id: &PermissionGroupId,
2124    ) -> Result<Option<PermissionGroup>, PermissionError> {
2125        self.authenticate_visibility_actor(actor)?;
2126        Ok(self
2127            .state
2128            .lock()
2129            .unwrap()
2130            .groups
2131            .get(group_id)
2132            .filter(|group| group.owner == GroupOwner::Flow(actor.run_id.clone()))
2133            .cloned())
2134    }
2135
2136    pub fn visible_group_list(
2137        &self,
2138        actor: &Arc<FlowIdentity>,
2139    ) -> Result<Vec<PermissionGroup>, PermissionError> {
2140        self.authenticate_visibility_actor(actor)?;
2141        let mut groups: Vec<_> = self
2142            .state
2143            .lock()
2144            .unwrap()
2145            .groups
2146            .values()
2147            .filter(|group| group.owner == GroupOwner::Flow(actor.run_id.clone()))
2148            .cloned()
2149            .collect();
2150        groups.sort_by_key(|group| group.created_at);
2151        Ok(groups)
2152    }
2153
2154    pub fn ungroup_requests(
2155        &self,
2156        actor: &Arc<FlowIdentity>,
2157        group_id: &PermissionGroupId,
2158        request_ids: &BTreeSet<PermissionRequestId>,
2159    ) -> Result<PermissionGroup, PermissionError> {
2160        self.authenticate_visibility_actor(actor)?;
2161        let mut state = self.state.lock().unwrap();
2162        let group = state
2163            .groups
2164            .get_mut(group_id)
2165            .ok_or(PermissionError::GroupNotFound)?;
2166        if group.owner != GroupOwner::Flow(actor.run_id.clone()) {
2167            return Err(PermissionError::GroupNotOwner);
2168        }
2169        let previous_len = group.request_ids.len();
2170        group.request_ids.retain(|id| !request_ids.contains(id));
2171        if group.request_ids.len() != previous_len {
2172            group.revision += 1;
2173        }
2174        let group = group.clone();
2175        let at = Utc::now();
2176        let mut records = vec![
2177            crate::permission_audit::PermissionAuditRecord::GroupUpdated(
2178                crate::permission_audit::PermissionGroupAudit::from_group(
2179                    &group,
2180                    &actor.session_id,
2181                    at,
2182                ),
2183            ),
2184        ];
2185        Self::append_resolved_group_audits(
2186            &mut state,
2187            &BTreeSet::from([group.group_id.clone()]),
2188            &mut records,
2189            at,
2190        );
2191        let sequence = Self::enqueue_audit_bundle(&mut state, records);
2192        drop(state);
2193        self.dispatch_audit_bundle(sequence);
2194        Ok(group)
2195    }
2196
2197    fn expand_selector(
2198        &self,
2199        state: &BrokerState,
2200        actor: &FlowIdentity,
2201        selector: &PermissionSelector,
2202        expected_group_revision: Option<u64>,
2203    ) -> Result<Vec<PermissionRequestId>, PermissionError> {
2204        let mut ids = BTreeSet::new();
2205        match selector {
2206            PermissionSelector::RequestIds(request_ids) => ids.extend(request_ids.iter().cloned()),
2207            PermissionSelector::Group(group_id) => {
2208                let group = state
2209                    .groups
2210                    .get(group_id)
2211                    .ok_or(PermissionError::GroupNotFound)?;
2212                if group.owner != GroupOwner::Flow(actor.run_id.clone()) {
2213                    return Err(PermissionError::GroupNotOwner);
2214                }
2215                if expected_group_revision.is_some_and(|revision| revision != group.revision) {
2216                    return Err(PermissionError::GroupRevisionConflict);
2217                }
2218                ids.extend(group.request_ids.iter().cloned());
2219            }
2220            PermissionSelector::DescendantRun(run_id) => {
2221                ids.extend(
2222                    state
2223                        .requests
2224                        .values()
2225                        .filter(|entry| {
2226                            let request = &entry.request;
2227                            request.session_id == actor.session_id
2228                                && request.requesting_run_id == *run_id
2229                                && self.flows.is_strict_ancestor(&actor.run_id, run_id)
2230                        })
2231                        .map(|entry| entry.request.request_id.clone()),
2232                );
2233            }
2234            selector => {
2235                ids.extend(state.requests.values().filter(|entry| {
2236                    let request = &entry.request;
2237                    if request.session_id != actor.session_id || !self.visible_to(actor, request) {
2238                        return false;
2239                    }
2240                    match selector {
2241                        PermissionSelector::ChildRun => request.parent_run_id.as_ref() == Some(&actor.run_id),
2242                        PermissionSelector::Tool(tool) => &request.intent.tool_name == tool,
2243                        PermissionSelector::Tier(tier) => request.intent.tier == *tier,
2244                        PermissionSelector::Risk(risk) => request.intent.risks.contains(risk),
2245                        PermissionSelector::PathPrefix(prefix) => request
2246                            .intent
2247                            .provenance
2248                            .authorized_targets()
2249                            .any(|path| path.starts_with(prefix)),
2250                        PermissionSelector::Target(target) => matches!(&request.state, PermissionRequestState::Pending { target: current } if current == target),
2251                        _ => false,
2252                    }
2253                }).map(|entry| entry.request.request_id.clone()));
2254            }
2255        }
2256        Ok(ids.into_iter().collect())
2257    }
2258
2259    /// Resolve a stable selector, deduplicating request IDs deterministically.
2260    #[allow(clippy::too_many_arguments)]
2261    pub fn resolve_batch(
2262        &self,
2263        actor: &Arc<FlowIdentity>,
2264        selector: PermissionSelector,
2265        action: PermissionAction,
2266        grant_scope: Option<GrantScope>,
2267        reason: Option<String>,
2268        mode: BatchMode,
2269        expected_group_revision: Option<u64>,
2270    ) -> Result<Vec<BatchResolution>, PermissionError> {
2271        let (results, sequence) = self.flows.with_lifecycle_arbitration(|| {
2272            self.authenticate_control_actor(actor)?;
2273            let authority = DecisionAuthority::Flow(self.flow_authority(Arc::clone(actor))?);
2274            let mut state = self.state.lock().unwrap();
2275            let ids = self.expand_selector(&state, actor, &selector, expected_group_revision)?;
2276            let request_ids = ids.iter().cloned().collect();
2277            let candidate_groups = Self::unresolved_groups_for_requests(&state, &request_ids);
2278            let mut prepared = if mode == BatchMode::Atomic {
2279                Some(
2280                    ids.iter()
2281                        .map(|id| {
2282                            self.prepare_resolution(
2283                                &state,
2284                                id,
2285                                &authority,
2286                                action,
2287                                grant_scope.as_ref(),
2288                            )
2289                        })
2290                        .collect::<Result<Vec<_>, _>>()?
2291                        .into_iter(),
2292                )
2293            } else {
2294                None
2295            };
2296            let mut results = Vec::with_capacity(ids.len());
2297            let mut audits = Vec::new();
2298            for request_id in ids {
2299                let outcome = if let Some(prepared) = &mut prepared {
2300                    Ok(self.commit_prepared_resolution(
2301                        &mut state,
2302                        prepared
2303                            .next()
2304                            .expect("one prepared resolution per request"),
2305                        action,
2306                        grant_scope.clone(),
2307                        reason.clone(),
2308                    ))
2309                } else {
2310                    self.resolve_locked_state(
2311                        &mut state,
2312                        &request_id,
2313                        &authority,
2314                        action,
2315                        grant_scope.clone(),
2316                        reason.clone(),
2317                    )
2318                };
2319                let batch = match outcome {
2320                    Ok(outcome) => {
2321                        let decision = match &outcome {
2322                            ResolveOutcome::Resolved(decision)
2323                            | ResolveOutcome::Deferred(decision) => decision,
2324                        };
2325                        audits.extend(resolution_audits_from_state(&state, &request_id, decision));
2326                        batch_outcome(action, outcome)
2327                    }
2328                    Err(error) if mode == BatchMode::BestEffort => batch_error_outcome(error),
2329                    Err(error) => return Err(error),
2330                };
2331                results.push(BatchResolution {
2332                    request_id,
2333                    outcome: batch,
2334                });
2335            }
2336            Self::append_resolved_group_audits(
2337                &mut state,
2338                &candidate_groups,
2339                &mut audits,
2340                Utc::now(),
2341            );
2342            let terminal_ids = results
2343                .iter()
2344                .map(|result| &result.request_id)
2345                .collect::<Vec<_>>();
2346            state.archive_terminals(terminal_ids);
2347            let sequence = Self::enqueue_audit_bundle(&mut state, audits);
2348            Ok((results, sequence))
2349        })?;
2350        self.dispatch_audit_bundle(sequence);
2351        Ok(results)
2352    }
2353
2354    pub fn delete_empty_group(
2355        &self,
2356        actor: &Arc<FlowIdentity>,
2357        group_id: &PermissionGroupId,
2358    ) -> Result<PermissionGroup, PermissionError> {
2359        self.authenticate_visibility_actor(actor)?;
2360        let mut state = self.state.lock().unwrap();
2361        let group = state
2362            .groups
2363            .get(group_id)
2364            .ok_or(PermissionError::GroupNotFound)?;
2365        if group.owner != GroupOwner::Flow(actor.run_id.clone()) {
2366            return Err(PermissionError::GroupNotOwner);
2367        }
2368        if !group.request_ids.is_empty() {
2369            return Err(PermissionError::GroupNotEmpty);
2370        }
2371        let group = state.groups.get(group_id).expect("group exists").clone();
2372        let mut records = Vec::new();
2373        Self::append_resolved_group_audits(
2374            &mut state,
2375            &BTreeSet::from([group_id.clone()]),
2376            &mut records,
2377            Utc::now(),
2378        );
2379        state.groups.remove(group_id).expect("group exists");
2380        let sequence = Self::enqueue_audit_bundle(&mut state, records);
2381        drop(state);
2382        self.dispatch_audit_bundle(sequence);
2383        Ok(group)
2384    }
2385
2386    fn authenticate_visibility_actor(
2387        &self,
2388        actor: &Arc<FlowIdentity>,
2389    ) -> Result<(), PermissionError> {
2390        self.authenticate_registered_identity(actor)?;
2391        if !matches!(actor.execution_state(), FlowExecutionState::Running) {
2392            return Err(PermissionError::ActorNotRunning);
2393        }
2394        if !actor.effective_authority.permission_management {
2395            return Err(PermissionError::PermissionManagementRequired);
2396        }
2397        Ok(())
2398    }
2399
2400    fn user_visible_request(request: &PermissionRequest, session_id: &str) -> bool {
2401        request.session_id == session_id
2402            && matches!(
2403                &request.state,
2404                PermissionRequestState::Pending {
2405                    target: ApprovalTarget::User
2406                }
2407            )
2408    }
2409
2410    fn visible_to(&self, actor: &FlowIdentity, request: &PermissionRequest) -> bool {
2411        request.session_id == actor.session_id
2412            && !request.state.is_terminal()
2413            && matches!(
2414                &request.state,
2415                PermissionRequestState::Pending {
2416                    target: ApprovalTarget::Flow(target),
2417                } if target == &actor.run_id
2418            )
2419            && self
2420                .flows
2421                .is_strict_ancestor(&actor.run_id, &request.requesting_run_id)
2422    }
2423
2424    pub fn grants(&self) -> Vec<PermissionGrant> {
2425        self.state.lock().unwrap().grants.clone()
2426    }
2427
2428    pub fn find_matching_grant(
2429        &self,
2430        identity: &FlowIdentity,
2431        intent: &PermissionIntent,
2432        shell: bool,
2433    ) -> Option<PermissionGrant> {
2434        let requirement = AuthorityRequirement::from_intent(intent, shell);
2435        self.flows.with_lifecycle_arbitration(|| {
2436            let state = self.state.lock().unwrap();
2437            matching_grant(&state.grants, identity, intent, &requirement)
2438        })
2439    }
2440
2441    fn authenticate_requester(
2442        &self,
2443        session_id: Option<&str>,
2444        run_id: Option<&FlowRunId>,
2445        requirement: &AuthorityRequirement,
2446        provenance: &ResourceProvenance,
2447    ) -> Result<Arc<FlowIdentity>, PermissionError> {
2448        let session_id = session_id.ok_or(PermissionError::MissingIdentity)?;
2449        let run_id = run_id.ok_or(PermissionError::MissingIdentity)?;
2450        let identity = self
2451            .flows
2452            .lookup_run(run_id)
2453            .ok_or(PermissionError::MissingIdentity)?;
2454        if identity.session_id != session_id {
2455            return Err(PermissionError::IdentityMismatch);
2456        }
2457        self.authenticate_registered_identity(&identity)?;
2458        if !matches!(identity.execution_state(), FlowExecutionState::Running) {
2459            return Err(PermissionError::ActorNotRunning);
2460        }
2461        if !authority_contains(&identity.effective_authority, requirement, provenance) {
2462            return Err(PermissionError::GrantExceedsAuthority);
2463        }
2464        Ok(identity)
2465    }
2466
2467    fn authenticate_registered_identity(
2468        &self,
2469        identity: &Arc<FlowIdentity>,
2470    ) -> Result<(), PermissionError> {
2471        let registered = self
2472            .flows
2473            .lookup_run(&identity.run_id)
2474            .ok_or(PermissionError::MissingIdentity)?;
2475        if !Arc::ptr_eq(&registered, identity) {
2476            return Err(PermissionError::IdentityMismatch);
2477        }
2478        Ok(())
2479    }
2480
2481    fn next_eligible_target(
2482        &self,
2483        requester: &FlowIdentity,
2484        requirement: &AuthorityRequirement,
2485        provenance: &ResourceProvenance,
2486        after: Option<&FlowRunId>,
2487    ) -> ApprovalTarget {
2488        let mut past_current = after.is_none();
2489        for ancestor in self.flows.strict_ancestors(&requester.run_id) {
2490            if !past_current {
2491                if after == Some(&ancestor.run_id) {
2492                    past_current = true;
2493                }
2494                continue;
2495            }
2496            if ancestor.session_id == requester.session_id
2497                && matches!(ancestor.execution_state(), FlowExecutionState::Running)
2498                && ancestor.effective_authority.permission_management
2499                && authority_contains(&ancestor.effective_authority, requirement, provenance)
2500            {
2501                return ApprovalTarget::Flow(ancestor.run_id.clone());
2502            }
2503        }
2504        ApprovalTarget::User
2505    }
2506
2507    fn authenticate_target(
2508        &self,
2509        requester: &FlowIdentity,
2510        requirement: &AuthorityRequirement,
2511        provenance: &ResourceProvenance,
2512        authority: ApprovalAuthority,
2513    ) -> Result<ApprovalTarget, PermissionError> {
2514        match authority {
2515            ApprovalAuthority::Flow(authority) => {
2516                self.authenticate_registered_identity(&authority.identity)?;
2517                let target = authority.identity;
2518                if target.session_id != requester.session_id
2519                    || !self
2520                        .flows
2521                        .is_strict_ancestor(&target.run_id, &requester.run_id)
2522                {
2523                    return Err(PermissionError::ActorNotAuthorized);
2524                }
2525                if !matches!(target.execution_state(), FlowExecutionState::Running) {
2526                    return Err(PermissionError::ActorNotRunning);
2527                }
2528                if !target.effective_authority.permission_management {
2529                    return Err(PermissionError::PermissionManagementRequired);
2530                }
2531                if !authority_contains(&target.effective_authority, requirement, provenance) {
2532                    return Err(PermissionError::GrantExceedsAuthority);
2533                }
2534                Ok(ApprovalTarget::Flow(target.run_id.clone()))
2535            }
2536            ApprovalAuthority::User(authority) => {
2537                if authority.broker_id != self.broker_id
2538                    || authority.session_id != requester.session_id
2539                {
2540                    return Err(PermissionError::ActorNotAuthorized);
2541                }
2542                Ok(ApprovalTarget::User)
2543            }
2544        }
2545    }
2546
2547    fn validate_authority(
2548        &self,
2549        request: &PermissionRequest,
2550        authority: &DecisionAuthority,
2551        action: PermissionAction,
2552        scope: Option<&GrantScope>,
2553    ) -> Result<DecisionActor, PermissionError> {
2554        if action != PermissionAction::Approve && scope.is_some() {
2555            return Err(PermissionError::ActorNotAuthorized);
2556        }
2557        if let Some(GrantScope::ChildRunSamePathRule { .. }) = scope {
2558            validate_same_path_scope(request, scope)?;
2559        }
2560        match authority {
2561            DecisionAuthority::Flow(authority) => {
2562                self.authenticate_registered_identity(&authority.identity)?;
2563                let identity = &authority.identity;
2564                if identity.session_id != request.session_id
2565                    || !self
2566                        .flows
2567                        .is_strict_ancestor(&identity.run_id, &request.requesting_run_id)
2568                {
2569                    return Err(PermissionError::ActorNotAuthorized);
2570                }
2571                if !matches!(identity.execution_state(), FlowExecutionState::Running) {
2572                    return Err(PermissionError::ActorNotRunning);
2573                }
2574                if !identity.effective_authority.permission_management {
2575                    return Err(PermissionError::PermissionManagementRequired);
2576                }
2577                if !authority_contains(
2578                    &identity.effective_authority,
2579                    &request.requirement,
2580                    &request.intent.provenance,
2581                ) {
2582                    return Err(PermissionError::GrantExceedsAuthority);
2583                }
2584                validate_same_tool_scope(request, scope)?;
2585                Ok(DecisionActor::Flow {
2586                    session_id: identity.session_id.clone(),
2587                    run_id: identity.run_id.clone(),
2588                })
2589            }
2590            DecisionAuthority::User(authority) => {
2591                if authority.broker_id != self.broker_id
2592                    || authority.session_id != request.session_id
2593                {
2594                    return Err(PermissionError::ActorNotAuthorized);
2595                }
2596                validate_same_tool_scope(request, scope)?;
2597                Ok(DecisionActor::User {
2598                    session_id: authority.session_id.clone(),
2599                    principal_id: authority.principal_id.clone(),
2600                })
2601            }
2602        }
2603    }
2604}
2605
2606fn request_audit_from_state(
2607    state: &BrokerState,
2608    request_id: &PermissionRequestId,
2609    decision: Option<&PermissionDecision>,
2610    at: DateTime<Utc>,
2611) -> Option<crate::permission_audit::PermissionRequestAudit> {
2612    let request = &state.requests.get(request_id)?.request;
2613    let group_ids = state
2614        .groups
2615        .values()
2616        .filter(|group| group.request_ids.contains(request_id))
2617        .map(|group| group.group_id.clone())
2618        .collect();
2619    Some(
2620        crate::permission_audit::PermissionRequestAudit::from_request(
2621            request, group_ids, decision, at,
2622        ),
2623    )
2624}
2625
2626fn submission_audits(
2627    request: &PermissionRequest,
2628    authorization: Option<&ImmediateAuthorization>,
2629    at: DateTime<Utc>,
2630) -> Vec<crate::permission_audit::PermissionAuditRecord> {
2631    let decision = authorization.map(|authorization| PermissionDecision {
2632        decision_id: match &request.state {
2633            PermissionRequestState::Approved { decision_id }
2634            | PermissionRequestState::Denied { decision_id } => decision_id.clone(),
2635            _ => unreachable!("immediate request is terminal"),
2636        },
2637        request_id: request.request_id.clone(),
2638        actor: match authorization {
2639            ImmediateAuthorization::Granted { grant } => grant.actor.clone(),
2640            _ => DecisionActor::Policy {
2641                policy_version: request.policy_reference.snapshot_id.clone(),
2642                rule_id: request.policy_reference.rule_id.clone(),
2643            },
2644        },
2645        action: if matches!(authorization, ImmediateAuthorization::Denied { .. }) {
2646            PermissionAction::Deny
2647        } else {
2648            PermissionAction::Approve
2649        },
2650        execution_boundary: match authorization {
2651            ImmediateAuthorization::Unrestricted => ExecutionBoundary::Direct,
2652            ImmediateAuthorization::Auto { execution_boundary } => *execution_boundary,
2653            ImmediateAuthorization::Granted { grant } => grant.execution_boundary,
2654            ImmediateAuthorization::Denied { .. } => ExecutionBoundary::Sandboxed,
2655        },
2656        grant_scope: match authorization {
2657            ImmediateAuthorization::Granted { grant } => Some(grant.scope.clone()),
2658            _ => None,
2659        },
2660        reason: match authorization {
2661            ImmediateAuthorization::Denied { reason } => Some(reason.clone()),
2662            _ => None,
2663        },
2664        decided_at: at,
2665    });
2666    let audit = crate::permission_audit::PermissionRequestAudit::from_request(
2667        request,
2668        Vec::new(),
2669        decision.as_ref(),
2670        at,
2671    );
2672    let mut records =
2673        vec![crate::permission_audit::PermissionAuditRecord::RequestCreated(audit.clone())];
2674    if let Some(authorization) = authorization {
2675        records.push(match authorization {
2676            ImmediateAuthorization::Denied { .. } => {
2677                crate::permission_audit::PermissionAuditRecord::RequestDenied(audit)
2678            }
2679            ImmediateAuthorization::Unrestricted => {
2680                crate::permission_audit::PermissionAuditRecord::UnrestrictedExecution(audit)
2681            }
2682            _ => crate::permission_audit::PermissionAuditRecord::RequestApproved(audit),
2683        });
2684    }
2685    records
2686}
2687
2688fn resolution_audits_from_state(
2689    state: &BrokerState,
2690    request_id: &PermissionRequestId,
2691    decision: &PermissionDecision,
2692) -> Vec<crate::permission_audit::PermissionAuditRecord> {
2693    let Some(entry) = state.requests.get(request_id) else {
2694        return Vec::new();
2695    };
2696    let Some(audit) =
2697        request_audit_from_state(state, request_id, Some(decision), decision.decided_at)
2698    else {
2699        return Vec::new();
2700    };
2701    let cancelled = matches!(
2702        entry.request.state,
2703        PermissionRequestState::Cancelled { .. }
2704    );
2705    let grant = state
2706        .grants
2707        .iter()
2708        .find(|grant| grant.request_id == *request_id && grant.granted_at == decision.decided_at);
2709    let mut records = vec![match decision.action {
2710        PermissionAction::Approve => {
2711            crate::permission_audit::PermissionAuditRecord::RequestApproved(audit)
2712        }
2713        PermissionAction::Deny => {
2714            crate::permission_audit::PermissionAuditRecord::RequestDenied(audit)
2715        }
2716        PermissionAction::Defer if cancelled => {
2717            crate::permission_audit::PermissionAuditRecord::RequestCancelled(audit)
2718        }
2719        PermissionAction::Defer => {
2720            crate::permission_audit::PermissionAuditRecord::RequestDeferred(audit)
2721        }
2722    }];
2723    if decision.action == PermissionAction::Defer
2724        && !cancelled
2725        && let Some(mut targeted) =
2726            request_audit_from_state(state, request_id, None, decision.decided_at)
2727    {
2728        targeted.actor = None;
2729        records.push(crate::permission_audit::PermissionAuditRecord::RequestTargeted(targeted));
2730    }
2731    if let Some(grant) = grant {
2732        records.push(
2733            crate::permission_audit::PermissionAuditRecord::GrantCreated(
2734                crate::permission_audit::PermissionGrantAudit::from_grant(
2735                    grant,
2736                    decision.reason.clone(),
2737                    decision.decided_at,
2738                ),
2739            ),
2740        );
2741    }
2742    records
2743}
2744
2745fn batch_outcome(action: PermissionAction, outcome: ResolveOutcome) -> BatchRequestOutcome {
2746    let decision = match outcome {
2747        ResolveOutcome::Resolved(decision) | ResolveOutcome::Deferred(decision) => decision,
2748    };
2749    match action {
2750        PermissionAction::Approve => BatchRequestOutcome::Approved(decision),
2751        PermissionAction::Deny => BatchRequestOutcome::Denied(decision),
2752        PermissionAction::Defer => BatchRequestOutcome::Deferred(decision),
2753    }
2754}
2755
2756fn batch_error_outcome(error: PermissionError) -> BatchRequestOutcome {
2757    match error {
2758        PermissionError::AlreadyResolved => BatchRequestOutcome::SkippedAlreadyResolved,
2759        PermissionError::ActorNotAuthorized => BatchRequestOutcome::RejectedNotAncestor,
2760        PermissionError::GrantExceedsAuthority => BatchRequestOutcome::RejectedOverAuthority,
2761        PermissionError::RequestNotFound => BatchRequestOutcome::RejectedNotFound,
2762        PermissionError::ActorNotRunning => BatchRequestOutcome::RejectedNotRunning,
2763        PermissionError::IdentityMismatch => BatchRequestOutcome::RejectedStale,
2764        PermissionError::PermissionManagementRequired => {
2765            BatchRequestOutcome::RejectedPermissionManagementRequired
2766        }
2767        PermissionError::UnsupportedGrantScope => {
2768            BatchRequestOutcome::RejectedUnsupportedGrantScope
2769        }
2770        PermissionError::NoEscalationTarget => BatchRequestOutcome::RejectedNoEscalationTarget,
2771        error => BatchRequestOutcome::Rejected(error),
2772    }
2773}
2774
2775fn cancel_entry(entry: &mut RequestEntry, component: &str, reason: String, at: DateTime<Utc>) {
2776    let target = match &entry.request.state {
2777        PermissionRequestState::Pending { target } => target.clone(),
2778        _ => ApprovalTarget::User,
2779    };
2780    entry.request.escalation_path.push(EscalationHop {
2781        target,
2782        actor: Some(DecisionActor::System {
2783            component: component.into(),
2784        }),
2785        action: None,
2786        reason: Some(reason.clone()),
2787        at,
2788    });
2789    entry.request.state = PermissionRequestState::Cancelled {
2790        reason: reason.clone(),
2791    };
2792    entry.request.revision += 1;
2793    if let Some(responder) = entry.responder.take() {
2794        let _ = responder.send(PermissionResolution::Cancelled { reason });
2795    }
2796}
2797
2798fn validate_same_tool_scope(
2799    request: &PermissionRequest,
2800    scope: Option<&GrantScope>,
2801) -> Result<(), PermissionError> {
2802    if let Some(GrantScope::ChildRunSameTool { run_id, tool_name }) = scope
2803        && (run_id != &request.requesting_run_id || tool_name != &request.intent.tool_name)
2804    {
2805        return Err(PermissionError::GrantExceedsAuthority);
2806    }
2807    Ok(())
2808}
2809
2810fn matching_grant(
2811    grants: &[PermissionGrant],
2812    identity: &FlowIdentity,
2813    intent: &PermissionIntent,
2814    requirement: &AuthorityRequirement,
2815) -> Option<PermissionGrant> {
2816    grants
2817        .iter()
2818        .filter(|grant| {
2819            grant.session_id == identity.session_id
2820                && grant.requesting_run_id == identity.run_id
2821                && requirement_contains(&grant.requirement, requirement)
2822                && match &grant.scope {
2823                    GrantScope::ChildRunSameTool { run_id, tool_name } => {
2824                        run_id == &identity.run_id && tool_name == &intent.tool_name
2825                    }
2826                    GrantScope::ChildRunSamePathRule {
2827                        run_id,
2828                        tool_name,
2829                        workspace_relative_path,
2830                    } => {
2831                        run_id == &identity.run_id
2832                            && tool_name == &intent.tool_name
2833                            && intent.provenance.authorized_targets().count() == 1
2834                            && intent
2835                                .provenance
2836                                .workspace_root
2837                                .as_deref()
2838                                .and_then(|root| {
2839                                    intent
2840                                        .provenance
2841                                        .authorized_targets()
2842                                        .next()
2843                                        .and_then(|path| path.strip_prefix(root).ok())
2844                                })
2845                                .is_some_and(|path| {
2846                                    path.to_str() == Some(workspace_relative_path.as_str())
2847                                })
2848                    }
2849                    GrantScope::CurrentCall => false,
2850                }
2851        })
2852        .max_by_key(|grant| grant.execution_boundary)
2853        .cloned()
2854}
2855
2856fn decision_execution_boundary(
2857    request: &PermissionRequest,
2858    actor: &DecisionActor,
2859    action: PermissionAction,
2860) -> ExecutionBoundary {
2861    if action == PermissionAction::Approve
2862        && request.intent.risks.contains(&RiskKind::ProcessSpawn)
2863        && matches!(actor, DecisionActor::User { .. })
2864    {
2865        ExecutionBoundary::Direct
2866    } else {
2867        ExecutionBoundary::Sandboxed
2868    }
2869}
2870
2871fn validate_same_path_scope(
2872    request: &PermissionRequest,
2873    scope: Option<&GrantScope>,
2874) -> Result<(), PermissionError> {
2875    let Some(GrantScope::ChildRunSamePathRule {
2876        run_id,
2877        tool_name,
2878        workspace_relative_path,
2879    }) = scope
2880    else {
2881        return Ok(());
2882    };
2883    if run_id != &request.requesting_run_id || tool_name != &request.intent.tool_name {
2884        return Err(PermissionError::GrantExceedsAuthority);
2885    }
2886    let provenance = &request.intent.provenance;
2887    if provenance.is_external()
2888        || provenance.is_unbound()
2889        || provenance.authorized_targets().count() != 1
2890        || workspace_relative_path.is_empty()
2891    {
2892        return Err(PermissionError::UnsupportedGrantScope);
2893    }
2894    let Some(root) = provenance.workspace_root.as_deref() else {
2895        return Err(PermissionError::UnsupportedGrantScope);
2896    };
2897    let Some(path) = provenance.authorized_targets().next() else {
2898        return Err(PermissionError::UnsupportedGrantScope);
2899    };
2900    let Ok(relative) = path.strip_prefix(root) else {
2901        return Err(PermissionError::GrantExceedsAuthority);
2902    };
2903    if relative.components().any(|component| {
2904        matches!(
2905            component,
2906            std::path::Component::ParentDir | std::path::Component::RootDir
2907        )
2908    }) || relative.as_os_str().is_empty()
2909        || relative.to_str() != Some(workspace_relative_path.as_str())
2910    {
2911        return Err(PermissionError::UnsupportedGrantScope);
2912    }
2913    Ok(())
2914}
2915
2916fn equivalent_grant(left: &PermissionGrant, right: &PermissionGrant) -> bool {
2917    left.session_id == right.session_id
2918        && left.requesting_run_id == right.requesting_run_id
2919        && left.requirement == right.requirement
2920        && left.execution_boundary == right.execution_boundary
2921        && left.scope == right.scope
2922}
2923
2924fn requirement_contains(grant: &AuthorityRequirement, requested: &AuthorityRequirement) -> bool {
2925    grant.tier == requested.tier
2926        && requested.risks.is_subset(&grant.risks)
2927        && (!requested.shell || grant.shell)
2928}
2929
2930fn actor_matches_target(actor: &DecisionActor, target: &ApprovalTarget) -> bool {
2931    matches!(
2932        (actor, target),
2933        (DecisionActor::Flow { run_id, .. }, ApprovalTarget::Flow(target_run)) if run_id == target_run
2934    ) || matches!(
2935        (actor, target),
2936        (DecisionActor::User { .. }, ApprovalTarget::User)
2937    )
2938}
2939
2940fn authority_contains(
2941    authority: &crate::flow_authority::EffectiveAuthority,
2942    requirement: &AuthorityRequirement,
2943    provenance: &ResourceProvenance,
2944) -> bool {
2945    let tier_index = match requirement.tier {
2946        Tier::Zero => 0,
2947        Tier::One => 1,
2948        Tier::Two => 2,
2949        Tier::Three => 3,
2950        Tier::Four => 4,
2951    };
2952    authority.allowed_tiers[tier_index]
2953        && requirement.risks.is_subset(&authority.allowed_risks)
2954        && (!requirement.shell || authority.shell)
2955        && provenance.authorized_targets().all(|target| {
2956            authority.workspace_root.as_deref().is_none_or(|root| {
2957                crate::fs_access::canonicalize_stable(target)
2958                    .strip_prefix(root)
2959                    .is_ok()
2960            })
2961        })
2962}
2963
2964#[cfg(test)]
2965mod tests {
2966    use super::*;
2967    use crate::flow_authority::{ChildWorkspaceAuthority, EffectiveAuthority, InvocationKind};
2968    use crate::stream::StreamFrame;
2969    use crate::trust::{EscalationPolicy, TrustMode};
2970
2971    fn authority(permission_management: bool) -> EffectiveAuthority {
2972        EffectiveAuthority {
2973            execution_policy: ExecutionPolicy::Controlled,
2974            allowed_tiers: [true; 5],
2975            allowed_risks: BTreeSet::from([
2976                RiskKind::WorkspaceExternal,
2977                RiskKind::Network,
2978                RiskKind::Irreversible,
2979                RiskKind::FilesystemWrite,
2980                RiskKind::ProcessSpawn,
2981                RiskKind::RepositoryMutation,
2982            ]),
2983            tier_ceiling: [PolicyAction::Auto; 5],
2984            risk_ceiling: [PolicyAction::Auto; 6],
2985            shell: true,
2986            permission_management,
2987            workspace_root: None,
2988        }
2989    }
2990
2991    fn register_root(
2992        flows: &Arc<FlowRegistry>,
2993        session_id: &str,
2994        permission_management: bool,
2995    ) -> Arc<FlowIdentity> {
2996        flows
2997            .register_root(
2998                session_id.into(),
2999                FlowRunId::now(),
3000                authority(permission_management),
3001            )
3002            .unwrap()
3003    }
3004
3005    fn child(flows: &Arc<FlowRegistry>, parent: &FlowIdentity) -> Arc<FlowIdentity> {
3006        flows
3007            .register_child(
3008                &parent.run_id,
3009                FlowRunId::now(),
3010                InvocationKind::InlineSubflow,
3011                true,
3012                ChildWorkspaceAuthority::Inherit,
3013            )
3014            .unwrap()
3015    }
3016
3017    fn ask_policy() -> TrustConfig {
3018        TrustConfig {
3019            mode: TrustMode::Steady,
3020            ..TrustConfig::default()
3021        }
3022    }
3023
3024    fn intent() -> PermissionIntent {
3025        PermissionIntent {
3026            tool_use_id: "call-1".into(),
3027            tool_name: "bash.spawn".into(),
3028            call_intent: None,
3029            tier: Tier::Two,
3030            risks: BTreeSet::new(),
3031            args_digest: "sha256:test".into(),
3032            preview: None,
3033            provenance: ResourceProvenance::none(),
3034        }
3035    }
3036
3037    fn process_intent() -> PermissionIntent {
3038        let mut intent = intent();
3039        intent.risks.insert(RiskKind::ProcessSpawn);
3040        intent.provenance = intent.provenance.with_risk(RiskKind::ProcessSpawn);
3041        intent
3042    }
3043
3044    fn submit_to_flow(
3045        broker: &PermissionBroker,
3046        requester: &FlowIdentity,
3047        target: Arc<FlowIdentity>,
3048    ) -> Result<SubmissionOutcome, PermissionError> {
3049        submit_to_flow_with_intent(broker, requester, target, intent())
3050    }
3051
3052    fn submit_to_flow_with_intent(
3053        broker: &PermissionBroker,
3054        requester: &FlowIdentity,
3055        target: Arc<FlowIdentity>,
3056        intent: PermissionIntent,
3057    ) -> Result<SubmissionOutcome, PermissionError> {
3058        let target = ApprovalAuthority::Flow(broker.flow_authority(target)?);
3059        broker.submit_to(
3060            Some(&requester.session_id),
3061            Some(&requester.run_id),
3062            intent,
3063            false,
3064            target,
3065            &ask_policy(),
3066        )
3067    }
3068
3069    fn set_blocked_on_child(parent: &FlowIdentity, child: &FlowIdentity) {
3070        *parent.execution_state.lock().unwrap() = FlowExecutionState::BlockedOnDescendants {
3071            child_run_counts: HashMap::from([(child.run_id.clone(), 1)]),
3072        };
3073    }
3074
3075    fn submit_user(
3076        broker: &PermissionBroker,
3077        requester: &FlowIdentity,
3078        intent: PermissionIntent,
3079    ) -> PendingPermission {
3080        let SubmissionOutcome::Pending(pending) = broker
3081            .submit(
3082                Some(&requester.session_id),
3083                Some(&requester.run_id),
3084                intent,
3085                false,
3086                &ask_policy(),
3087            )
3088            .unwrap()
3089        else {
3090            panic!("expected pending request");
3091        };
3092        *pending
3093    }
3094
3095    fn submit_to_user(broker: &PermissionBroker, requester: &FlowIdentity) -> PendingPermission {
3096        submit_to_user_with_intent(broker, requester, intent())
3097    }
3098
3099    fn submit_to_user_with_intent(
3100        broker: &PermissionBroker,
3101        requester: &FlowIdentity,
3102        intent: PermissionIntent,
3103    ) -> PendingPermission {
3104        let target = ApprovalAuthority::User(broker.user_authority(&requester.session_id, None));
3105        let SubmissionOutcome::Pending(pending) = broker
3106            .submit_to(
3107                Some(&requester.session_id),
3108                Some(&requester.run_id),
3109                intent,
3110                false,
3111                target,
3112                &ask_policy(),
3113            )
3114            .unwrap()
3115        else {
3116            panic!("expected pending request");
3117        };
3118        *pending
3119    }
3120
3121    fn user_decision(broker: &PermissionBroker, session_id: &str) -> DecisionAuthority {
3122        DecisionAuthority::User(broker.user_authority(session_id, None))
3123    }
3124
3125    fn assert_immediate_is_auditable(broker: &PermissionBroker, submission: &ImmediateSubmission) {
3126        assert_eq!(
3127            broker.get(&submission.request.request_id),
3128            Some(submission.request.clone())
3129        );
3130        assert!(!matches!(
3131            submission.request.state,
3132            PermissionRequestState::Pending { .. }
3133        ));
3134    }
3135
3136    fn attach_audit(broker: &PermissionBroker) -> tokio::sync::broadcast::Receiver<StreamFrame> {
3137        let (stream, receiver) = tokio::sync::broadcast::channel(64);
3138        broker.set_audit_projector(crate::permission_audit::PermissionAuditProjector::new(
3139            crate::event::EventSink::new(),
3140            stream,
3141        ));
3142        receiver
3143    }
3144
3145    fn drain_audit(
3146        receiver: &mut tokio::sync::broadcast::Receiver<StreamFrame>,
3147    ) -> Vec<StreamFrame> {
3148        std::iter::from_fn(|| receiver.try_recv().ok()).collect()
3149    }
3150
3151    #[test]
3152    fn invalid_flow_targets_do_not_create_requests() {
3153        let cases = [
3154            "self",
3155            "sibling",
3156            "cross-session",
3157            "terminal",
3158            "no-management",
3159        ];
3160        for case in cases {
3161            let flows = Arc::new(FlowRegistry::default());
3162            let root = register_root(&flows, "session", case != "no-management");
3163            let requester = child(&flows, &root);
3164            let target = match case {
3165                "self" => Arc::clone(&requester),
3166                "sibling" => child(&flows, &root),
3167                "cross-session" => register_root(&flows, "other", true),
3168                "terminal" => {
3169                    flows.mark_terminal(&root.run_id);
3170                    Arc::clone(&root)
3171                }
3172                "no-management" => Arc::clone(&root),
3173                _ => unreachable!(),
3174            };
3175            let broker = PermissionBroker::new(flows);
3176
3177            assert!(
3178                submit_to_flow(&broker, &requester, target).is_err(),
3179                "{case}"
3180            );
3181            assert!(broker.list().is_empty(), "{case} allocated a request");
3182        }
3183    }
3184
3185    #[test]
3186    fn running_managing_strict_ancestor_can_be_initial_target() {
3187        let flows = Arc::new(FlowRegistry::default());
3188        let root = register_root(&flows, "session", true);
3189        let requester = child(&flows, &root);
3190        let broker = PermissionBroker::new(flows);
3191
3192        let SubmissionOutcome::Pending(pending) =
3193            submit_to_flow(&broker, &requester, root).unwrap()
3194        else {
3195            panic!("expected pending request");
3196        };
3197        assert!(matches!(
3198            pending.request.state,
3199            PermissionRequestState::Pending {
3200                target: ApprovalTarget::Flow(_)
3201            }
3202        ));
3203    }
3204
3205    #[test]
3206    fn running_sync_parent_remains_an_escalation_target() {
3207        let flows = Arc::new(FlowRegistry::default());
3208        let root = register_root(&flows, "session", true);
3209        let sync_parent = flows
3210            .register_child(
3211                &root.run_id,
3212                FlowRunId::now(),
3213                InvocationKind::SpawnSync,
3214                true,
3215                ChildWorkspaceAuthority::Inherit,
3216            )
3217            .unwrap();
3218        let requester = child(&flows, &sync_parent);
3219        let broker = PermissionBroker::new(Arc::clone(&flows));
3220
3221        let pending = submit_user(&broker, &requester, intent());
3222        assert!(matches!(
3223            pending.request.state,
3224            PermissionRequestState::Pending {
3225                target: ApprovalTarget::Flow(ref run_id)
3226            } if run_id == &sync_parent.run_id
3227        ));
3228    }
3229
3230    #[test]
3231    fn blocked_sync_parent_is_skipped_during_escalation() {
3232        let flows = Arc::new(FlowRegistry::default());
3233        let root = register_root(&flows, "session", true);
3234        let sync_parent = flows
3235            .register_child(
3236                &root.run_id,
3237                FlowRunId::now(),
3238                InvocationKind::SpawnSync,
3239                true,
3240                ChildWorkspaceAuthority::Inherit,
3241            )
3242            .unwrap();
3243        let requester = child(&flows, &sync_parent);
3244        set_blocked_on_child(&sync_parent, &requester);
3245        let broker = PermissionBroker::new(Arc::clone(&flows));
3246
3247        let pending = submit_user(&broker, &requester, intent());
3248        assert!(matches!(
3249            pending.request.state,
3250            PermissionRequestState::Pending {
3251                target: ApprovalTarget::Flow(ref run_id)
3252            } if run_id == &root.run_id
3253        ));
3254    }
3255
3256    #[test]
3257    fn unrestricted_still_requires_authenticated_identity() {
3258        let broker = PermissionBroker::new(Arc::new(FlowRegistry::default()));
3259        let policy = TrustConfig {
3260            mode: TrustMode::Reckless,
3261            ..TrustConfig::default()
3262        };
3263
3264        assert!(matches!(
3265            broker.submit(None, None, intent(), false, &policy),
3266            Err(PermissionError::MissingIdentity)
3267        ));
3268    }
3269
3270    #[test]
3271    fn submissions_emit_only_canonical_lifecycle_transitions() {
3272        let flows = Arc::new(FlowRegistry::default());
3273        let controlled = register_root(&flows, "session", false);
3274        let unrestricted_policy = TrustConfig {
3275            mode: TrustMode::Reckless,
3276            ..TrustConfig::default()
3277        };
3278        let unrestricted = flows
3279            .register_root(
3280                "session".into(),
3281                FlowRunId::now(),
3282                EffectiveAuthority::root(&unrestricted_policy, true, None),
3283            )
3284            .unwrap();
3285        let broker = PermissionBroker::new(flows);
3286        let mut audit = attach_audit(&broker);
3287
3288        let pending = submit_to_user(&broker, &controlled);
3289        assert!(matches!(
3290            drain_audit(&mut audit).as_slice(),
3291            [StreamFrame::PermissionRequestCreated { .. }]
3292        ));
3293        broker
3294            .cancel(&pending.request.request_id, "test cleanup")
3295            .unwrap();
3296        drain_audit(&mut audit);
3297
3298        let auto_policy = TrustConfig {
3299            mode: TrustMode::Eager,
3300            escalation: EscalationPolicy::Allow,
3301            ..TrustConfig::default()
3302        };
3303        assert!(matches!(
3304            broker.submit(
3305                Some(&controlled.session_id),
3306                Some(&controlled.run_id),
3307                intent(),
3308                false,
3309                &auto_policy,
3310            ),
3311            Ok(SubmissionOutcome::Immediate(_))
3312        ));
3313        assert!(matches!(
3314            drain_audit(&mut audit).as_slice(),
3315            [
3316                StreamFrame::PermissionRequestCreated { .. },
3317                StreamFrame::PermissionRequestApproved { .. }
3318            ]
3319        ));
3320
3321        let deny_policy = TrustConfig {
3322            mode: TrustMode::Eager,
3323            escalation: EscalationPolicy::Deny,
3324            ..TrustConfig::default()
3325        };
3326        assert!(matches!(
3327            broker.submit(
3328                Some(&controlled.session_id),
3329                Some(&controlled.run_id),
3330                intent(),
3331                false,
3332                &deny_policy,
3333            ),
3334            Ok(SubmissionOutcome::Immediate(_))
3335        ));
3336        assert!(matches!(
3337            drain_audit(&mut audit).as_slice(),
3338            [
3339                StreamFrame::PermissionRequestCreated { .. },
3340                StreamFrame::PermissionRequestDenied { .. }
3341            ]
3342        ));
3343
3344        assert!(matches!(
3345            broker.submit(
3346                Some(&unrestricted.session_id),
3347                Some(&unrestricted.run_id),
3348                intent(),
3349                false,
3350                &unrestricted_policy,
3351            ),
3352            Ok(SubmissionOutcome::Immediate(_))
3353        ));
3354        assert!(matches!(
3355            drain_audit(&mut audit).as_slice(),
3356            [
3357                StreamFrame::PermissionRequestCreated { .. },
3358                StreamFrame::UnrestrictedExecution { .. }
3359            ]
3360        ));
3361    }
3362
3363    #[test]
3364    fn deny_is_immediate_and_persists_terminal_request() {
3365        let flows = Arc::new(FlowRegistry::default());
3366        let requester = register_root(&flows, "session", false);
3367        let broker = PermissionBroker::new(flows);
3368        let policy = TrustConfig {
3369            mode: TrustMode::Eager,
3370            escalation: EscalationPolicy::Deny,
3371            ..TrustConfig::default()
3372        };
3373
3374        let outcome = broker
3375            .submit(
3376                Some(&requester.session_id),
3377                Some(&requester.run_id),
3378                intent(),
3379                false,
3380                &policy,
3381            )
3382            .unwrap();
3383
3384        let SubmissionOutcome::Immediate(submission) = outcome else {
3385            panic!("expected immediate denial");
3386        };
3387        assert!(matches!(
3388            submission.authorization,
3389            ImmediateAuthorization::Denied { .. }
3390        ));
3391        assert_immediate_is_auditable(&broker, &submission);
3392        assert!(matches!(
3393            submission.request.state,
3394            PermissionRequestState::Denied { .. }
3395        ));
3396    }
3397
3398    #[test]
3399    fn submit_uses_each_policy_snapshot_without_cross_request_state() {
3400        let flows = Arc::new(FlowRegistry::default());
3401        let requester = register_root(&flows, "session", false);
3402        let broker = PermissionBroker::new(Arc::clone(&flows));
3403        let auto = TrustConfig {
3404            mode: TrustMode::Eager,
3405            escalation: EscalationPolicy::Allow,
3406            ..TrustConfig::default()
3407        };
3408        let ask = ask_policy();
3409        let deny = TrustConfig {
3410            mode: TrustMode::Eager,
3411            escalation: EscalationPolicy::Deny,
3412            ..TrustConfig::default()
3413        };
3414        let unrestricted = TrustConfig {
3415            mode: TrustMode::Reckless,
3416            ..TrustConfig::default()
3417        };
3418        let unrestricted_requester = flows
3419            .register_root(
3420                "session".into(),
3421                FlowRunId::now(),
3422                EffectiveAuthority::root(&unrestricted, true, None),
3423            )
3424            .unwrap();
3425
3426        let SubmissionOutcome::Immediate(auto_submission) = broker
3427            .submit(
3428                Some(&requester.session_id),
3429                Some(&requester.run_id),
3430                intent(),
3431                false,
3432                &auto,
3433            )
3434            .unwrap()
3435        else {
3436            panic!("expected immediate auto authorization");
3437        };
3438        assert!(matches!(
3439            auto_submission.authorization,
3440            ImmediateAuthorization::Auto { .. }
3441        ));
3442        assert_immediate_is_auditable(&broker, &auto_submission);
3443
3444        let SubmissionOutcome::Pending(pending) = broker
3445            .submit(
3446                Some(&requester.session_id),
3447                Some(&requester.run_id),
3448                intent(),
3449                false,
3450                &ask,
3451            )
3452            .unwrap()
3453        else {
3454            panic!("expected pending request");
3455        };
3456
3457        let SubmissionOutcome::Immediate(denied_submission) = broker
3458            .submit(
3459                Some(&requester.session_id),
3460                Some(&requester.run_id),
3461                intent(),
3462                false,
3463                &deny,
3464            )
3465            .unwrap()
3466        else {
3467            panic!("expected immediate denial");
3468        };
3469        assert!(matches!(
3470            denied_submission.authorization,
3471            ImmediateAuthorization::Denied { .. }
3472        ));
3473        assert_immediate_is_auditable(&broker, &denied_submission);
3474
3475        let SubmissionOutcome::Immediate(unrestricted_submission) = broker
3476            .submit(
3477                Some(&unrestricted_requester.session_id),
3478                Some(&unrestricted_requester.run_id),
3479                intent(),
3480                false,
3481                &unrestricted,
3482            )
3483            .unwrap()
3484        else {
3485            panic!("expected immediate unrestricted authorization");
3486        };
3487        assert!(matches!(
3488            unrestricted_submission.authorization,
3489            ImmediateAuthorization::Unrestricted
3490        ));
3491        assert_immediate_is_auditable(&broker, &unrestricted_submission);
3492        assert_eq!(broker.list().len(), 4);
3493        assert_eq!(
3494            broker
3495                .list()
3496                .into_iter()
3497                .filter(|request| matches!(request.state, PermissionRequestState::Pending { .. }))
3498                .count(),
3499            1
3500        );
3501
3502        broker
3503            .cancel(&pending.request.request_id, "test cleanup")
3504            .unwrap();
3505        assert!(matches!(
3506            broker.get(&pending.request.request_id).unwrap().state,
3507            PermissionRequestState::Cancelled { .. }
3508        ));
3509        assert_eq!(broker.list().len(), 4);
3510    }
3511
3512    #[test]
3513    fn immediate_terminal_retention_is_bounded() {
3514        let flows = Arc::new(FlowRegistry::default());
3515        let requester = register_root(&flows, "session", false);
3516        let broker = PermissionBroker::new(flows);
3517        let policy = TrustConfig {
3518            mode: TrustMode::Eager,
3519            escalation: EscalationPolicy::Allow,
3520            ..TrustConfig::default()
3521        };
3522        let mut first_id = None;
3523        let mut last_id = None;
3524
3525        for _ in 0..10_000 {
3526            let SubmissionOutcome::Immediate(submission) = broker
3527                .submit(
3528                    Some(&requester.session_id),
3529                    Some(&requester.run_id),
3530                    intent(),
3531                    false,
3532                    &policy,
3533                )
3534                .unwrap()
3535            else {
3536                panic!("expected immediate authorization");
3537            };
3538            first_id.get_or_insert_with(|| submission.request.request_id.clone());
3539            last_id = Some(submission.request.request_id.clone());
3540        }
3541
3542        let state = broker.state.lock().unwrap();
3543        assert!(state.requests.is_empty());
3544        assert_eq!(
3545            state.recent_terminal_requests.len(),
3546            RECENT_TERMINAL_REQUEST_LIMIT
3547        );
3548        drop(state);
3549        assert_eq!(broker.list().len(), RECENT_TERMINAL_REQUEST_LIMIT);
3550        let first_id = first_id.unwrap();
3551        assert!(broker.get(&first_id).is_none());
3552        assert!(matches!(
3553            broker.cancel(&first_id, "already evicted"),
3554            Err(PermissionError::RequestNotFound)
3555        ));
3556        assert!(matches!(
3557            broker.get(&last_id.unwrap()).map(|request| request.state),
3558            Some(PermissionRequestState::Approved { .. })
3559        ));
3560    }
3561
3562    #[test]
3563    fn resolved_terminal_retention_is_bounded() {
3564        let flows = Arc::new(FlowRegistry::default());
3565        let requester = register_root(&flows, "session", false);
3566        let broker = PermissionBroker::new(flows);
3567        let authority = user_decision(&broker, "session");
3568        let mut first_id = None;
3569        let mut last_id = None;
3570
3571        for _ in 0..10_000 {
3572            let pending = submit_to_user(&broker, &requester);
3573            let request_id = pending.request.request_id.clone();
3574            first_id.get_or_insert_with(|| request_id.clone());
3575            broker
3576                .resolve(
3577                    &request_id,
3578                    &authority,
3579                    PermissionAction::Approve,
3580                    Some(GrantScope::CurrentCall),
3581                    None,
3582                )
3583                .unwrap();
3584            last_id = Some(request_id);
3585        }
3586
3587        let state = broker.state.lock().unwrap();
3588        assert!(state.requests.is_empty());
3589        assert_eq!(
3590            state.recent_terminal_requests.len(),
3591            RECENT_TERMINAL_REQUEST_LIMIT
3592        );
3593        drop(state);
3594        assert_eq!(broker.list().len(), RECENT_TERMINAL_REQUEST_LIMIT);
3595        let first_id = first_id.unwrap();
3596        assert!(broker.get(&first_id).is_none());
3597        assert!(matches!(
3598            broker.resolve(
3599                &first_id,
3600                &authority,
3601                PermissionAction::Approve,
3602                Some(GrantScope::CurrentCall),
3603                None,
3604            ),
3605            Err(PermissionError::RequestNotFound)
3606        ));
3607        assert!(matches!(
3608            broker.get(&last_id.unwrap()).map(|request| request.state),
3609            Some(PermissionRequestState::Approved { .. })
3610        ));
3611    }
3612
3613    #[test]
3614    fn recent_terminal_requests_keep_resolved_operation_semantics() {
3615        let flows = Arc::new(FlowRegistry::default());
3616        let requester = register_root(&flows, "session", false);
3617        let broker = PermissionBroker::new(flows);
3618        let pending = submit_to_user(&broker, &requester);
3619        let request_id = pending.request.request_id.clone();
3620        let authority = user_decision(&broker, "session");
3621        broker
3622            .resolve(
3623                &request_id,
3624                &authority,
3625                PermissionAction::Approve,
3626                Some(GrantScope::CurrentCall),
3627                None,
3628            )
3629            .unwrap();
3630
3631        assert!(matches!(
3632            broker.resolve(
3633                &request_id,
3634                &authority,
3635                PermissionAction::Approve,
3636                Some(GrantScope::CurrentCall),
3637                None,
3638            ),
3639            Err(PermissionError::AlreadyResolved)
3640        ));
3641        assert!(matches!(
3642            broker.cancel(&request_id, "already resolved"),
3643            Err(PermissionError::AlreadyResolved)
3644        ));
3645        assert!(
3646            !broker
3647                .defer_timed_out_target(&request_id, &requester.run_id)
3648                .unwrap()
3649        );
3650    }
3651
3652    #[test]
3653    fn process_boundary_distinguishes_policy_flow_and_user_authority() {
3654        let flows = Arc::new(FlowRegistry::default());
3655        let root = register_root(&flows, "session", true);
3656        let requester = child(&flows, &root);
3657        let broker = PermissionBroker::new(Arc::clone(&flows));
3658
3659        let eager_deny = TrustConfig {
3660            mode: TrustMode::Eager,
3661            escalation: EscalationPolicy::Deny,
3662            ..Default::default()
3663        };
3664        let SubmissionOutcome::Immediate(sandboxed) = broker
3665            .submit(
3666                Some(&requester.session_id),
3667                Some(&requester.run_id),
3668                process_intent(),
3669                false,
3670                &eager_deny,
3671            )
3672            .unwrap()
3673        else {
3674            panic!("expected sandboxed process authorization");
3675        };
3676        assert!(matches!(
3677            sandboxed.authorization,
3678            ImmediateAuthorization::Auto {
3679                execution_boundary: ExecutionBoundary::Sandboxed
3680            }
3681        ));
3682
3683        let automatic = TrustConfig {
3684            mode: TrustMode::Eager,
3685            tiers: crate::trust::TierPolicyConfig {
3686                eager: crate::trust::TierPolicyOverrides {
3687                    tier2: Some(PolicyAction::Auto),
3688                    ..Default::default()
3689                },
3690            },
3691            risks: crate::trust::RiskPolicyConfig {
3692                eager: crate::trust::RiskPolicyOverrides {
3693                    process_spawn: Some(PolicyAction::Auto),
3694                    ..Default::default()
3695                },
3696            },
3697            ..Default::default()
3698        };
3699        let SubmissionOutcome::Immediate(automatic) = broker
3700            .submit(
3701                Some(&requester.session_id),
3702                Some(&requester.run_id),
3703                process_intent(),
3704                false,
3705                &automatic,
3706            )
3707            .unwrap()
3708        else {
3709            panic!("expected automatic authorization");
3710        };
3711        assert!(matches!(
3712            automatic.authorization,
3713            ImmediateAuthorization::Auto {
3714                execution_boundary: ExecutionBoundary::Sandboxed
3715            }
3716        ));
3717
3718        let eager_allow = TrustConfig {
3719            mode: TrustMode::Eager,
3720            escalation: EscalationPolicy::Allow,
3721            ..Default::default()
3722        };
3723        let SubmissionOutcome::Immediate(allowed) = broker
3724            .submit(
3725                Some(&requester.session_id),
3726                Some(&requester.run_id),
3727                process_intent(),
3728                false,
3729                &eager_allow,
3730            )
3731            .unwrap()
3732        else {
3733            panic!("expected eager allow authorization");
3734        };
3735        assert!(matches!(
3736            allowed.authorization,
3737            ImmediateAuthorization::Auto {
3738                execution_boundary: ExecutionBoundary::Direct
3739            }
3740        ));
3741
3742        let user_pending = submit_to_user_with_intent(&broker, &requester, process_intent());
3743        assert_eq!(
3744            crate::permission_audit::PermissionRequestAudit::from_request(
3745                &user_pending.request,
3746                Vec::new(),
3747                None,
3748                Utc::now(),
3749            )
3750            .execution_boundary,
3751            Some(ExecutionBoundary::Direct)
3752        );
3753        let ResolveOutcome::Resolved(user_decision) = broker
3754            .resolve(
3755                &user_pending.request.request_id,
3756                &user_decision(&broker, "session"),
3757                PermissionAction::Approve,
3758                None,
3759                None,
3760            )
3761            .unwrap()
3762        else {
3763            panic!("user approval must resolve");
3764        };
3765        assert_eq!(user_decision.execution_boundary, ExecutionBoundary::Direct);
3766
3767        let SubmissionOutcome::Pending(flow_pending) =
3768            submit_to_flow_with_intent(&broker, &requester, Arc::clone(&root), process_intent())
3769                .unwrap()
3770        else {
3771            panic!("expected flow approval");
3772        };
3773        let ResolveOutcome::Resolved(flow_decision) = broker
3774            .resolve(
3775                &flow_pending.request.request_id,
3776                &DecisionAuthority::Flow(broker.flow_authority(root).unwrap()),
3777                PermissionAction::Approve,
3778                None,
3779                None,
3780            )
3781            .unwrap()
3782        else {
3783            panic!("flow approval must resolve");
3784        };
3785        assert_eq!(
3786            flow_decision.execution_boundary,
3787            ExecutionBoundary::Sandboxed
3788        );
3789    }
3790
3791    #[test]
3792    fn user_process_grant_reuses_direct_boundary() {
3793        let flows = Arc::new(FlowRegistry::default());
3794        let requester = register_root(&flows, "session", false);
3795        let broker = PermissionBroker::new(flows);
3796        let pending = submit_to_user_with_intent(&broker, &requester, process_intent());
3797        let scope = GrantScope::ChildRunSameTool {
3798            run_id: requester.run_id.clone(),
3799            tool_name: "bash.spawn".into(),
3800        };
3801
3802        broker
3803            .resolve(
3804                &pending.request.request_id,
3805                &user_decision(&broker, "session"),
3806                PermissionAction::Approve,
3807                Some(scope),
3808                None,
3809            )
3810            .unwrap();
3811
3812        let SubmissionOutcome::Immediate(granted) = broker
3813            .submit(
3814                Some(&requester.session_id),
3815                Some(&requester.run_id),
3816                process_intent(),
3817                false,
3818                &ask_policy(),
3819            )
3820            .unwrap()
3821        else {
3822            panic!("expected reusable grant");
3823        };
3824        assert!(matches!(
3825            granted.authorization,
3826            ImmediateAuthorization::Granted { grant }
3827                if grant.execution_boundary == ExecutionBoundary::Direct
3828        ));
3829    }
3830
3831    #[test]
3832    fn registered_root_applies_live_session_policy_changes() {
3833        let flows = Arc::new(FlowRegistry::default());
3834        let started_controlled = TrustConfig::default();
3835        let requester = flows
3836            .register_root(
3837                "session".into(),
3838                FlowRunId::now(),
3839                EffectiveAuthority::root(&started_controlled, true, None),
3840            )
3841            .unwrap();
3842        let broker = PermissionBroker::new(flows);
3843        let reckless = TrustConfig {
3844            mode: TrustMode::Reckless,
3845            ..TrustConfig::default()
3846        };
3847
3848        let SubmissionOutcome::Immediate(unrestricted) = broker
3849            .submit(
3850                Some(&requester.session_id),
3851                Some(&requester.run_id),
3852                intent(),
3853                false,
3854                &reckless,
3855            )
3856            .unwrap()
3857        else {
3858            panic!("live Reckless policy must authorize the registered root");
3859        };
3860        assert!(matches!(
3861            unrestricted.authorization,
3862            ImmediateAuthorization::Unrestricted
3863        ));
3864
3865        let SubmissionOutcome::Pending(pending) = broker
3866            .submit(
3867                Some(&requester.session_id),
3868                Some(&requester.run_id),
3869                intent(),
3870                false,
3871                &ask_policy(),
3872            )
3873            .unwrap()
3874        else {
3875            panic!("the same root must return to controlled policy");
3876        };
3877        broker
3878            .cancel(&pending.request.request_id, "test cleanup")
3879            .unwrap();
3880    }
3881
3882    #[test]
3883    fn broker_bound_user_authority_cannot_be_reused() {
3884        let flows = Arc::new(FlowRegistry::default());
3885        let requester = register_root(&flows, "session", false);
3886        let first = PermissionBroker::new(Arc::clone(&flows));
3887        let second = PermissionBroker::new(flows);
3888        let pending = submit_user(&second, &requester, intent());
3889        let foreign = user_decision(&first, "session");
3890
3891        assert!(matches!(
3892            second.resolve(
3893                &pending.request.request_id,
3894                &foreign,
3895                PermissionAction::Approve,
3896                None,
3897                None,
3898            ),
3899            Err(PermissionError::ActorNotAuthorized)
3900        ));
3901    }
3902
3903    #[test]
3904    fn cancellation_wakes_waiter_with_fail_closed_resolution() {
3905        let flows = Arc::new(FlowRegistry::default());
3906        let requester = register_root(&flows, "session", false);
3907        let broker = PermissionBroker::new(flows);
3908        let pending = submit_user(&broker, &requester, intent());
3909
3910        broker
3911            .cancel(&pending.request.request_id, "flow stopped")
3912            .unwrap();
3913
3914        assert_eq!(
3915            pending.resolution.blocking_recv().unwrap(),
3916            PermissionResolution::Cancelled {
3917                reason: "flow stopped".into()
3918            }
3919        );
3920    }
3921
3922    #[test]
3923    fn terminal_run_cleanup_cancels_pending_and_removes_grants() {
3924        let flows = Arc::new(FlowRegistry::default());
3925        let requester = register_root(&flows, "session", false);
3926        let broker = PermissionBroker::new(Arc::clone(&flows));
3927        let first = submit_user(&broker, &requester, intent());
3928        let authority = user_decision(&broker, "session");
3929        let scope = GrantScope::ChildRunSameTool {
3930            run_id: requester.run_id.clone(),
3931            tool_name: "bash.spawn".into(),
3932        };
3933        broker
3934            .resolve(
3935                &first.request.request_id,
3936                &authority,
3937                PermissionAction::Approve,
3938                Some(scope),
3939                None,
3940            )
3941            .unwrap();
3942        let pending = submit_user(
3943            &broker,
3944            &requester,
3945            PermissionIntent {
3946                tool_name: "fs.write".into(),
3947                ..intent()
3948            },
3949        );
3950
3951        flows.mark_terminal(&requester.run_id);
3952        assert_eq!(broker.expire_terminal("flow stopped"), 1);
3953        assert!(broker.grants().is_empty());
3954        assert!(matches!(
3955            pending.resolution.blocking_recv().unwrap(),
3956            PermissionResolution::Cancelled { .. }
3957        ));
3958    }
3959
3960    #[test]
3961    fn same_tool_grant_covers_narrower_risk_but_not_escalation() {
3962        let flows = Arc::new(FlowRegistry::default());
3963        let requester = register_root(&flows, "session", false);
3964        let broker = PermissionBroker::new(flows);
3965        let mut approved = intent();
3966        approved.risks.insert(RiskKind::Network);
3967        let pending = submit_user(&broker, &requester, approved);
3968        let scope = GrantScope::ChildRunSameTool {
3969            run_id: requester.run_id.clone(),
3970            tool_name: "bash.spawn".into(),
3971        };
3972        broker
3973            .resolve(
3974                &pending.request.request_id,
3975                &user_decision(&broker, "session"),
3976                PermissionAction::Approve,
3977                Some(scope),
3978                None,
3979            )
3980            .unwrap();
3981
3982        let SubmissionOutcome::Immediate(granted_submission) = broker
3983            .submit(
3984                Some(&requester.session_id),
3985                Some(&requester.run_id),
3986                intent(),
3987                false,
3988                &ask_policy(),
3989            )
3990            .unwrap()
3991        else {
3992            panic!("expected immediate grant authorization");
3993        };
3994        assert!(matches!(
3995            granted_submission.authorization,
3996            ImmediateAuthorization::Granted { .. }
3997        ));
3998        assert_immediate_is_auditable(&broker, &granted_submission);
3999        let mut escalated = intent();
4000        escalated.risks.insert(RiskKind::Network);
4001        escalated.risks.insert(RiskKind::Irreversible);
4002        assert!(matches!(
4003            broker.submit(
4004                Some(&requester.session_id),
4005                Some(&requester.run_id),
4006                escalated,
4007                false,
4008                &ask_policy(),
4009            ),
4010            Ok(SubmissionOutcome::Pending(_))
4011        ));
4012    }
4013
4014    #[test]
4015    fn current_call_is_not_persisted_and_equivalent_grants_are_deduplicated() {
4016        let flows = Arc::new(FlowRegistry::default());
4017        let requester = register_root(&flows, "session", false);
4018        let broker = PermissionBroker::new(flows);
4019        let authority = user_decision(&broker, "session");
4020
4021        let current = submit_user(&broker, &requester, intent());
4022        broker
4023            .resolve(
4024                &current.request.request_id,
4025                &authority,
4026                PermissionAction::Approve,
4027                None,
4028                None,
4029            )
4030            .unwrap();
4031        assert!(broker.grants().is_empty());
4032
4033        let pending = submit_user(
4034            &broker,
4035            &requester,
4036            PermissionIntent {
4037                tool_use_id: "call-2".into(),
4038                ..intent()
4039            },
4040        );
4041        broker
4042            .resolve(
4043                &pending.request.request_id,
4044                &authority,
4045                PermissionAction::Approve,
4046                Some(GrantScope::ChildRunSameTool {
4047                    run_id: requester.run_id.clone(),
4048                    tool_name: "bash.spawn".into(),
4049                }),
4050                None,
4051            )
4052            .unwrap();
4053        assert!(matches!(
4054            broker.submit(
4055                Some(&requester.session_id),
4056                Some(&requester.run_id),
4057                PermissionIntent {
4058                    tool_use_id: "call-3".into(),
4059                    ..intent()
4060                },
4061                false,
4062                &ask_policy(),
4063            ),
4064            Ok(SubmissionOutcome::Immediate(value)) if matches!(value.authorization, ImmediateAuthorization::Granted { .. })
4065        ));
4066        assert_eq!(broker.grants().len(), 1);
4067    }
4068
4069    #[test]
4070    fn terminal_transition_cancels_pending_without_explicit_sweep() {
4071        let flows = Arc::new(FlowRegistry::default());
4072        let requester = register_root(&flows, "session", false);
4073        let broker = PermissionBroker::shared(Arc::clone(&flows));
4074        let first = submit_user(&broker, &requester, intent());
4075        let second = submit_user(
4076            &broker,
4077            &requester,
4078            PermissionIntent {
4079                tool_use_id: "call-2".into(),
4080                tool_name: "fs.write".into(),
4081                ..intent()
4082            },
4083        );
4084
4085        flows.mark_terminal(&requester.run_id);
4086
4087        for pending in [first, second] {
4088            assert!(matches!(
4089                pending.resolution.blocking_recv().unwrap(),
4090                PermissionResolution::Cancelled { .. }
4091            ));
4092        }
4093        assert!(broker.grants().is_empty());
4094        assert_eq!(broker.expire_terminal("late sweep"), 0);
4095    }
4096
4097    #[test]
4098    fn terminal_before_resolve_leaves_no_grant_and_rejects_approval() {
4099        let flows = Arc::new(FlowRegistry::default());
4100        let requester = register_root(&flows, "session", false);
4101        let broker = PermissionBroker::shared(Arc::clone(&flows));
4102        let pending = submit_user(&broker, &requester, intent());
4103        let authority = user_decision(&broker, "session");
4104
4105        flows.mark_terminal(&requester.run_id);
4106
4107        assert!(matches!(
4108            broker.resolve(
4109                &pending.request.request_id,
4110                &authority,
4111                PermissionAction::Approve,
4112                Some(GrantScope::ChildRunSameTool {
4113                    run_id: requester.run_id.clone(),
4114                    tool_name: "bash.spawn".into(),
4115                }),
4116                None,
4117            ),
4118            Err(PermissionError::AlreadyResolved)
4119        ));
4120        assert!(broker.grants().is_empty());
4121        assert!(matches!(
4122            broker.get(&pending.request.request_id).unwrap().state,
4123            PermissionRequestState::Cancelled { .. }
4124        ));
4125        assert!(matches!(
4126            pending.resolution.blocking_recv().unwrap(),
4127            PermissionResolution::Cancelled { .. }
4128        ));
4129    }
4130
4131    #[test]
4132    fn resolve_before_terminal_keeps_decision_but_revokes_grant() {
4133        let flows = Arc::new(FlowRegistry::default());
4134        let requester = register_root(&flows, "session", false);
4135        let broker = PermissionBroker::shared(Arc::clone(&flows));
4136        let pending = submit_user(&broker, &requester, intent());
4137
4138        broker
4139            .resolve(
4140                &pending.request.request_id,
4141                &user_decision(&broker, "session"),
4142                PermissionAction::Approve,
4143                Some(GrantScope::ChildRunSameTool {
4144                    run_id: requester.run_id.clone(),
4145                    tool_name: "bash.spawn".into(),
4146                }),
4147                None,
4148            )
4149            .unwrap();
4150        assert_eq!(broker.grants().len(), 1);
4151
4152        flows.mark_terminal(&requester.run_id);
4153
4154        assert!(broker.grants().is_empty());
4155        assert!(matches!(
4156            broker.get(&pending.request.request_id).unwrap().state,
4157            PermissionRequestState::Approved { .. }
4158        ));
4159        assert!(matches!(
4160            pending.resolution.blocking_recv().unwrap(),
4161            PermissionResolution::Decision(_)
4162        ));
4163    }
4164
4165    #[test]
4166    fn terminal_racing_approve_never_leaves_a_terminal_run_authorized() {
4167        for _ in 0..64 {
4168            let flows = Arc::new(FlowRegistry::default());
4169            let requester = register_root(&flows, "session", false);
4170            let broker = PermissionBroker::shared(Arc::clone(&flows));
4171            let pending = submit_user(&broker, &requester, intent());
4172            let request_id = pending.request.request_id.clone();
4173            let barrier = Arc::new(std::sync::Barrier::new(2));
4174
4175            let approver = {
4176                let broker = Arc::clone(&broker);
4177                let barrier = Arc::clone(&barrier);
4178                let authority = user_decision(&broker, "session");
4179                let run_id = requester.run_id.clone();
4180                std::thread::spawn(move || {
4181                    barrier.wait();
4182                    broker.resolve(
4183                        &request_id,
4184                        &authority,
4185                        PermissionAction::Approve,
4186                        Some(GrantScope::ChildRunSameTool {
4187                            run_id,
4188                            tool_name: "bash.spawn".into(),
4189                        }),
4190                        None,
4191                    )
4192                })
4193            };
4194            let terminator = {
4195                let flows = Arc::clone(&flows);
4196                let barrier = Arc::clone(&barrier);
4197                let run_id = requester.run_id.clone();
4198                std::thread::spawn(move || {
4199                    barrier.wait();
4200                    flows.mark_terminal(&run_id);
4201                })
4202            };
4203
4204            let approved = approver.join().unwrap();
4205            terminator.join().unwrap();
4206
4207            // Both interleavings are legal, but a terminal run must never retain
4208            // authorization: grants are revoked either way.
4209            assert!(broker.grants().is_empty());
4210            match approved {
4211                Ok(ResolveOutcome::Resolved(_)) => assert!(matches!(
4212                    pending.resolution.blocking_recv().unwrap(),
4213                    PermissionResolution::Decision(_)
4214                )),
4215                Err(PermissionError::AlreadyResolved) => assert!(matches!(
4216                    pending.resolution.blocking_recv().unwrap(),
4217                    PermissionResolution::Cancelled { .. }
4218                )),
4219                other => panic!("unexpected outcome: {other:?}"),
4220            }
4221        }
4222    }
4223
4224    #[test]
4225    fn s8_deterministic_approve_terminal_interleavings_preserve_state_and_audit_order() {
4226        enum Schedule {
4227            ApproveThenTerminal,
4228            TerminalThenApprove,
4229        }
4230        for schedule in [Schedule::ApproveThenTerminal, Schedule::TerminalThenApprove] {
4231            let flows = Arc::new(FlowRegistry::default());
4232            let requester = register_root(&flows, "session", false);
4233            let broker = PermissionBroker::shared(Arc::clone(&flows));
4234            let mut audit = attach_audit(&broker);
4235            let pending = submit_to_user(&broker, &requester);
4236            drain_audit(&mut audit);
4237            let request_id = pending.request.request_id.clone();
4238            let scope = GrantScope::ChildRunSameTool {
4239                run_id: requester.run_id.clone(),
4240                tool_name: pending.request.intent.tool_name.clone(),
4241            };
4242            let (first_done_tx, first_done_rx) = std::sync::mpsc::channel();
4243            let (continue_tx, continue_rx) = std::sync::mpsc::channel();
4244
4245            let approver = {
4246                let broker = Arc::clone(&broker);
4247                let authority = user_decision(&broker, "session");
4248                let first_done_tx = first_done_tx.clone();
4249                let approve_first = matches!(schedule, Schedule::ApproveThenTerminal);
4250                std::thread::spawn(move || {
4251                    if !approve_first {
4252                        continue_rx.recv().unwrap();
4253                    }
4254                    let outcome = broker.resolve(
4255                        &request_id,
4256                        &authority,
4257                        PermissionAction::Approve,
4258                        Some(scope),
4259                        Some("accepted".into()),
4260                    );
4261                    if approve_first {
4262                        first_done_tx.send(()).unwrap();
4263                    }
4264                    outcome
4265                })
4266            };
4267            let terminator = {
4268                let flows = Arc::clone(&flows);
4269                let run_id = requester.run_id.clone();
4270                let first_done_tx = first_done_tx;
4271                let approve_first = matches!(schedule, Schedule::ApproveThenTerminal);
4272                std::thread::spawn(move || {
4273                    if approve_first {
4274                        first_done_rx.recv().unwrap();
4275                    }
4276                    flows.mark_terminal(&run_id);
4277                    if !approve_first {
4278                        first_done_tx.send(()).unwrap();
4279                        continue_tx.send(()).unwrap();
4280                    }
4281                })
4282            };
4283
4284            let approval = approver.join().unwrap();
4285            terminator.join().unwrap();
4286            assert!(broker.grants().is_empty());
4287            let frames = drain_audit(&mut audit);
4288            match schedule {
4289                Schedule::ApproveThenTerminal => {
4290                    assert!(matches!(approval, Ok(ResolveOutcome::Resolved(_))));
4291                    assert!(matches!(
4292                        pending.resolution.blocking_recv().unwrap(),
4293                        PermissionResolution::Decision(_)
4294                    ));
4295                    assert!(matches!(
4296                        broker.get(&pending.request.request_id).unwrap().state,
4297                        PermissionRequestState::Approved { .. }
4298                    ));
4299                    assert!(matches!(
4300                        frames.as_slice(),
4301                        [
4302                            StreamFrame::PermissionRequestApproved { .. },
4303                            StreamFrame::PermissionGrantCreated { .. },
4304                            StreamFrame::PermissionGrantExpired { .. }
4305                        ]
4306                    ));
4307                }
4308                Schedule::TerminalThenApprove => {
4309                    assert!(matches!(approval, Err(PermissionError::AlreadyResolved)));
4310                    assert!(matches!(
4311                        pending.resolution.blocking_recv().unwrap(),
4312                        PermissionResolution::Cancelled { .. }
4313                    ));
4314                    assert!(matches!(
4315                        broker.get(&pending.request.request_id).unwrap().state,
4316                        PermissionRequestState::Cancelled { .. }
4317                    ));
4318                    assert!(matches!(
4319                        frames.as_slice(),
4320                        [StreamFrame::PermissionRequestCancelled { .. }]
4321                    ));
4322                }
4323            }
4324        }
4325    }
4326
4327    #[test]
4328    fn s8_audit_dispatch_preserves_linearized_bundle_order_when_callers_arrive_out_of_order() {
4329        let flows = Arc::new(FlowRegistry::default());
4330        let requester = register_root(&flows, "session", false);
4331        let broker = Arc::new(PermissionBroker::new(flows));
4332        let mut audit = attach_audit(&broker);
4333        let entered = Arc::new(std::sync::Barrier::new(2));
4334        let release = Arc::new(std::sync::Barrier::new(2));
4335        broker.set_audit_dispatch_hook({
4336            let entered = Arc::clone(&entered);
4337            let release = Arc::clone(&release);
4338            Arc::new(move |sequence| {
4339                if sequence == 0 {
4340                    entered.wait();
4341                    release.wait();
4342                }
4343            })
4344        });
4345
4346        let first = {
4347            let broker = Arc::clone(&broker);
4348            let requester = Arc::clone(&requester);
4349            std::thread::spawn(move || submit_to_user(&broker, &requester))
4350        };
4351        entered.wait();
4352        let second = submit_to_user(&broker, &requester);
4353        release.wait();
4354        let first = first.join().unwrap();
4355
4356        let frames = drain_audit(&mut audit);
4357        let expected = [first.request.request_id, second.request.request_id];
4358        let actual = frames
4359            .iter()
4360            .map(|frame| match frame {
4361                StreamFrame::PermissionRequestCreated { payload, .. } => {
4362                    payload.request_id.clone().unwrap()
4363                }
4364                _ => panic!("unexpected frame: {frame:?}"),
4365            })
4366            .collect::<Vec<_>>();
4367        assert_eq!(actual, expected);
4368        assert!(matches!(
4369            frames.as_slice(),
4370            [
4371                StreamFrame::PermissionRequestCreated { .. },
4372                StreamFrame::PermissionRequestCreated { .. }
4373            ]
4374        ));
4375    }
4376
4377    #[test]
4378    fn terminal_observer_does_not_keep_broker_alive() {
4379        let flows = Arc::new(FlowRegistry::default());
4380        let requester = register_root(&flows, "session", false);
4381        let broker = PermissionBroker::shared(Arc::clone(&flows));
4382        let weak = Arc::downgrade(&broker);
4383        drop(broker);
4384
4385        assert!(weak.upgrade().is_none());
4386        // A terminal transition after the broker is gone must not panic on the dead weak.
4387        flows.mark_terminal(&requester.run_id);
4388    }
4389
4390    #[test]
4391    fn unrestricted_still_enforces_authority_and_target() {
4392        let flows = Arc::new(FlowRegistry::default());
4393        let root = register_root(&flows, "session", true);
4394        let requester = child(&flows, &root);
4395        let broker = PermissionBroker::new(Arc::clone(&flows));
4396        let policy = TrustConfig {
4397            mode: TrustMode::Reckless,
4398            ..TrustConfig::default()
4399        };
4400
4401        let mut beyond = intent();
4402        beyond.tier = Tier::Four;
4403        beyond.risks.insert(RiskKind::Irreversible);
4404        assert!(matches!(
4405            broker.submit(
4406                Some(&requester.session_id),
4407                Some(&requester.run_id),
4408                beyond,
4409                true,
4410                &policy,
4411            ),
4412            Ok(SubmissionOutcome::Immediate(value)) if matches!(value.authorization, ImmediateAuthorization::Auto { .. })
4413        ));
4414        assert!(submit_to_flow(&broker, &requester, Arc::clone(&requester)).is_err());
4415        assert_eq!(broker.list().len(), 1);
4416    }
4417
4418    #[test]
4419    fn resolve_rejects_actor_that_is_not_the_current_target() {
4420        let flows = Arc::new(FlowRegistry::default());
4421        let root = register_root(&flows, "session", true);
4422        let requester = child(&flows, &root);
4423        let broker = PermissionBroker::new(Arc::clone(&flows));
4424        let SubmissionOutcome::Pending(pending) =
4425            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
4426        else {
4427            panic!("expected pending request");
4428        };
4429
4430        assert!(matches!(
4431            broker.resolve(
4432                &pending.request.request_id,
4433                &user_decision(&broker, "session"),
4434                PermissionAction::Approve,
4435                None,
4436                None,
4437            ),
4438            Err(PermissionError::ActorNotAuthorized)
4439        ));
4440        assert!(matches!(
4441            broker.get(&pending.request.request_id).unwrap().state,
4442            PermissionRequestState::Pending {
4443                target: ApprovalTarget::Flow(_)
4444            }
4445        ));
4446    }
4447
4448    #[test]
4449    fn same_path_grant_scope_stays_unsupported_for_unbound_resources() {
4450        let flows = Arc::new(FlowRegistry::default());
4451        let requester = register_root(&flows, "session", false);
4452        let broker = PermissionBroker::new(flows);
4453        let pending = submit_user(&broker, &requester, intent());
4454
4455        assert!(matches!(
4456            broker.resolve(
4457                &pending.request.request_id,
4458                &user_decision(&broker, "session"),
4459                PermissionAction::Approve,
4460                Some(GrantScope::ChildRunSamePathRule {
4461                    run_id: requester.run_id.clone(),
4462                    tool_name: "bash.spawn".into(),
4463                    workspace_relative_path: "src/lib.rs".into(),
4464                }),
4465                None,
4466            ),
4467            Err(PermissionError::UnsupportedGrantScope)
4468        ));
4469        assert!(broker.grants().is_empty());
4470    }
4471
4472    #[test]
4473    fn grant_scope_cannot_target_another_run_or_tool() {
4474        let flows = Arc::new(FlowRegistry::default());
4475        let requester = register_root(&flows, "session", false);
4476        let other = register_root(&flows, "session", false);
4477        let broker = PermissionBroker::new(flows);
4478        let authority = user_decision(&broker, "session");
4479
4480        for scope in [
4481            GrantScope::ChildRunSameTool {
4482                run_id: other.run_id.clone(),
4483                tool_name: "bash.spawn".into(),
4484            },
4485            GrantScope::ChildRunSameTool {
4486                run_id: requester.run_id.clone(),
4487                tool_name: "fs.write".into(),
4488            },
4489        ] {
4490            let pending = submit_user(&broker, &requester, intent());
4491            assert!(matches!(
4492                broker.resolve(
4493                    &pending.request.request_id,
4494                    &authority,
4495                    PermissionAction::Approve,
4496                    Some(scope),
4497                    None,
4498                ),
4499                Err(PermissionError::GrantExceedsAuthority)
4500            ));
4501        }
4502        assert!(broker.grants().is_empty());
4503    }
4504
4505    #[test]
4506    fn concurrent_resolvers_have_exactly_one_winner() {
4507        let flows = Arc::new(FlowRegistry::default());
4508        let requester = register_root(&flows, "session", false);
4509        let broker = Arc::new(PermissionBroker::new(flows));
4510        let pending = submit_user(&broker, &requester, intent());
4511        let request_id = pending.request.request_id.clone();
4512        let threads: Vec<_> = [PermissionAction::Approve, PermissionAction::Deny]
4513            .into_iter()
4514            .map(|action| {
4515                let broker = Arc::clone(&broker);
4516                let request_id = request_id.clone();
4517                let authority = user_decision(&broker, "session");
4518                std::thread::spawn(move || {
4519                    broker.resolve(&request_id, &authority, action, None, None)
4520                })
4521            })
4522            .collect();
4523        let outcomes: Vec<_> = threads
4524            .into_iter()
4525            .map(|thread| thread.join().unwrap())
4526            .collect();
4527
4528        assert_eq!(outcomes.iter().filter(|outcome| outcome.is_ok()).count(), 1);
4529        assert_eq!(
4530            outcomes
4531                .iter()
4532                .filter(|outcome| matches!(outcome, Err(PermissionError::AlreadyResolved)))
4533                .count(),
4534            1
4535        );
4536        assert!(matches!(
4537            pending.resolution.blocking_recv().unwrap(),
4538            PermissionResolution::Decision(_)
4539        ));
4540    }
4541    #[test]
4542    fn root_capability_does_not_amplify_child_authority() {
4543        let trust = TrustConfig::default();
4544        let root = EffectiveAuthority::root(&trust, false, None);
4545        assert!(root.permission_management);
4546        let restricted = EffectiveAuthority {
4547            permission_management: false,
4548            ..root.clone()
4549        };
4550        let child = root.for_child(&restricted, false, None).unwrap();
4551        assert!(!child.permission_management);
4552    }
4553
4554    #[test]
4555    fn nearest_target_skips_terminal_ancestor() {
4556        let flows = Arc::new(FlowRegistry::default());
4557        let root = register_root(&flows, "session", true);
4558        let middle = child(&flows, &root);
4559        let requester = child(&flows, &middle);
4560        flows.mark_terminal(&middle.run_id);
4561        let broker = PermissionBroker::new(Arc::clone(&flows));
4562        let requirement = AuthorityRequirement::from_intent(&intent(), false);
4563        assert_eq!(
4564            broker.next_eligible_target(&requester, &requirement, &intent().provenance, None,),
4565            ApprovalTarget::Flow(root.run_id.clone())
4566        );
4567    }
4568
4569    #[test]
4570    fn defer_records_flow_actor_and_moves_to_next_hop() {
4571        let flows = Arc::new(FlowRegistry::default());
4572        let root = register_root(&flows, "session", true);
4573        let requester = child(&flows, &root);
4574        let broker = PermissionBroker::new(Arc::clone(&flows));
4575        let mut audit = attach_audit(&broker);
4576        let SubmissionOutcome::Pending(pending) =
4577            submit_to_flow(&broker, &requester, root.clone()).unwrap()
4578        else {
4579            panic!("expected pending");
4580        };
4581        drain_audit(&mut audit);
4582        let authority = DecisionAuthority::Flow(broker.flow_authority(root.clone()).unwrap());
4583        assert!(matches!(
4584            broker.resolve(
4585                &pending.request.request_id,
4586                &authority,
4587                PermissionAction::Defer,
4588                None,
4589                Some("not mine".into())
4590            ),
4591            Ok(ResolveOutcome::Deferred(_))
4592        ));
4593        let request = broker.get(&pending.request.request_id).unwrap();
4594        assert!(
4595            matches!(request.escalation_path[1].actor, Some(DecisionActor::Flow { ref run_id, .. }) if run_id == &root.run_id)
4596        );
4597        assert!(matches!(
4598            request.state,
4599            PermissionRequestState::Pending {
4600                target: ApprovalTarget::User
4601            }
4602        ));
4603        let frames = drain_audit(&mut audit);
4604        assert!(matches!(
4605            frames.as_slice(),
4606            [
4607                StreamFrame::PermissionRequestDeferred { payload: deferred, .. },
4608                StreamFrame::PermissionRequestTargeted { payload: targeted, .. }
4609            ] if deferred.actor == Some(crate::permission_audit::PermissionProjectionActor::Flow {
4610                session_id: root.session_id.clone(),
4611                run_id: root.run_id.clone(),
4612            }) && targeted.actor.is_none()
4613                && targeted.target == crate::permission_audit::PermissionAuditTarget::User
4614        ));
4615    }
4616
4617    #[test]
4618    fn stale_timeout_cannot_retarget_resolved_request() {
4619        let flows = Arc::new(FlowRegistry::default());
4620        let root = register_root(&flows, "session", true);
4621        let requester = child(&flows, &root);
4622        let broker = PermissionBroker::new(Arc::clone(&flows));
4623        let SubmissionOutcome::Pending(pending) =
4624            submit_to_flow(&broker, &requester, root.clone()).unwrap()
4625        else {
4626            panic!("expected pending");
4627        };
4628        broker
4629            .resolve(
4630                &pending.request.request_id,
4631                &DecisionAuthority::Flow(broker.flow_authority(root).unwrap()),
4632                PermissionAction::Approve,
4633                None,
4634                None,
4635            )
4636            .unwrap();
4637        assert!(
4638            !broker
4639                .defer_timed_out_target(&pending.request.request_id, &requester.run_id)
4640                .unwrap()
4641        );
4642        assert!(matches!(
4643            broker.get(&pending.request.request_id).unwrap().state,
4644            PermissionRequestState::Approved { .. }
4645        ));
4646    }
4647
4648    #[test]
4649    fn terminal_target_retargets_pending_request() {
4650        let flows = Arc::new(FlowRegistry::default());
4651        let root = register_root(&flows, "session", true);
4652        let middle = child(&flows, &root);
4653        let requester = child(&flows, &middle);
4654        let broker = PermissionBroker::shared(Arc::clone(&flows));
4655        let mut audit = attach_audit(&broker);
4656        let SubmissionOutcome::Pending(pending) =
4657            submit_to_flow(&broker, &requester, middle.clone()).unwrap()
4658        else {
4659            panic!("expected pending");
4660        };
4661        drain_audit(&mut audit);
4662        flows.mark_terminal(&middle.run_id);
4663        assert_eq!(
4664            pending.target_changes.borrow().clone(),
4665            ApprovalTarget::Flow(root.run_id.clone())
4666        );
4667        assert_eq!(
4668            broker.get(&pending.request.request_id).unwrap().state,
4669            PermissionRequestState::Pending {
4670                target: ApprovalTarget::Flow(root.run_id.clone())
4671            }
4672        );
4673        assert!(matches!(
4674            drain_audit(&mut audit).as_slice(),
4675            [
4676                StreamFrame::PermissionRequestDeferred { payload: deferred, .. },
4677                StreamFrame::PermissionRequestTargeted { payload: targeted, .. }
4678            ] if deferred.request_id == targeted.request_id
4679                && targeted.target == crate::permission_audit::PermissionAuditTarget::Flow {
4680                    run_id: root.run_id.clone(),
4681                }
4682        ));
4683    }
4684
4685    #[test]
4686    fn s8_timeout_defer_emits_deferred_then_targeted_from_real_broker() {
4687        let flows = Arc::new(FlowRegistry::default());
4688        let root = register_root(&flows, "session", true);
4689        let middle = child(&flows, &root);
4690        let requester = child(&flows, &middle);
4691        let broker = PermissionBroker::new(flows);
4692        let mut audit = attach_audit(&broker);
4693        let SubmissionOutcome::Pending(pending) =
4694            submit_to_flow(&broker, &requester, middle.clone()).unwrap()
4695        else {
4696            panic!("expected pending");
4697        };
4698        drain_audit(&mut audit);
4699        assert!(
4700            broker
4701                .defer_timed_out_target(&pending.request.request_id, &middle.run_id)
4702                .unwrap()
4703        );
4704        assert!(matches!(
4705            drain_audit(&mut audit).as_slice(),
4706            [
4707                StreamFrame::PermissionRequestDeferred { payload: deferred, .. },
4708                StreamFrame::PermissionRequestTargeted { payload: targeted, .. }
4709            ] if deferred.request_id == targeted.request_id
4710                && targeted.target == crate::permission_audit::PermissionAuditTarget::Flow {
4711                    run_id: root.run_id.clone(),
4712                }
4713        ));
4714    }
4715
4716    #[test]
4717    fn visible_apis_reject_unknown_forged_terminal_blocked_and_unprivileged_actors() {
4718        let cases = ["unknown", "forged", "terminal", "blocked", "unprivileged"];
4719        for case in cases {
4720            let flows = Arc::new(FlowRegistry::default());
4721            let actor = register_root(&flows, "session", case != "unprivileged");
4722            let requester = child(&flows, &actor);
4723            let broker = PermissionBroker::new(Arc::clone(&flows));
4724            let pending = if case == "unprivileged" {
4725                submit_user(&broker, &requester, intent())
4726            } else {
4727                let SubmissionOutcome::Pending(pending) =
4728                    submit_to_flow(&broker, &requester, Arc::clone(&actor)).unwrap()
4729                else {
4730                    panic!("expected pending");
4731                };
4732                *pending
4733            };
4734            let tested_actor = match case {
4735                "unknown" => Arc::new(FlowIdentity {
4736                    session_id: actor.session_id.clone(),
4737                    run_id: FlowRunId::now(),
4738                    parent_run_id: None,
4739                    root_run_id: actor.root_run_id.clone(),
4740                    invocation: InvocationKind::Root,
4741                    effective_authority: authority(true),
4742                    execution_state: Mutex::new(FlowExecutionState::Running),
4743                }),
4744                "forged" => Arc::new(FlowIdentity {
4745                    session_id: actor.session_id.clone(),
4746                    run_id: actor.run_id.clone(),
4747                    parent_run_id: actor.parent_run_id.clone(),
4748                    root_run_id: actor.root_run_id.clone(),
4749                    invocation: actor.invocation,
4750                    effective_authority: actor.effective_authority.clone(),
4751                    execution_state: Mutex::new(FlowExecutionState::Running),
4752                }),
4753                "terminal" => {
4754                    flows.mark_terminal(&actor.run_id);
4755                    Arc::clone(&actor)
4756                }
4757                "blocked" => {
4758                    *actor.execution_state.lock().unwrap() =
4759                        FlowExecutionState::BlockedOnDescendants {
4760                            child_run_counts: std::collections::HashMap::from([(
4761                                requester.run_id.clone(),
4762                                1,
4763                            )]),
4764                        };
4765                    Arc::clone(&actor)
4766                }
4767                "unprivileged" => Arc::clone(&actor),
4768                _ => unreachable!(),
4769            };
4770
4771            assert!(broker.visible_list(&tested_actor).is_err(), "{case}");
4772            assert!(
4773                broker
4774                    .visible_get(&tested_actor, &pending.request.request_id)
4775                    .is_err(),
4776                "{case}"
4777            );
4778        }
4779    }
4780
4781    #[test]
4782    fn visible_apis_only_return_pending_requests_targeted_to_the_actor() {
4783        let flows = Arc::new(FlowRegistry::default());
4784        let root = register_root(&flows, "session", true);
4785        let middle = child(&flows, &root);
4786        let requester = child(&flows, &middle);
4787        let broker = PermissionBroker::new(Arc::clone(&flows));
4788        let SubmissionOutcome::Pending(pending) =
4789            submit_to_flow(&broker, &requester, Arc::clone(&middle)).unwrap()
4790        else {
4791            panic!("expected pending");
4792        };
4793
4794        assert_eq!(broker.visible_list(&middle).unwrap().len(), 1);
4795        assert!(
4796            broker
4797                .visible_get(&root, &pending.request.request_id)
4798                .unwrap()
4799                .is_none()
4800        );
4801    }
4802
4803    #[test]
4804    fn flow_groups_are_owner_scoped_and_delete_only_when_empty() {
4805        let flows = Arc::new(FlowRegistry::default());
4806        let root = register_root(&flows, "session", true);
4807        let sibling_owner = child(&flows, &root);
4808        let requester = child(&flows, &root);
4809        let broker = PermissionBroker::new(Arc::clone(&flows));
4810        let SubmissionOutcome::Pending(pending) =
4811            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
4812        else {
4813            panic!("expected pending");
4814        };
4815        let request_id = pending.request.request_id.clone();
4816        let group = broker
4817            .create_group(&root, BTreeSet::from([request_id.clone()]), "review".into())
4818            .unwrap();
4819
4820        assert_eq!(
4821            broker.visible_group_list(&root).unwrap(),
4822            vec![group.clone()]
4823        );
4824        assert!(
4825            broker
4826                .visible_group_get(&sibling_owner, &group.group_id)
4827                .unwrap()
4828                .is_none()
4829        );
4830        assert!(matches!(
4831            broker.delete_empty_group(&root, &group.group_id),
4832            Err(PermissionError::GroupNotEmpty)
4833        ));
4834        let emptied = broker
4835            .ungroup_requests(&root, &group.group_id, &BTreeSet::from([request_id]))
4836            .unwrap();
4837        assert!(emptied.request_ids.is_empty());
4838        assert_eq!(emptied.revision, 1);
4839        broker.delete_empty_group(&root, &group.group_id).unwrap();
4840        assert!(broker.visible_group_list(&root).unwrap().is_empty());
4841    }
4842
4843    #[test]
4844    fn group_creation_rejects_hidden_and_unknown_request_ids() {
4845        let flows = Arc::new(FlowRegistry::default());
4846        let root = register_root(&flows, "session", true);
4847        let middle = child(&flows, &root);
4848        let requester = child(&flows, &middle);
4849        let broker = PermissionBroker::new(Arc::clone(&flows));
4850        let SubmissionOutcome::Pending(pending) =
4851            submit_to_flow(&broker, &requester, Arc::clone(&middle)).unwrap()
4852        else {
4853            panic!("expected pending");
4854        };
4855
4856        for request_id in [pending.request.request_id, PermissionRequestId::now()] {
4857            assert!(matches!(
4858                broker.create_group(&root, BTreeSet::from([request_id]), "hidden".into()),
4859                Err(PermissionError::ActorNotAuthorized)
4860            ));
4861        }
4862        assert!(broker.visible_group_list(&root).unwrap().is_empty());
4863    }
4864
4865    #[test]
4866    fn delegated_workspace_ancestor_is_skipped_when_resource_is_outside_its_root() {
4867        let temp = tempfile::tempdir().unwrap();
4868        let root_path = temp.path().join("root");
4869        let delegated_path = temp.path().join("delegated");
4870        std::fs::create_dir_all(&root_path).unwrap();
4871        std::fs::create_dir_all(&delegated_path).unwrap();
4872        let flows = Arc::new(FlowRegistry::default());
4873        let root = flows
4874            .register_root(
4875                "session".into(),
4876                FlowRunId::now(),
4877                EffectiveAuthority {
4878                    workspace_root: Some(crate::fs_access::canonicalize_stable(&root_path)),
4879                    ..authority(true)
4880                },
4881            )
4882            .unwrap();
4883        let middle = flows
4884            .register_child(
4885                &root.run_id,
4886                FlowRunId::now(),
4887                InvocationKind::InlineSubflow,
4888                true,
4889                ChildWorkspaceAuthority::TrustedDelegation(delegated_path.clone()),
4890            )
4891            .unwrap();
4892        let requester = child(&flows, &middle);
4893        let broker = PermissionBroker::new(Arc::clone(&flows));
4894        let delegated_intent = PermissionIntent {
4895            provenance: ResourceProvenance {
4896                path: Some(delegated_path.join("file.txt")),
4897                path_origin: Some(PathOrigin::ExplicitInside),
4898                workspace_root: Some(delegated_path),
4899                ..ResourceProvenance::default()
4900            },
4901            ..intent()
4902        };
4903        let SubmissionOutcome::Pending(pending) = broker
4904            .submit(
4905                Some(&requester.session_id),
4906                Some(&requester.run_id),
4907                delegated_intent,
4908                false,
4909                &ask_policy(),
4910            )
4911            .unwrap()
4912        else {
4913            panic!("expected pending");
4914        };
4915        assert_eq!(
4916            pending.request.state,
4917            PermissionRequestState::Pending {
4918                target: ApprovalTarget::Flow(middle.run_id.clone())
4919            }
4920        );
4921        broker
4922            .resolve(
4923                &pending.request.request_id,
4924                &DecisionAuthority::Flow(broker.flow_authority(middle).unwrap()),
4925                PermissionAction::Defer,
4926                None,
4927                None,
4928            )
4929            .unwrap();
4930        assert_eq!(
4931            broker.get(&pending.request.request_id).unwrap().state,
4932            PermissionRequestState::Pending {
4933                target: ApprovalTarget::User
4934            }
4935        );
4936    }
4937
4938    #[test]
4939    fn delegated_requester_cannot_submit_resource_outside_delegated_workspace() {
4940        let temp = tempfile::tempdir().unwrap();
4941        let parent_path = temp.path().join("parent");
4942        let delegated_path = temp.path().join("delegated");
4943        std::fs::create_dir_all(&parent_path).unwrap();
4944        std::fs::create_dir_all(&delegated_path).unwrap();
4945        let parent_root = crate::fs_access::canonicalize_stable(&parent_path);
4946        let delegated_root = crate::fs_access::canonicalize_stable(&delegated_path);
4947        let flows = Arc::new(FlowRegistry::default());
4948        let parent = flows
4949            .register_root(
4950                "session".into(),
4951                FlowRunId::now(),
4952                EffectiveAuthority {
4953                    workspace_root: Some(parent_root.clone()),
4954                    ..authority(true)
4955                },
4956            )
4957            .unwrap();
4958        let delegated = flows
4959            .register_child(
4960                &parent.run_id,
4961                FlowRunId::now(),
4962                InvocationKind::InlineSubflow,
4963                true,
4964                ChildWorkspaceAuthority::TrustedDelegation(delegated_root.clone()),
4965            )
4966            .unwrap();
4967        let broker = PermissionBroker::new(Arc::clone(&flows));
4968        let outside_intent = PermissionIntent {
4969            provenance: ResourceProvenance {
4970                path: Some(parent_root.join("parent.txt")),
4971                path_origin: Some(PathOrigin::ExplicitInside),
4972                workspace_root: Some(parent_root),
4973                ..ResourceProvenance::default()
4974            },
4975            ..intent()
4976        };
4977        let policies = [
4978            ask_policy(),
4979            TrustConfig {
4980                mode: TrustMode::Eager,
4981                escalation: EscalationPolicy::Allow,
4982                ..TrustConfig::default()
4983            },
4984            TrustConfig {
4985                mode: TrustMode::Reckless,
4986                ..TrustConfig::default()
4987            },
4988        ];
4989        for policy in policies {
4990            assert!(matches!(
4991                broker.submit(
4992                    Some(&delegated.session_id),
4993                    Some(&delegated.run_id),
4994                    outside_intent.clone(),
4995                    false,
4996                    &policy,
4997                ),
4998                Err(PermissionError::GrantExceedsAuthority)
4999            ));
5000        }
5001        assert!(broker.list().is_empty());
5002        assert!(broker.grants().is_empty());
5003    }
5004
5005    #[test]
5006    fn automatic_and_denied_submissions_do_not_require_approval_target_selection() {
5007        let flows = Arc::new(FlowRegistry::default());
5008        let root = register_root(&flows, "session", true);
5009        let requester = child(&flows, &root);
5010        flows.mark_terminal(&root.run_id);
5011        let broker = PermissionBroker::new(flows);
5012        for policy in [
5013            TrustConfig {
5014                mode: TrustMode::Eager,
5015                escalation: EscalationPolicy::Allow,
5016                ..TrustConfig::default()
5017            },
5018            TrustConfig {
5019                mode: TrustMode::Eager,
5020                escalation: EscalationPolicy::Deny,
5021                ..TrustConfig::default()
5022            },
5023        ] {
5024            assert!(matches!(
5025                broker.submit(
5026                    Some(&requester.session_id),
5027                    Some(&requester.run_id),
5028                    intent(),
5029                    false,
5030                    &policy,
5031                ),
5032                Ok(SubmissionOutcome::Immediate(_))
5033            ));
5034        }
5035    }
5036
5037    #[test]
5038    fn ask_target_selection_is_lifecycle_arbitrated_with_terminal_transition() {
5039        for _ in 0..64 {
5040            let flows = Arc::new(FlowRegistry::default());
5041            let root = register_root(&flows, "session", true);
5042            let requester = child(&flows, &root);
5043            let broker = PermissionBroker::shared(Arc::clone(&flows));
5044            let barrier = Arc::new(std::sync::Barrier::new(2));
5045            let submitter = {
5046                let broker = Arc::clone(&broker);
5047                let barrier = Arc::clone(&barrier);
5048                let requester = Arc::clone(&requester);
5049                std::thread::spawn(move || {
5050                    barrier.wait();
5051                    broker.submit(
5052                        Some(&requester.session_id),
5053                        Some(&requester.run_id),
5054                        intent(),
5055                        false,
5056                        &ask_policy(),
5057                    )
5058                })
5059            };
5060            let terminator = {
5061                let flows = Arc::clone(&flows);
5062                let barrier = Arc::clone(&barrier);
5063                let root_run_id = root.run_id.clone();
5064                std::thread::spawn(move || {
5065                    barrier.wait();
5066                    flows.mark_terminal(&root_run_id);
5067                })
5068            };
5069
5070            let outcome = submitter.join().unwrap().unwrap();
5071            terminator.join().unwrap();
5072            let SubmissionOutcome::Pending(pending) = outcome else {
5073                panic!("Ask must remain pending");
5074            };
5075            assert_eq!(
5076                broker.get(&pending.request.request_id).unwrap().state,
5077                PermissionRequestState::Pending {
5078                    target: ApprovalTarget::User
5079                }
5080            );
5081        }
5082    }
5083
5084    #[test]
5085    fn timeout_is_per_hop_stale_safe_and_user_is_terminal_fallback() {
5086        let flows = Arc::new(FlowRegistry::default());
5087        let root = register_root(&flows, "session", true);
5088        let middle = child(&flows, &root);
5089        let requester = child(&flows, &middle);
5090        let broker = PermissionBroker::new(Arc::clone(&flows));
5091        let pending = match broker
5092            .submit(
5093                Some(&requester.session_id),
5094                Some(&requester.run_id),
5095                intent(),
5096                false,
5097                &ask_policy(),
5098            )
5099            .unwrap()
5100        {
5101            SubmissionOutcome::Pending(pending) => pending,
5102            _ => panic!("expected pending"),
5103        };
5104        assert!(
5105            broker
5106                .defer_timed_out_target(&pending.request.request_id, &middle.run_id)
5107                .unwrap()
5108        );
5109        assert_eq!(
5110            pending.target_changes.borrow().clone(),
5111            ApprovalTarget::Flow(root.run_id.clone())
5112        );
5113        assert!(
5114            !broker
5115                .defer_timed_out_target(&pending.request.request_id, &middle.run_id)
5116                .unwrap()
5117        );
5118        assert!(
5119            broker
5120                .defer_timed_out_target(&pending.request.request_id, &root.run_id)
5121                .unwrap()
5122        );
5123        assert_eq!(
5124            pending.target_changes.borrow().clone(),
5125            ApprovalTarget::User
5126        );
5127        assert!(
5128            !broker
5129                .defer_timed_out_target(&pending.request.request_id, &root.run_id)
5130                .unwrap()
5131        );
5132    }
5133
5134    #[test]
5135    fn persistent_same_path_grant_matches_only_same_run_tool_path_and_narrower_risk() {
5136        let temp = tempfile::tempdir().unwrap();
5137        let root_path = crate::fs_access::canonicalize_stable(temp.path());
5138        let file = root_path.join("file.txt");
5139        let other = root_path.join("other.txt");
5140        let flows = Arc::new(FlowRegistry::default());
5141        let requester = register_root(&flows, "session", false);
5142        let other_run = register_root(&flows, "session", false);
5143        let broker = PermissionBroker::new(flows);
5144        let approved_intent = PermissionIntent {
5145            risks: BTreeSet::from([RiskKind::FilesystemWrite]),
5146            provenance: ResourceProvenance {
5147                path: Some(file.clone()),
5148                path_origin: Some(PathOrigin::ExplicitInside),
5149                workspace_root: Some(root_path.clone()),
5150                risks: BTreeSet::from([RiskKind::FilesystemWrite]),
5151                ..ResourceProvenance::default()
5152            },
5153            ..intent()
5154        };
5155        let pending = submit_user(&broker, &requester, approved_intent.clone());
5156        broker
5157            .resolve(
5158                &pending.request.request_id,
5159                &user_decision(&broker, "session"),
5160                PermissionAction::Approve,
5161                Some(GrantScope::ChildRunSamePathRule {
5162                    run_id: requester.run_id.clone(),
5163                    tool_name: "bash.spawn".into(),
5164                    workspace_relative_path: "file.txt".into(),
5165                }),
5166                None,
5167            )
5168            .unwrap();
5169
5170        let mut narrower = approved_intent.clone();
5171        narrower.risks.clear();
5172        narrower.provenance.risks.clear();
5173        assert!(
5174            broker
5175                .find_matching_grant(&requester, &narrower, false)
5176                .is_some()
5177        );
5178        let mut wrong_tool = narrower.clone();
5179        wrong_tool.tool_name = "fs.write".into();
5180        assert!(
5181            broker
5182                .find_matching_grant(&requester, &wrong_tool, false)
5183                .is_none()
5184        );
5185        let mut wrong_path = narrower.clone();
5186        wrong_path.provenance.path = Some(other);
5187        assert!(
5188            broker
5189                .find_matching_grant(&requester, &wrong_path, false)
5190                .is_none()
5191        );
5192        assert!(
5193            broker
5194                .find_matching_grant(&other_run, &narrower, false)
5195                .is_none()
5196        );
5197        let mut escalated = approved_intent;
5198        escalated.risks.insert(RiskKind::Irreversible);
5199        assert!(
5200            broker
5201                .find_matching_grant(&requester, &escalated, false)
5202                .is_none()
5203        );
5204    }
5205
5206    #[cfg(unix)]
5207    #[test]
5208    fn persistent_same_path_grant_rejects_non_utf8_relative_path() {
5209        use std::os::unix::ffi::OsStringExt;
5210
5211        let temp = tempfile::tempdir().unwrap();
5212        let root_path = crate::fs_access::canonicalize_stable(temp.path());
5213        let path = root_path.join(std::ffi::OsString::from_vec(vec![0xff]));
5214        let flows = Arc::new(FlowRegistry::default());
5215        let requester = register_root(&flows, "session", false);
5216        let broker = PermissionBroker::new(flows);
5217        let pending = submit_user(
5218            &broker,
5219            &requester,
5220            PermissionIntent {
5221                provenance: ResourceProvenance {
5222                    path: Some(path),
5223                    path_origin: Some(PathOrigin::ExplicitInside),
5224                    workspace_root: Some(root_path),
5225                    ..ResourceProvenance::default()
5226                },
5227                ..intent()
5228            },
5229        );
5230        assert!(matches!(
5231            broker.resolve(
5232                &pending.request.request_id,
5233                &user_decision(&broker, "session"),
5234                PermissionAction::Approve,
5235                Some(GrantScope::ChildRunSamePathRule {
5236                    run_id: requester.run_id.clone(),
5237                    tool_name: "bash.spawn".into(),
5238                    workspace_relative_path: "�".into(),
5239                }),
5240                None,
5241            ),
5242            Err(PermissionError::UnsupportedGrantScope)
5243        ));
5244    }
5245
5246    #[test]
5247    fn root_authority_retains_permission_management_capability() {
5248        let root = EffectiveAuthority::root(&TrustConfig::default(), true, None);
5249        assert!(root.permission_management);
5250        let child = root
5251            .inherited_child(true, ChildWorkspaceAuthority::Inherit)
5252            .unwrap();
5253        assert!(child.permission_management);
5254    }
5255
5256    #[test]
5257    fn built_in_batch_selectors_match_stable_request_identity() {
5258        for case in 0..6 {
5259            let temp = tempfile::tempdir().unwrap();
5260            let path = temp.path().join("nested/file.txt");
5261            let flows = Arc::new(FlowRegistry::default());
5262            let root = register_root(&flows, "session", true);
5263            let requester = child(&flows, &root);
5264            let broker = PermissionBroker::new(flows);
5265            let selected_intent = PermissionIntent {
5266                risks: BTreeSet::from([RiskKind::Network]),
5267                provenance: ResourceProvenance {
5268                    path: Some(path),
5269                    ..ResourceProvenance::default()
5270                },
5271                ..intent()
5272            };
5273            let SubmissionOutcome::Pending(pending) =
5274                submit_to_flow_with_intent(&broker, &requester, Arc::clone(&root), selected_intent)
5275                    .unwrap()
5276            else {
5277                panic!("expected pending request");
5278            };
5279            let selector = match case {
5280                0 => PermissionSelector::ChildRun,
5281                1 => PermissionSelector::Tool("bash.spawn".into()),
5282                2 => PermissionSelector::Tier(Tier::Two),
5283                3 => PermissionSelector::Risk(RiskKind::Network),
5284                4 => PermissionSelector::PathPrefix(temp.path().join("nested")),
5285                5 => PermissionSelector::Target(ApprovalTarget::Flow(root.run_id.clone())),
5286                _ => unreachable!(),
5287            };
5288            let results = broker
5289                .resolve_batch(
5290                    &root,
5291                    selector,
5292                    PermissionAction::Approve,
5293                    None,
5294                    None,
5295                    BatchMode::BestEffort,
5296                    None,
5297                )
5298                .unwrap();
5299            assert_eq!(results.len(), 1);
5300            assert_eq!(results[0].request_id, pending.request.request_id);
5301            assert!(matches!(
5302                results[0].outcome,
5303                BatchRequestOutcome::Approved(_)
5304            ));
5305        }
5306    }
5307
5308    #[test]
5309    fn path_prefix_selector_matches_every_authorized_provenance_target() {
5310        let temp = tempfile::tempdir().unwrap();
5311        let root_path = crate::fs_access::canonicalize_stable(temp.path());
5312        for (case, prefix) in [
5313            (0, root_path.join("primary")),
5314            (1, root_path.join("cwd")),
5315            (2, root_path.join("extra")),
5316            (3, root_path.join("missing")),
5317        ] {
5318            let flows = Arc::new(FlowRegistry::default());
5319            let root = register_root(&flows, "session", true);
5320            let requester = child(&flows, &root);
5321            let broker = PermissionBroker::new(flows);
5322            let SubmissionOutcome::Pending(pending) = submit_to_flow_with_intent(
5323                &broker,
5324                &requester,
5325                Arc::clone(&root),
5326                PermissionIntent {
5327                    provenance: ResourceProvenance {
5328                        path: Some(root_path.join("primary/file.txt")),
5329                        cwd: Some(root_path.join("cwd/work")),
5330                        extra_targets: vec![root_path.join("extra/other.txt")],
5331                        ..ResourceProvenance::default()
5332                    },
5333                    ..intent()
5334                },
5335            )
5336            .unwrap() else {
5337                panic!("expected pending request");
5338            };
5339
5340            let results = broker
5341                .resolve_batch(
5342                    &root,
5343                    PermissionSelector::PathPrefix(prefix),
5344                    PermissionAction::Approve,
5345                    None,
5346                    None,
5347                    BatchMode::BestEffort,
5348                    None,
5349                )
5350                .unwrap();
5351            if case == 3 {
5352                assert!(results.is_empty());
5353                assert!(matches!(
5354                    broker.get(&pending.request.request_id).unwrap().state,
5355                    PermissionRequestState::Pending { .. }
5356                ));
5357            } else {
5358                assert_eq!(results.len(), 1);
5359                assert_eq!(results[0].request_id, pending.request.request_id);
5360                assert!(matches!(
5361                    results[0].outcome,
5362                    BatchRequestOutcome::Approved(_)
5363                ));
5364            }
5365        }
5366    }
5367
5368    #[test]
5369    fn best_effort_batch_preserves_prior_success_before_later_member_rejection() {
5370        let temp = tempfile::tempdir().unwrap();
5371        let root_path = crate::fs_access::canonicalize_stable(temp.path());
5372        let flows = Arc::new(FlowRegistry::default());
5373        let root = register_root(&flows, "session", true);
5374        let requester = child(&flows, &root);
5375        let broker = PermissionBroker::new(flows);
5376        let SubmissionOutcome::Pending(first) = submit_to_flow_with_intent(
5377            &broker,
5378            &requester,
5379            Arc::clone(&root),
5380            PermissionIntent {
5381                provenance: ResourceProvenance {
5382                    path: Some(root_path.join("first.txt")),
5383                    path_origin: Some(PathOrigin::ExplicitInside),
5384                    workspace_root: Some(root_path),
5385                    ..ResourceProvenance::default()
5386                },
5387                ..intent()
5388            },
5389        )
5390        .unwrap() else {
5391            panic!("expected pending request");
5392        };
5393        let SubmissionOutcome::Pending(second) =
5394            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5395        else {
5396            panic!("expected pending request");
5397        };
5398        let scope = GrantScope::ChildRunSamePathRule {
5399            run_id: requester.run_id.clone(),
5400            tool_name: "bash.spawn".into(),
5401            workspace_relative_path: "first.txt".into(),
5402        };
5403
5404        let results = broker
5405            .resolve_batch(
5406                &root,
5407                PermissionSelector::RequestIds(vec![
5408                    first.request.request_id.clone(),
5409                    second.request.request_id.clone(),
5410                ]),
5411                PermissionAction::Approve,
5412                Some(scope),
5413                None,
5414                BatchMode::BestEffort,
5415                None,
5416            )
5417            .unwrap();
5418
5419        assert_eq!(results.len(), 2);
5420        assert_eq!(results[0].request_id, first.request.request_id);
5421        assert!(matches!(
5422            results[0].outcome,
5423            BatchRequestOutcome::Approved(_)
5424        ));
5425        assert_eq!(results[1].request_id, second.request.request_id);
5426        assert!(matches!(
5427            results[1].outcome,
5428            BatchRequestOutcome::RejectedUnsupportedGrantScope
5429        ));
5430        assert!(matches!(
5431            broker.get(&first.request.request_id).unwrap().state,
5432            PermissionRequestState::Approved { .. }
5433        ));
5434        assert!(matches!(
5435            broker.get(&second.request.request_id).unwrap().state,
5436            PermissionRequestState::Pending { .. }
5437        ));
5438    }
5439
5440    #[test]
5441    fn best_effort_batch_reports_mixed_outcomes_and_independent_decisions() {
5442        let flows = Arc::new(FlowRegistry::default());
5443        let root = register_root(&flows, "session", true);
5444        let requesters = [child(&flows, &root), child(&flows, &root)];
5445        let broker = PermissionBroker::new(flows);
5446        let pending = requesters.map(|requester| {
5447            let SubmissionOutcome::Pending(pending) =
5448                submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5449            else {
5450                panic!("expected pending request");
5451            };
5452            *pending
5453        });
5454        broker
5455            .resolve(
5456                &pending[0].request.request_id,
5457                &DecisionAuthority::Flow(broker.flow_authority(Arc::clone(&root)).unwrap()),
5458                PermissionAction::Approve,
5459                None,
5460                None,
5461            )
5462            .unwrap();
5463        let missing = PermissionRequestId::now();
5464        let results = broker
5465            .resolve_batch(
5466                &root,
5467                PermissionSelector::RequestIds(vec![
5468                    pending[0].request.request_id.clone(),
5469                    pending[1].request.request_id.clone(),
5470                    missing.clone(),
5471                ]),
5472                PermissionAction::Approve,
5473                None,
5474                None,
5475                BatchMode::BestEffort,
5476                None,
5477            )
5478            .unwrap();
5479        assert!(
5480            results
5481                .iter()
5482                .any(|result| result.request_id == pending[0].request.request_id
5483                    && matches!(result.outcome, BatchRequestOutcome::SkippedAlreadyResolved))
5484        );
5485        assert!(
5486            results
5487                .iter()
5488                .any(|result| result.request_id == pending[1].request.request_id
5489                    && matches!(result.outcome, BatchRequestOutcome::Approved(_)))
5490        );
5491        assert!(results.iter().any(|result| result.request_id == missing
5492            && matches!(result.outcome, BatchRequestOutcome::RejectedNotFound)));
5493        let decision_id =
5494            |request_id: &PermissionRequestId| match broker.get(request_id).unwrap().state {
5495                PermissionRequestState::Approved { decision_id, .. } => decision_id,
5496                state => panic!("unexpected state: {state:?}"),
5497            };
5498        assert_ne!(
5499            decision_id(&pending[0].request.request_id),
5500            decision_id(&pending[1].request.request_id)
5501        );
5502    }
5503
5504    #[test]
5505    fn atomic_batch_rolls_back_and_group_revision_is_checked() {
5506        let flows = Arc::new(FlowRegistry::default());
5507        let root = register_root(&flows, "session", true);
5508        let requester = child(&flows, &root);
5509        let broker = PermissionBroker::new(flows);
5510        let SubmissionOutcome::Pending(pending) =
5511            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5512        else {
5513            panic!("expected pending request");
5514        };
5515        let request_id = pending.request.request_id.clone();
5516        assert!(matches!(
5517            broker.resolve_batch(
5518                &root,
5519                PermissionSelector::RequestIds(vec![
5520                    request_id.clone(),
5521                    PermissionRequestId::now()
5522                ]),
5523                PermissionAction::Approve,
5524                None,
5525                None,
5526                BatchMode::Atomic,
5527                None,
5528            ),
5529            Err(PermissionError::RequestNotFound)
5530        ));
5531        assert!(matches!(
5532            broker.get(&request_id).unwrap().state,
5533            PermissionRequestState::Pending { .. }
5534        ));
5535        let group = broker
5536            .create_group(&root, BTreeSet::from([request_id]), "batch".into())
5537            .unwrap();
5538        assert!(matches!(
5539            broker.resolve_batch(
5540                &root,
5541                PermissionSelector::Group(group.group_id),
5542                PermissionAction::Approve,
5543                None,
5544                None,
5545                BatchMode::Atomic,
5546                Some(group.revision + 1),
5547            ),
5548            Err(PermissionError::GroupRevisionConflict)
5549        ));
5550    }
5551
5552    #[test]
5553    fn group_selector_expands_membership_and_revision_from_the_same_state() {
5554        let flows = Arc::new(FlowRegistry::default());
5555        let root = register_root(&flows, "session", true);
5556        let requester = child(&flows, &root);
5557        let broker = PermissionBroker::new(flows);
5558        let SubmissionOutcome::Pending(first) =
5559            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5560        else {
5561            panic!("expected pending request");
5562        };
5563        let SubmissionOutcome::Pending(second) =
5564            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5565        else {
5566            panic!("expected pending request");
5567        };
5568        let group = broker
5569            .create_group(
5570                &root,
5571                BTreeSet::from([first.request.request_id.clone()]),
5572                "atomic".into(),
5573            )
5574            .unwrap();
5575        let mut state = broker.state.lock().unwrap();
5576        let current = state.groups.get_mut(&group.group_id).unwrap();
5577        current
5578            .request_ids
5579            .insert(second.request.request_id.clone());
5580        current.revision += 1;
5581        assert!(matches!(
5582            broker.expand_selector(
5583                &state,
5584                &root,
5585                &PermissionSelector::Group(group.group_id.clone()),
5586                Some(group.revision)
5587            ),
5588            Err(PermissionError::GroupRevisionConflict)
5589        ));
5590        let expected = BTreeSet::from([first.request.request_id, second.request.request_id])
5591            .into_iter()
5592            .collect::<Vec<_>>();
5593        assert_eq!(
5594            broker
5595                .expand_selector(
5596                    &state,
5597                    &root,
5598                    &PermissionSelector::Group(group.group_id),
5599                    Some(group.revision + 1)
5600                )
5601                .unwrap(),
5602            expected
5603        );
5604    }
5605
5606    #[test]
5607    fn s8_broker_group_emissions_keep_request_finals_before_single_resolved_transition() {
5608        let flows = Arc::new(FlowRegistry::default());
5609        let root = register_root(&flows, "session", true);
5610        let requesters = [
5611            child(&flows, &root),
5612            child(&flows, &root),
5613            child(&flows, &root),
5614        ];
5615        let broker = PermissionBroker::new(flows);
5616        let mut audit = attach_audit(&broker);
5617        let pending = requesters.each_ref().map(|requester| {
5618            let SubmissionOutcome::Pending(pending) =
5619                submit_to_flow(&broker, requester, Arc::clone(&root)).unwrap()
5620            else {
5621                panic!("expected pending request");
5622            };
5623            *pending
5624        });
5625        drain_audit(&mut audit);
5626        let ids = pending
5627            .each_ref()
5628            .map(|pending| pending.request.request_id.clone());
5629        let group = broker
5630            .create_group(&root, BTreeSet::from(ids.clone()), "review".into())
5631            .unwrap();
5632        broker
5633            .ungroup_requests(&root, &group.group_id, &BTreeSet::from([ids[2].clone()]))
5634            .unwrap();
5635        let current = broker
5636            .visible_group_get(&root, &group.group_id)
5637            .unwrap()
5638            .unwrap();
5639        broker
5640            .resolve_batch(
5641                &root,
5642                PermissionSelector::Group(group.group_id.clone()),
5643                PermissionAction::Approve,
5644                None,
5645                Some("accepted".into()),
5646                BatchMode::Atomic,
5647                Some(current.revision),
5648            )
5649            .unwrap();
5650
5651        let frames = drain_audit(&mut audit);
5652        assert!(matches!(
5653            frames.as_slice(),
5654            [
5655                StreamFrame::PermissionGroupCreated { .. },
5656                StreamFrame::PermissionGroupUpdated { .. },
5657                StreamFrame::PermissionRequestApproved { .. },
5658                StreamFrame::PermissionRequestApproved { .. },
5659                StreamFrame::PermissionGroupResolved { .. }
5660            ]
5661        ));
5662        let final_ids = frames[2..4]
5663            .iter()
5664            .map(|frame| match frame {
5665                StreamFrame::PermissionRequestApproved { payload, .. } => {
5666                    payload.request_id.clone().unwrap()
5667                }
5668                _ => unreachable!(),
5669            })
5670            .collect::<BTreeSet<_>>();
5671        assert_eq!(final_ids, BTreeSet::from([ids[0].clone(), ids[1].clone()]));
5672        assert!(matches!(
5673            broker.get(&ids[2]).unwrap().state,
5674            PermissionRequestState::Pending { .. }
5675        ));
5676
5677        broker
5678            .resolve_batch(
5679                &root,
5680                PermissionSelector::Group(group.group_id),
5681                PermissionAction::Approve,
5682                None,
5683                None,
5684                BatchMode::BestEffort,
5685                Some(current.revision),
5686            )
5687            .unwrap();
5688        assert!(drain_audit(&mut audit).is_empty());
5689    }
5690
5691    #[test]
5692    fn s8_empty_group_delete_emits_terminal_snapshot_and_live_replay_remove_it() {
5693        use crate::projection::message_window::{TranscriptEntry, replay_transcript_from};
5694        use crate::workflow::WorkflowGraph;
5695
5696        let flows = Arc::new(FlowRegistry::default());
5697        let root = register_root(&flows, "session", true);
5698        let requester = child(&flows, &root);
5699        let broker = PermissionBroker::new(flows);
5700        let mut audit = attach_audit(&broker);
5701        let SubmissionOutcome::Pending(pending) =
5702            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5703        else {
5704            panic!("expected pending request");
5705        };
5706        drain_audit(&mut audit);
5707        let request_id = pending.request.request_id;
5708        let group = broker
5709            .create_group(
5710                &root,
5711                BTreeSet::from([request_id.clone()]),
5712                "delete-empty".into(),
5713            )
5714            .unwrap();
5715        broker
5716            .ungroup_requests(&root, &group.group_id, &BTreeSet::from([request_id]))
5717            .unwrap();
5718        let deleted = broker.delete_empty_group(&root, &group.group_id).unwrap();
5719        assert!(deleted.request_ids.is_empty());
5720        assert!(
5721            broker
5722                .visible_group_get(&root, &group.group_id)
5723                .unwrap()
5724                .is_none()
5725        );
5726
5727        let frames = drain_audit(&mut audit);
5728        assert!(matches!(
5729            frames.as_slice(),
5730            [
5731                StreamFrame::PermissionGroupCreated { .. },
5732                StreamFrame::PermissionGroupUpdated { payload: updated, .. },
5733                StreamFrame::PermissionGroupResolved { payload: resolved, .. }
5734            ] if updated.request_ids.is_empty()
5735                && resolved.request_ids.is_empty()
5736                && resolved.revision == updated.revision
5737        ));
5738        let mut live = WorkflowGraph::new(crate::event::TurnId::now());
5739        for frame in &frames {
5740            live.apply_stream_frame(frame);
5741        }
5742        assert!(!live.permission_groups.contains_key(&group.group_id));
5743
5744        let lines = frames
5745            .iter()
5746            .enumerate()
5747            .map(|(index, frame)| {
5748                let (kind, payload) = match frame {
5749                    StreamFrame::PermissionGroupCreated { payload, .. } => {
5750                        ("permission_group_created", payload)
5751                    }
5752                    StreamFrame::PermissionGroupUpdated { payload, .. } => {
5753                        ("permission_group_updated", payload)
5754                    }
5755                    StreamFrame::PermissionGroupResolved { payload, .. } => {
5756                        ("permission_group_resolved", payload)
5757                    }
5758                    _ => unreachable!(),
5759                };
5760                serde_json::json!({"type": kind, "seq": index + 1, "payload": payload}).to_string()
5761            })
5762            .collect::<Vec<_>>();
5763        let dir = tempfile::tempdir().unwrap();
5764        let path = dir.path().join("events.jsonl");
5765        std::fs::write(&path, lines.join("\n")).unwrap();
5766        let mut replay = WorkflowGraph::new(crate::event::TurnId::now());
5767        for entry in replay_transcript_from(&path).unwrap() {
5768            if let TranscriptEntry::PermissionGroup { payload, resolved } = entry {
5769                replay.apply_permission_group(&payload, resolved);
5770            }
5771        }
5772        assert_eq!(replay.permission_groups, live.permission_groups);
5773        assert!(!replay.permission_groups.contains_key(&group.group_id));
5774    }
5775
5776    #[test]
5777    fn s8_batch_resolve_versus_ungroup_uses_one_locked_membership_snapshot() {
5778        enum Schedule {
5779            BatchThenUngroup,
5780            UngroupThenBatch,
5781        }
5782        for schedule in [Schedule::BatchThenUngroup, Schedule::UngroupThenBatch] {
5783            let flows = Arc::new(FlowRegistry::default());
5784            let root = register_root(&flows, "session", true);
5785            let requester = child(&flows, &root);
5786            let broker = Arc::new(PermissionBroker::new(flows));
5787            let mut audit = attach_audit(&broker);
5788            let SubmissionOutcome::Pending(pending) =
5789                submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5790            else {
5791                panic!("expected pending request");
5792            };
5793            drain_audit(&mut audit);
5794            let request_id = pending.request.request_id.clone();
5795            let group = broker
5796                .create_group(
5797                    &root,
5798                    BTreeSet::from([request_id.clone()]),
5799                    "batch-linearized".into(),
5800                )
5801                .unwrap();
5802            drain_audit(&mut audit);
5803            let (first_done_tx, first_done_rx) = std::sync::mpsc::channel();
5804            let (continue_tx, continue_rx) = std::sync::mpsc::channel();
5805            let batch_first = matches!(schedule, Schedule::BatchThenUngroup);
5806            let resolver = {
5807                let broker = Arc::clone(&broker);
5808                let root = Arc::clone(&root);
5809                let group_id = group.group_id.clone();
5810                let first_done_tx = first_done_tx.clone();
5811                std::thread::spawn(move || {
5812                    if !batch_first {
5813                        continue_rx.recv().unwrap();
5814                    }
5815                    let result = broker.resolve_batch(
5816                        &root,
5817                        PermissionSelector::Group(group_id),
5818                        PermissionAction::Approve,
5819                        None,
5820                        None,
5821                        BatchMode::Atomic,
5822                        None,
5823                    );
5824                    if batch_first {
5825                        first_done_tx.send(()).unwrap();
5826                    }
5827                    result
5828                })
5829            };
5830            let ungrouper = {
5831                let broker = Arc::clone(&broker);
5832                let root = Arc::clone(&root);
5833                let group_id = group.group_id.clone();
5834                let request_id = request_id.clone();
5835                std::thread::spawn(move || {
5836                    if batch_first {
5837                        first_done_rx.recv().unwrap();
5838                    }
5839                    broker
5840                        .ungroup_requests(&root, &group_id, &BTreeSet::from([request_id]))
5841                        .unwrap();
5842                    if !batch_first {
5843                        first_done_tx.send(()).unwrap();
5844                        continue_tx.send(()).unwrap();
5845                    }
5846                })
5847            };
5848            let resolutions = resolver.join().unwrap().unwrap();
5849            ungrouper.join().unwrap();
5850            let frames = drain_audit(&mut audit);
5851
5852            match schedule {
5853                Schedule::BatchThenUngroup => {
5854                    assert_eq!(resolutions.len(), 1);
5855                    assert!(matches!(
5856                        broker.get(&request_id).unwrap().state,
5857                        PermissionRequestState::Approved { .. }
5858                    ));
5859                    let StreamFrame::PermissionRequestApproved { payload, .. } = &frames[0] else {
5860                        panic!("expected approved frame");
5861                    };
5862                    assert_eq!(payload.group_ids, vec![group.group_id.clone()]);
5863                    assert!(matches!(
5864                        frames.as_slice(),
5865                        [
5866                            StreamFrame::PermissionRequestApproved { .. },
5867                            StreamFrame::PermissionGroupResolved { .. },
5868                            StreamFrame::PermissionGroupUpdated { .. }
5869                        ]
5870                    ));
5871                }
5872                Schedule::UngroupThenBatch => {
5873                    assert!(resolutions.is_empty());
5874                    assert!(matches!(
5875                        broker.get(&request_id).unwrap().state,
5876                        PermissionRequestState::Pending { .. }
5877                    ));
5878                    assert!(matches!(
5879                        frames.as_slice(),
5880                        [
5881                            StreamFrame::PermissionGroupUpdated { payload, .. },
5882                            StreamFrame::PermissionGroupResolved { .. }
5883                        ] if payload.request_ids.is_empty()
5884                    ));
5885                }
5886            }
5887        }
5888    }
5889
5890    #[test]
5891    fn s8_non_group_terminal_paths_emit_one_resolved_after_request_finals() {
5892        use crate::workflow::WorkflowGraph;
5893
5894        enum Path {
5895            RequestIdsBatch,
5896            UserCancel,
5897            RunCleanup,
5898        }
5899        for path in [Path::RequestIdsBatch, Path::UserCancel, Path::RunCleanup] {
5900            let flows = Arc::new(FlowRegistry::default());
5901            let root = register_root(&flows, "session", true);
5902            let requester = child(&flows, &root);
5903            let broker = PermissionBroker::new(flows);
5904            let mut audit = attach_audit(&broker);
5905            let SubmissionOutcome::Pending(pending) =
5906                submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5907            else {
5908                panic!("expected pending request");
5909            };
5910            drain_audit(&mut audit);
5911            let request_id = pending.request.request_id.clone();
5912            let group = broker
5913                .create_group(
5914                    &root,
5915                    BTreeSet::from([request_id.clone()]),
5916                    "non-group-terminal".into(),
5917                )
5918                .unwrap();
5919            drain_audit(&mut audit);
5920
5921            match path {
5922                Path::RequestIdsBatch => {
5923                    broker
5924                        .resolve_batch(
5925                            &root,
5926                            PermissionSelector::RequestIds(vec![request_id]),
5927                            PermissionAction::Approve,
5928                            None,
5929                            None,
5930                            BatchMode::Atomic,
5931                            None,
5932                        )
5933                        .unwrap();
5934                }
5935                Path::UserCancel => broker.cancel(&request_id, "cancelled").unwrap(),
5936                Path::RunCleanup => {
5937                    assert_eq!(
5938                        broker.cancel_for_run("session", &requester.run_id, "terminal"),
5939                        1
5940                    );
5941                }
5942            }
5943
5944            let frames = drain_audit(&mut audit);
5945            assert_eq!(
5946                frames
5947                    .iter()
5948                    .filter(|frame| matches!(frame, StreamFrame::PermissionGroupResolved { .. }))
5949                    .count(),
5950                1
5951            );
5952            assert!(matches!(
5953                frames.last(),
5954                Some(StreamFrame::PermissionGroupResolved { payload, .. })
5955                    if payload.group_id == group.group_id
5956            ));
5957            assert!(matches!(
5958                frames.first(),
5959                Some(
5960                    StreamFrame::PermissionRequestApproved { .. }
5961                        | StreamFrame::PermissionRequestCancelled { .. }
5962                )
5963            ));
5964            let mut live = WorkflowGraph::new(crate::event::TurnId::now());
5965            live.apply_permission_group(
5966                &crate::permission_audit::PermissionGroupAudit::from_group(
5967                    &group,
5968                    "session",
5969                    group.created_at,
5970                ),
5971                false,
5972            );
5973            for frame in &frames {
5974                live.apply_stream_frame(frame);
5975            }
5976            assert!(!live.permission_groups.contains_key(&group.group_id));
5977        }
5978
5979        let flows = Arc::new(FlowRegistry::default());
5980        let root = register_root(&flows, "session", true);
5981        let requester = child(&flows, &root);
5982        let broker = PermissionBroker::new(flows);
5983        let mut audit = attach_audit(&broker);
5984        let SubmissionOutcome::Pending(pending) =
5985            submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
5986        else {
5987            panic!("expected pending request");
5988        };
5989        drain_audit(&mut audit);
5990        let request_id = pending.request.request_id.clone();
5991        broker.cancel(&request_id, "already terminal").unwrap();
5992        drain_audit(&mut audit);
5993        let group = broker
5994            .create_group(&root, BTreeSet::from([request_id]), "terminal-only".into())
5995            .unwrap();
5996        assert!(matches!(
5997            drain_audit(&mut audit).as_slice(),
5998            [
5999                StreamFrame::PermissionGroupCreated { payload: created, .. },
6000                StreamFrame::PermissionGroupResolved { payload: resolved, .. }
6001            ] if created.group_id == group.group_id && resolved.group_id == group.group_id
6002        ));
6003    }
6004
6005    #[test]
6006    fn s8_resolve_versus_ungroup_linearization_controls_final_membership_snapshot() {
6007        enum Schedule {
6008            ResolveThenUngroup,
6009            UngroupThenResolve,
6010        }
6011        for schedule in [Schedule::ResolveThenUngroup, Schedule::UngroupThenResolve] {
6012            let flows = Arc::new(FlowRegistry::default());
6013            let root = register_root(&flows, "session", true);
6014            let requester = child(&flows, &root);
6015            let broker = Arc::new(PermissionBroker::new(flows));
6016            let mut audit = attach_audit(&broker);
6017            let SubmissionOutcome::Pending(pending) =
6018                submit_to_flow(&broker, &requester, Arc::clone(&root)).unwrap()
6019            else {
6020                panic!("expected pending request");
6021            };
6022            drain_audit(&mut audit);
6023            let request_id = pending.request.request_id.clone();
6024            let group = broker
6025                .create_group(
6026                    &root,
6027                    BTreeSet::from([request_id.clone()]),
6028                    "linearized".into(),
6029                )
6030                .unwrap();
6031            drain_audit(&mut audit);
6032            let (first_done_tx, first_done_rx) = std::sync::mpsc::channel();
6033            let (continue_tx, continue_rx) = std::sync::mpsc::channel();
6034            let resolve_first = matches!(schedule, Schedule::ResolveThenUngroup);
6035            let resolver = {
6036                let broker = Arc::clone(&broker);
6037                let root = Arc::clone(&root);
6038                let request_id = request_id.clone();
6039                let first_done_tx = first_done_tx.clone();
6040                std::thread::spawn(move || {
6041                    if !resolve_first {
6042                        continue_rx.recv().unwrap();
6043                    }
6044                    let outcome = broker.resolve(
6045                        &request_id,
6046                        &DecisionAuthority::Flow(broker.flow_authority(root).unwrap()),
6047                        PermissionAction::Approve,
6048                        None,
6049                        None,
6050                    );
6051                    if resolve_first {
6052                        first_done_tx.send(()).unwrap();
6053                    }
6054                    outcome
6055                })
6056            };
6057            let ungrouper = {
6058                let broker = Arc::clone(&broker);
6059                let root = Arc::clone(&root);
6060                let group_id = group.group_id.clone();
6061                let request_id = request_id.clone();
6062                std::thread::spawn(move || {
6063                    if resolve_first {
6064                        first_done_rx.recv().unwrap();
6065                    }
6066                    broker
6067                        .ungroup_requests(&root, &group_id, &BTreeSet::from([request_id]))
6068                        .unwrap();
6069                    if !resolve_first {
6070                        first_done_tx.send(()).unwrap();
6071                        continue_tx.send(()).unwrap();
6072                    }
6073                })
6074            };
6075            assert!(resolver.join().unwrap().is_ok());
6076            ungrouper.join().unwrap();
6077            let frames = drain_audit(&mut audit);
6078            let final_payload = frames
6079                .iter()
6080                .find_map(|frame| match frame {
6081                    StreamFrame::PermissionRequestApproved { payload, .. } => Some(payload),
6082                    _ => None,
6083                })
6084                .unwrap();
6085            match schedule {
6086                Schedule::ResolveThenUngroup => {
6087                    assert_eq!(final_payload.group_ids, vec![group.group_id.clone()]);
6088                    assert!(matches!(
6089                        frames.as_slice(),
6090                        [
6091                            StreamFrame::PermissionRequestApproved { .. },
6092                            StreamFrame::PermissionGroupResolved { .. },
6093                            StreamFrame::PermissionGroupUpdated { .. }
6094                        ]
6095                    ));
6096                }
6097                Schedule::UngroupThenResolve => {
6098                    assert!(final_payload.group_ids.is_empty());
6099                    assert!(matches!(
6100                        frames.as_slice(),
6101                        [
6102                            StreamFrame::PermissionGroupUpdated { .. },
6103                            StreamFrame::PermissionGroupResolved { .. },
6104                            StreamFrame::PermissionRequestApproved { .. }
6105                        ]
6106                    ));
6107                }
6108            }
6109        }
6110    }
6111
6112    #[test]
6113    fn atomic_batch_invalid_target_member_leaves_every_request_unchanged() {
6114        let flows = Arc::new(FlowRegistry::default());
6115        let root = register_root(&flows, "session", true);
6116        let requesters = [child(&flows, &root), child(&flows, &root)];
6117        let broker = PermissionBroker::new(flows);
6118        let SubmissionOutcome::Pending(flow_pending) =
6119            submit_to_flow(&broker, &requesters[0], Arc::clone(&root)).unwrap()
6120        else {
6121            panic!("expected pending request");
6122        };
6123        let user_pending = submit_to_user(&broker, &requesters[1]);
6124        let ids = [
6125            flow_pending.request.request_id.clone(),
6126            user_pending.request.request_id.clone(),
6127        ];
6128        let result = broker.resolve_batch(
6129            &root,
6130            PermissionSelector::RequestIds(ids.to_vec()),
6131            PermissionAction::Approve,
6132            None,
6133            None,
6134            BatchMode::Atomic,
6135            None,
6136        );
6137        assert!(
6138            matches!(result, Err(PermissionError::ActorNotAuthorized)),
6139            "unexpected result: {result:?}"
6140        );
6141        for id in ids {
6142            let request = broker.get(&id).unwrap();
6143            assert!(matches!(
6144                request.state,
6145                PermissionRequestState::Pending { .. }
6146            ));
6147            assert!(
6148                request
6149                    .escalation_path
6150                    .iter()
6151                    .all(|hop| hop.actor.is_none())
6152            );
6153        }
6154    }
6155
6156    #[test]
6157    fn atomic_batch_terminal_requester_member_leaves_every_request_unchanged() {
6158        let flows = Arc::new(FlowRegistry::default());
6159        let root = register_root(&flows, "session", true);
6160        let requesters = [child(&flows, &root), child(&flows, &root)];
6161        let broker = PermissionBroker::new(flows);
6162        let pending = requesters.each_ref().map(|requester| {
6163            let SubmissionOutcome::Pending(pending) =
6164                submit_to_flow(&broker, requester, Arc::clone(&root)).unwrap()
6165            else {
6166                panic!("expected pending request");
6167            };
6168            *pending
6169        });
6170        let mut audit = attach_audit(&broker);
6171        *requesters[1].execution_state.lock().unwrap() = FlowExecutionState::Terminal;
6172        let ids = pending
6173            .each_ref()
6174            .map(|pending| pending.request.request_id.clone());
6175        assert!(matches!(
6176            broker.resolve_batch(
6177                &root,
6178                PermissionSelector::RequestIds(ids.to_vec()),
6179                PermissionAction::Approve,
6180                None,
6181                None,
6182                BatchMode::Atomic,
6183                None
6184            ),
6185            Err(PermissionError::ActorNotRunning)
6186        ));
6187        assert!(drain_audit(&mut audit).is_empty());
6188        for id in ids {
6189            let request = broker.get(&id).unwrap();
6190            assert!(matches!(
6191                request.state,
6192                PermissionRequestState::Pending { .. }
6193            ));
6194            assert!(
6195                request
6196                    .escalation_path
6197                    .iter()
6198                    .all(|hop| hop.actor.is_none())
6199            );
6200        }
6201    }
6202
6203    #[test]
6204    fn cancellation_audit_attributes_the_system_component() {
6205        let flows = Arc::new(FlowRegistry::default());
6206        let requester = register_root(&flows, "session", false);
6207        let broker = PermissionBroker::new(flows);
6208        let mut audit = attach_audit(&broker);
6209        let pending = submit_to_user(&broker, &requester);
6210        drain_audit(&mut audit);
6211
6212        broker
6213            .cancel(&pending.request.request_id, "operator cancelled")
6214            .unwrap();
6215
6216        let frames = drain_audit(&mut audit);
6217        let [StreamFrame::PermissionRequestCancelled { payload, .. }] = frames.as_slice() else {
6218            panic!("expected one cancellation audit");
6219        };
6220        assert_eq!(
6221            payload.actor,
6222            Some(crate::permission_audit::PermissionProjectionActor::System {
6223                component: "permission.user_cancel".into(),
6224            })
6225        );
6226        assert_eq!(payload.reason.as_deref(), Some("operator cancelled"));
6227    }
6228
6229    #[test]
6230    fn approval_audit_precedes_its_persistent_grant() {
6231        let flows = Arc::new(FlowRegistry::default());
6232        let requester = register_root(&flows, "session", false);
6233        let broker = PermissionBroker::new(flows);
6234        let mut audit = attach_audit(&broker);
6235        let pending = submit_to_user(&broker, &requester);
6236        drain_audit(&mut audit);
6237
6238        broker
6239            .resolve(
6240                &pending.request.request_id,
6241                &user_decision(&broker, &requester.session_id),
6242                PermissionAction::Approve,
6243                Some(GrantScope::ChildRunSameTool {
6244                    run_id: requester.run_id.clone(),
6245                    tool_name: pending.request.intent.tool_name.clone(),
6246                }),
6247                Some("approved".into()),
6248            )
6249            .unwrap();
6250
6251        let frames = drain_audit(&mut audit);
6252        assert!(matches!(
6253            frames.as_slice(),
6254            [
6255                StreamFrame::PermissionRequestApproved { .. },
6256                StreamFrame::PermissionGrantCreated { .. }
6257            ]
6258        ));
6259    }
6260
6261    #[test]
6262    fn terminal_cleanup_emits_after_the_lifecycle_transition() {
6263        let flows = Arc::new(FlowRegistry::default());
6264        let requester = register_root(&flows, "session", false);
6265        let broker = PermissionBroker::shared(Arc::clone(&flows));
6266        let mut audit = attach_audit(&broker);
6267        submit_to_user(&broker, &requester);
6268        drain_audit(&mut audit);
6269
6270        flows.mark_terminal(&requester.run_id);
6271
6272        let frames = drain_audit(&mut audit);
6273        let [StreamFrame::PermissionRequestCancelled { payload, .. }] = frames.as_slice() else {
6274            panic!("expected synchronous terminal cancellation audit");
6275        };
6276        assert_eq!(
6277            payload.actor,
6278            Some(crate::permission_audit::PermissionProjectionActor::System {
6279                component: "permission.run_cleanup".into(),
6280            })
6281        );
6282    }
6283
6284    #[test]
6285    fn user_facade_hides_flow_targeted_and_terminal_requests() {
6286        let flows = Arc::new(FlowRegistry::default());
6287        let flow_target = register_root(&flows, "session", true);
6288        let requester = child(&flows, &flow_target);
6289        let broker = PermissionBroker::new(flows);
6290
6291        let SubmissionOutcome::Pending(flow_pending) =
6292            submit_to_flow(&broker, &requester, Arc::clone(&flow_target)).unwrap()
6293        else {
6294            panic!("expected pending request");
6295        };
6296        let flow_pending = *flow_pending;
6297        let (requests, groups) = broker.user_list(&requester.session_id);
6298        assert!(requests.is_empty());
6299        assert!(groups.is_empty());
6300        assert!(matches!(
6301            broker.user_create_group(
6302                &requester.session_id,
6303                BTreeSet::from([flow_pending.request.request_id.clone()]),
6304                "hidden flow request".into(),
6305                &HashMap::from([(
6306                    flow_pending.request.request_id.clone(),
6307                    flow_pending.request.revision
6308                )]),
6309            ),
6310            Err(PermissionError::GroupRevisionConflict)
6311        ));
6312        assert!(matches!(
6313            broker.user_resolve(
6314                &requester.session_id,
6315                None,
6316                vec![flow_pending.request.request_id.clone()],
6317                &HashMap::from([(
6318                    flow_pending.request.request_id.clone(),
6319                    flow_pending.request.revision
6320                )]),
6321                None,
6322                PermissionAction::Deny,
6323                None,
6324                None,
6325            ),
6326            Err(PermissionError::GroupRevisionConflict)
6327        ));
6328
6329        let terminal = submit_to_user(&broker, &requester);
6330        let terminal_id = terminal.request.request_id.clone();
6331        broker.cancel(&terminal_id, "test terminal").unwrap();
6332        let (requests, groups) = broker.user_list(&requester.session_id);
6333        assert!(requests.is_empty());
6334        assert!(groups.is_empty());
6335        assert!(matches!(
6336            broker.user_create_group(
6337                &requester.session_id,
6338                BTreeSet::from([terminal_id.clone()]),
6339                "hidden terminal request".into(),
6340                &HashMap::from([(terminal_id.clone(), terminal.request.revision + 1)]),
6341            ),
6342            Err(PermissionError::GroupRevisionConflict)
6343        ));
6344        assert!(matches!(
6345            broker.user_resolve(
6346                &requester.session_id,
6347                None,
6348                vec![terminal_id.clone()],
6349                &HashMap::from([(terminal_id, terminal.request.revision + 1)]),
6350                None,
6351                PermissionAction::Deny,
6352                None,
6353                None,
6354            ),
6355            Err(PermissionError::GroupRevisionConflict)
6356        ));
6357    }
6358
6359    #[test]
6360    fn user_group_facade_creates_and_resolves_with_current_revision() {
6361        let flows = Arc::new(FlowRegistry::default());
6362        let requester = register_root(&flows, "session", false);
6363        let broker = PermissionBroker::new(flows);
6364        let pending = submit_to_user(&broker, &requester);
6365        let request_id = pending.request.request_id.clone();
6366        let expected_revisions = HashMap::from([(request_id.clone(), pending.request.revision)]);
6367
6368        let group = broker
6369            .user_create_group(
6370                &requester.session_id,
6371                BTreeSet::from([request_id.clone()]),
6372                "user review".into(),
6373                &expected_revisions,
6374            )
6375            .unwrap();
6376
6377        assert_eq!(group.owner, GroupOwner::User);
6378        assert_eq!(group.request_ids, BTreeSet::from([request_id.clone()]));
6379        assert_eq!(
6380            broker.user_list(&requester.session_id).1,
6381            vec![group.clone()]
6382        );
6383
6384        let resolutions = broker
6385            .user_resolve(
6386                &requester.session_id,
6387                Some("principal".into()),
6388                Vec::new(),
6389                &HashMap::new(),
6390                Some((group.group_id, group.revision)),
6391                PermissionAction::Deny,
6392                None,
6393                Some("denied by user".into()),
6394            )
6395            .unwrap();
6396
6397        assert_eq!(resolutions.len(), 1);
6398        assert_eq!(resolutions[0].request_id, request_id);
6399        assert!(matches!(
6400            resolutions[0].outcome,
6401            BatchRequestOutcome::Denied(_)
6402        ));
6403        let resolved = broker.get(&resolutions[0].request_id).unwrap();
6404        assert!(matches!(
6405            resolved.state,
6406            PermissionRequestState::Denied { .. }
6407        ));
6408        assert_eq!(resolved.revision, pending.request.revision + 1);
6409        assert!(matches!(
6410            pending.resolution.blocking_recv().unwrap(),
6411            PermissionResolution::Decision(_)
6412        ));
6413    }
6414
6415    #[test]
6416    fn user_group_facade_rejects_stale_group_revision_without_resolving() {
6417        let flows = Arc::new(FlowRegistry::default());
6418        let requester = register_root(&flows, "session", false);
6419        let broker = PermissionBroker::new(flows);
6420        let pending = submit_to_user(&broker, &requester);
6421        let request_id = pending.request.request_id.clone();
6422        let group = broker
6423            .user_create_group(
6424                &requester.session_id,
6425                BTreeSet::from([request_id.clone()]),
6426                "user review".into(),
6427                &HashMap::from([(request_id.clone(), pending.request.revision)]),
6428            )
6429            .unwrap();
6430        assert!(matches!(
6431            broker.user_resolve(
6432                &requester.session_id,
6433                None,
6434                Vec::new(),
6435                &HashMap::new(),
6436                Some((group.group_id, group.revision + 1)),
6437                PermissionAction::Deny,
6438                None,
6439                None,
6440            ),
6441            Err(PermissionError::GroupRevisionConflict)
6442        ));
6443        assert_eq!(broker.get(&request_id).unwrap(), pending.request);
6444    }
6445
6446    #[test]
6447    fn user_request_facade_rejects_stale_and_duplicate_resolutions() {
6448        let flows = Arc::new(FlowRegistry::default());
6449        let requester = register_root(&flows, "session", false);
6450        let broker = PermissionBroker::new(flows);
6451        let pending = submit_to_user(&broker, &requester);
6452        let request_id = pending.request.request_id.clone();
6453
6454        assert!(matches!(
6455            broker.user_resolve(
6456                &requester.session_id,
6457                None,
6458                vec![request_id.clone()],
6459                &HashMap::from([(request_id.clone(), pending.request.revision + 1)]),
6460                None,
6461                PermissionAction::Deny,
6462                None,
6463                None,
6464            ),
6465            Err(PermissionError::GroupRevisionConflict)
6466        ));
6467
6468        let expected_revisions = HashMap::from([(request_id.clone(), pending.request.revision)]);
6469        let first = broker
6470            .user_resolve(
6471                &requester.session_id,
6472                None,
6473                vec![request_id.clone()],
6474                &expected_revisions,
6475                None,
6476                PermissionAction::Deny,
6477                None,
6478                None,
6479            )
6480            .unwrap();
6481        assert_eq!(first.len(), 1);
6482        assert!(matches!(first[0].outcome, BatchRequestOutcome::Denied(_)));
6483
6484        let after_first = broker.get(&request_id).unwrap();
6485        assert_eq!(after_first.revision, pending.request.revision + 1);
6486        assert!(matches!(
6487            broker.user_resolve(
6488                &requester.session_id,
6489                None,
6490                vec![request_id.clone()],
6491                &expected_revisions,
6492                None,
6493                PermissionAction::Deny,
6494                None,
6495                None,
6496            ),
6497            Err(PermissionError::GroupRevisionConflict)
6498        ));
6499        assert_eq!(broker.get(&request_id).unwrap(), after_first);
6500    }
6501
6502    #[test]
6503    fn concurrent_user_resolutions_have_exactly_one_winner() {
6504        let flows = Arc::new(FlowRegistry::default());
6505        let requester = register_root(&flows, "session", false);
6506        let broker = Arc::new(PermissionBroker::new(flows));
6507        let pending = submit_to_user(&broker, &requester);
6508        let request_id = pending.request.request_id.clone();
6509        let revision = pending.request.revision;
6510        let barrier = Arc::new(std::sync::Barrier::new(3));
6511
6512        let resolve = |broker: Arc<PermissionBroker>, barrier: Arc<std::sync::Barrier>| {
6513            let session_id = requester.session_id.clone();
6514            let request_id = request_id.clone();
6515            std::thread::spawn(move || {
6516                barrier.wait();
6517                broker.user_resolve(
6518                    &session_id,
6519                    None,
6520                    vec![request_id.clone()],
6521                    &HashMap::from([(request_id, revision)]),
6522                    None,
6523                    PermissionAction::Deny,
6524                    None,
6525                    None,
6526                )
6527            })
6528        };
6529
6530        let first = resolve(Arc::clone(&broker), Arc::clone(&barrier));
6531        let second = resolve(Arc::clone(&broker), Arc::clone(&barrier));
6532        barrier.wait();
6533        let results = [first.join().unwrap(), second.join().unwrap()];
6534
6535        assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
6536        assert_eq!(
6537            results
6538                .iter()
6539                .filter(|result| matches!(result, Err(PermissionError::GroupRevisionConflict)))
6540                .count(),
6541            1
6542        );
6543        let resolved = broker.get(&request_id).unwrap();
6544        assert!(matches!(
6545            resolved.state,
6546            PermissionRequestState::Denied { .. }
6547        ));
6548        assert_eq!(resolved.revision, revision + 1);
6549        assert!(matches!(
6550            pending.resolution.blocking_recv().unwrap(),
6551            PermissionResolution::Decision(_)
6552        ));
6553    }
6554
6555    #[test]
6556    fn exact_path_authorization_rejects_sibling_path() {
6557        let root = tempfile::tempdir().unwrap();
6558        let file = root.path().join("file.txt");
6559        let sibling = root.path().join("file.txt.bak");
6560        let other = root.path().join("other.txt");
6561        std::fs::write(&file, "content").unwrap();
6562        std::fs::write(&sibling, "content").unwrap();
6563        std::fs::write(&other, "content").unwrap();
6564        let authorization = InvocationAuthorization::new(
6565            PermissionRequestId::now(),
6566            "call-1",
6567            "fs.write",
6568            ResourceProvenance {
6569                path: Some(crate::fs_access::canonicalize_stable(&file)),
6570                path_origin: Some(PathOrigin::ExplicitInside),
6571                workspace_root: Some(crate::fs_access::canonicalize_stable(root.path())),
6572                ..ResourceProvenance::default()
6573            },
6574            ExecutionBoundary::Sandboxed,
6575        );
6576        assert!(authorization.covers("fs.write", &file));
6577        assert!(!authorization.covers("fs.write", &sibling));
6578        assert!(!authorization.covers("fs.write", &other));
6579    }
6580}