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 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 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 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#[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 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#[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 observer: Mutex<Option<Arc<dyn crate::tools::agent_ctrl::FlowTerminalObserver>>>,
683}
684
685struct 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 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 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 flows.register_terminal_observer(Arc::downgrade(&observer));
879 broker.observer.lock().unwrap().replace(observer);
880 broker
881 }
882
883 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 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 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 #[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(®istered, 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 ¤t.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 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 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}