1use std::collections::{BTreeSet, HashMap};
10use std::sync::{Arc, RwLock};
11
12use aion::Engine;
13use aion_core::{ScheduleId, SearchAttributeValue, WorkflowId, search_attributes_from_events};
14use aion_proto::WireError;
15use async_trait::async_trait;
16
17use crate::config::{NamespaceConfig, NamespaceMode};
18use crate::error::ServerError;
19
20use super::schedule_source::{HistoryScheduleNamespaceSource, ScheduleNamespaceSource};
21
22pub const NAMESPACE_ATTRIBUTE: &str = "aion.namespace";
25
26pub const TASK_QUEUE_ATTRIBUTE: &str = aion_core::START_TIME_TASK_QUEUE_ATTRIBUTE;
37
38#[derive(Clone, Copy, Debug, Eq, PartialEq)]
42pub(crate) enum GrantSource {
43 NamespacesHeader,
46 TokenClaim,
48 Operator,
54}
55
56impl GrantSource {
57 pub(crate) const fn label(self) -> &'static str {
59 match self {
60 Self::NamespacesHeader => "header",
61 Self::TokenClaim => "token_claim",
62 Self::Operator => "operator",
63 }
64 }
65}
66
67#[derive(Clone, Debug, Eq, PartialEq)]
69pub struct CallerIdentity {
70 subject: String,
71 namespaces: BTreeSet<String>,
72 denial_reason: Option<String>,
73 grant_source: GrantSource,
74 deploy: bool,
77 all_namespaces: bool,
82}
83
84impl CallerIdentity {
85 #[must_use]
88 pub fn new(subject: impl Into<String>, namespaces: impl IntoIterator<Item = String>) -> Self {
89 Self {
90 subject: subject.into(),
91 namespaces: namespaces.into_iter().collect(),
92 denial_reason: None,
93 grant_source: GrantSource::NamespacesHeader,
94 deploy: false,
95 all_namespaces: false,
96 }
97 }
98
99 #[must_use]
103 pub fn from_token_claims(
104 subject: impl Into<String>,
105 namespaces: impl IntoIterator<Item = String>,
106 ) -> Self {
107 Self {
108 subject: subject.into(),
109 namespaces: namespaces.into_iter().collect(),
110 denial_reason: None,
111 grant_source: GrantSource::TokenClaim,
112 deploy: false,
113 all_namespaces: false,
114 }
115 }
116
117 #[must_use]
127 pub fn operator(subject: impl Into<String>) -> Self {
128 Self {
129 subject: subject.into(),
130 namespaces: BTreeSet::new(),
131 denial_reason: None,
132 grant_source: GrantSource::Operator,
133 deploy: true,
134 all_namespaces: true,
135 }
136 }
137
138 #[must_use]
145 pub fn with_deploy(mut self, deploy: bool) -> Self {
146 self.deploy = deploy;
147 self
148 }
149
150 #[must_use]
152 pub const fn deploy_granted(&self) -> bool {
153 self.deploy
154 }
155
156 #[must_use]
158 pub fn denied(subject: impl Into<String>, reason: impl Into<String>) -> Self {
159 Self {
160 subject: subject.into(),
161 namespaces: BTreeSet::new(),
162 denial_reason: Some(reason.into()),
163 grant_source: GrantSource::NamespacesHeader,
164 deploy: false,
165 all_namespaces: false,
166 }
167 }
168
169 #[must_use]
171 pub fn subject(&self) -> &str {
172 &self.subject
173 }
174
175 #[must_use]
180 pub fn namespaces(&self) -> Vec<String> {
181 self.namespaces.iter().cloned().collect()
182 }
183
184 #[must_use]
189 pub const fn all_namespaces(&self) -> bool {
190 self.all_namespaces
191 }
192
193 pub(crate) fn can_access(&self, namespace: &str) -> bool {
202 self.all_namespaces || self.namespaces.contains(namespace)
203 }
204
205 pub(crate) fn denial_reason(&self) -> Option<&str> {
206 self.denial_reason.as_deref()
207 }
208
209 pub(crate) const fn grant_source(&self) -> GrantSource {
212 self.grant_source
213 }
214}
215
216#[derive(Clone)]
218pub struct ScopedEngine {
219 namespace: String,
220 engine: Option<Arc<Engine>>,
221}
222
223impl ScopedEngine {
224 #[must_use]
226 pub fn namespace(&self) -> &str {
227 &self.namespace
228 }
229
230 pub fn engine(&self) -> Result<&Arc<Engine>, ServerError> {
237 self.engine.as_ref().ok_or_else(|| ServerError::Config {
238 message: "namespace resolver has no engine handle".to_owned(),
239 })
240 }
241}
242
243#[derive(Clone, Debug, Eq, PartialEq)]
250pub struct WorkflowAttribution {
251 pub namespace: String,
253 pub workflow_type: Option<String>,
256}
257
258#[async_trait]
264pub trait WorkflowNamespaceSource: Send + Sync {
265 async fn workflow_attribution(
272 &self,
273 workflow_id: &WorkflowId,
274 ) -> Result<Option<WorkflowAttribution>, ServerError>;
275}
276
277struct HistoryNamespaceSource {
281 engine: Arc<Engine>,
282}
283
284#[async_trait]
285impl WorkflowNamespaceSource for HistoryNamespaceSource {
286 async fn workflow_attribution(
287 &self,
288 workflow_id: &WorkflowId,
289 ) -> Result<Option<WorkflowAttribution>, ServerError> {
290 let history = self
291 .engine
292 .store()
293 .read_history(workflow_id)
294 .await
295 .map_err(ServerError::from)?;
296 let namespace = match search_attributes_from_events(&history).remove(NAMESPACE_ATTRIBUTE) {
297 Some(SearchAttributeValue::String(namespace)) => namespace,
298 Some(other) => {
299 return Err(ServerError::Config {
300 message: format!(
301 "workflow {workflow_id} recorded a non-string {NAMESPACE_ATTRIBUTE} search attribute: {other:?}"
302 ),
303 });
304 }
305 None => return Ok(None),
306 };
307 let workflow_type = history.iter().rev().find_map(|event| match event {
310 aion_core::Event::WorkflowStarted { workflow_type, .. } => Some(workflow_type.clone()),
311 _ => None,
312 });
313 Ok(Some(WorkflowAttribution {
314 namespace,
315 workflow_type,
316 }))
317 }
318}
319
320#[derive(Clone, Default)]
323pub struct StaticWorkflowNamespaces {
324 inner: Arc<RwLock<HashMap<WorkflowId, WorkflowAttribution>>>,
325}
326
327impl StaticWorkflowNamespaces {
328 pub fn record(&self, workflow_id: WorkflowId, namespace: &str) -> Result<(), ServerError> {
336 self.insert(
337 workflow_id,
338 WorkflowAttribution {
339 namespace: namespace.to_owned(),
340 workflow_type: None,
341 },
342 )
343 }
344
345 pub fn record_with_type(
352 &self,
353 workflow_id: WorkflowId,
354 namespace: &str,
355 workflow_type: &str,
356 ) -> Result<(), ServerError> {
357 self.insert(
358 workflow_id,
359 WorkflowAttribution {
360 namespace: namespace.to_owned(),
361 workflow_type: Some(workflow_type.to_owned()),
362 },
363 )
364 }
365
366 fn insert(
367 &self,
368 workflow_id: WorkflowId,
369 attribution: WorkflowAttribution,
370 ) -> Result<(), ServerError> {
371 let mut ownership = self
372 .inner
373 .write()
374 .map_err(|_| ServerError::lock_poisoned("namespace workflow ownership"))?;
375 ownership.insert(workflow_id, attribution);
376 Ok(())
377 }
378}
379
380#[async_trait]
381impl WorkflowNamespaceSource for StaticWorkflowNamespaces {
382 async fn workflow_attribution(
383 &self,
384 workflow_id: &WorkflowId,
385 ) -> Result<Option<WorkflowAttribution>, ServerError> {
386 let ownership = self
387 .inner
388 .read()
389 .map_err(|_| ServerError::lock_poisoned("namespace workflow ownership"))?;
390 Ok(ownership.get(workflow_id).cloned())
391 }
392}
393
394#[derive(Clone)]
396pub struct NamespaceResolver {
397 mode: NamespaceMode,
398 engine: Option<Arc<Engine>>,
399 ownership: Arc<dyn WorkflowNamespaceSource>,
400 schedule_ownership: Arc<dyn ScheduleNamespaceSource>,
401}
402
403impl NamespaceResolver {
404 #[must_use]
407 pub fn from_config(config: NamespaceConfig, engine: Arc<Engine>) -> Self {
408 Self {
409 mode: config.mode,
410 ownership: Arc::new(HistoryNamespaceSource {
411 engine: Arc::clone(&engine),
412 }),
413 schedule_ownership: Arc::new(HistoryScheduleNamespaceSource::new(Arc::clone(&engine))),
414 engine: Some(engine),
415 }
416 }
417
418 #[must_use]
420 pub fn from_parts(
421 mode: NamespaceMode,
422 engine: Option<Arc<Engine>>,
423 ownership: Arc<dyn WorkflowNamespaceSource>,
424 schedule_ownership: Arc<dyn ScheduleNamespaceSource>,
425 ) -> Self {
426 Self {
427 mode,
428 engine,
429 ownership,
430 schedule_ownership,
431 }
432 }
433
434 #[must_use]
439 pub fn authorization_only(
440 mode: NamespaceMode,
441 ownership: impl WorkflowNamespaceSource + 'static,
442 schedule_ownership: impl ScheduleNamespaceSource + 'static,
443 ) -> Self {
444 Self::from_parts(
445 mode,
446 None,
447 Arc::new(ownership),
448 Arc::new(schedule_ownership),
449 )
450 }
451
452 #[must_use]
454 pub const fn mode(&self) -> &NamespaceMode {
455 &self.mode
456 }
457
458 pub(crate) fn engine(&self) -> Result<&Arc<Engine>, ServerError> {
465 self.engine.as_ref().ok_or_else(|| ServerError::Config {
466 message: "namespace resolver has no engine handle".to_owned(),
467 })
468 }
469
470 pub fn shutdown_engine(&self) -> Result<(), ServerError> {
477 self.engine
478 .as_ref()
479 .ok_or_else(|| ServerError::Config {
480 message: "namespace resolver has no engine handle".to_owned(),
481 })?
482 .shutdown()
483 .map_err(ServerError::from)
484 }
485
486 pub(super) fn resolve(
494 &self,
495 caller: &CallerIdentity,
496 requested_namespace: &str,
497 ) -> Result<ScopedEngine, ServerError> {
498 if requested_namespace.is_empty() {
499 return Err(ServerError::namespace_denied(
500 "requested namespace must not be empty",
501 ));
502 }
503
504 if let Some(reason) = caller.denial_reason() {
505 return Err(ServerError::namespace_denied(reason));
506 }
507
508 match &self.mode {
509 NamespaceMode::SingleTenant { namespace } if namespace == requested_namespace => {
510 Ok(self.scoped(requested_namespace))
511 }
512 NamespaceMode::SharedEngine if caller.can_access(requested_namespace) => {
513 Ok(self.scoped(requested_namespace))
514 }
515 NamespaceMode::SingleTenant { .. } | NamespaceMode::SharedEngine => {
516 Err(namespace_denied(caller, requested_namespace))
517 }
518 }
519 }
520
521 pub async fn verify_workflow_ownership(
539 &self,
540 namespace: &str,
541 workflow_id: &WorkflowId,
542 ) -> Result<(), ServerError> {
543 match self.workflow_attribution(namespace, workflow_id).await? {
544 Some(_) => Ok(()),
545 None => Err(ServerError::Wire {
546 wire: WireError::not_found(format!("workflow not found in namespace {namespace}")),
547 }),
548 }
549 }
550
551 pub async fn workflow_attribution(
567 &self,
568 namespace: &str,
569 workflow_id: &WorkflowId,
570 ) -> Result<Option<WorkflowAttribution>, ServerError> {
571 Ok(self
572 .ownership
573 .workflow_attribution(workflow_id)
574 .await?
575 .filter(|attribution| attribution.namespace == namespace))
576 }
577
578 pub async fn recorded_workflow_attribution(
603 &self,
604 workflow_id: &WorkflowId,
605 ) -> Result<Option<WorkflowAttribution>, ServerError> {
606 self.ownership.workflow_attribution(workflow_id).await
607 }
608
609 pub async fn verify_schedule_ownership(
627 &self,
628 namespace: &str,
629 schedule_id: &ScheduleId,
630 ) -> Result<(), ServerError> {
631 match self
632 .schedule_ownership
633 .schedule_namespace(schedule_id)
634 .await?
635 {
636 Some(owner) if owner == namespace => Ok(()),
637 Some(_) | None => Err(ServerError::Wire {
640 wire: WireError::not_found(format!("schedule not found in namespace {namespace}")),
641 }),
642 }
643 }
644
645 fn scoped(&self, namespace: &str) -> ScopedEngine {
646 ScopedEngine {
647 namespace: namespace.to_owned(),
648 engine: self.engine.clone(),
649 }
650 }
651}
652
653fn namespace_denied(caller: &CallerIdentity, requested_namespace: &str) -> ServerError {
654 let hint = match caller.grant_source {
655 GrantSource::NamespacesHeader => format!(
656 "add {requested_namespace} to x-aion-namespaces for subject `{}` or request a namespace listed in that header",
657 caller.subject()
658 ),
659 GrantSource::TokenClaim => format!(
660 "grant {requested_namespace} in the namespace claim of the token minted for subject `{}` or request a namespace the token grants",
661 caller.subject()
662 ),
663 GrantSource::Operator => format!(
667 "subject `{}` is the operator and already holds every namespace",
668 caller.subject()
669 ),
670 };
671 ServerError::namespace_denied(format!(
672 "subject not authorized for namespace {requested_namespace}; {hint}"
673 ))
674}
675
676#[cfg(test)]
677mod tests {
678 use super::{
679 CallerIdentity, NamespaceResolver, StaticWorkflowNamespaces, WorkflowNamespaceSource,
680 };
681 use crate::config::NamespaceMode;
682 use crate::namespace::StaticScheduleNamespaces;
683 use aion_core::{ScheduleId, WorkflowId};
684
685 fn resolver(mode: NamespaceMode) -> NamespaceResolver {
686 NamespaceResolver::authorization_only(
687 mode,
688 StaticWorkflowNamespaces::default(),
689 StaticScheduleNamespaces::default(),
690 )
691 }
692
693 #[test]
694 fn shared_engine_authorizes_explicit_caller_grant() -> Result<(), Box<dyn std::error::Error>> {
695 let resolver = resolver(NamespaceMode::SharedEngine);
696 let caller = CallerIdentity::new("alice", [String::from("tenant-a")]);
697
698 let scoped = resolver.resolve(&caller, "tenant-a")?;
699
700 assert_eq!(scoped.namespace(), "tenant-a");
701 Ok(())
702 }
703
704 #[test]
709 fn operator_is_authorized_for_any_namespace() -> Result<(), Box<dyn std::error::Error>> {
710 let resolver = resolver(NamespaceMode::SharedEngine);
711 let operator = CallerIdentity::operator("operator");
712
713 assert!(operator.all_namespaces());
714 assert!(operator.deploy_granted());
715 assert!(operator.namespaces().is_empty());
716
717 assert_eq!(
718 resolver.resolve(&operator, "tenant-a")?.namespace(),
719 "tenant-a"
720 );
721 assert_eq!(
722 resolver.resolve(&operator, "tenant-z")?.namespace(),
723 "tenant-z"
724 );
725 Ok(())
726 }
727
728 #[test]
729 fn shared_engine_denies_missing_caller_grant() {
730 let resolver = resolver(NamespaceMode::SharedEngine);
731 let caller = CallerIdentity::new("alice", [String::from("tenant-a")]);
732
733 let denied = resolver.resolve(&caller, "tenant-b");
734
735 assert!(denied.is_err());
736 }
737
738 #[test]
739 fn single_tenant_authorizes_only_configured_namespace() -> Result<(), Box<dyn std::error::Error>>
740 {
741 let resolver = resolver(NamespaceMode::SingleTenant {
742 namespace: String::from("tenant-a"),
743 });
744 let caller = CallerIdentity::new("alice", [String::from("tenant-b")]);
745
746 let scoped = resolver.resolve(&caller, "tenant-a")?;
747 let denied = resolver.resolve(&caller, "tenant-b");
748
749 assert_eq!(scoped.namespace(), "tenant-a");
750 assert!(denied.is_err());
751 Ok(())
752 }
753
754 #[test]
759 fn denial_hint_names_the_grant_source() -> Result<(), Box<dyn std::error::Error>> {
760 let resolver = resolver(NamespaceMode::SharedEngine);
761
762 let header_caller = CallerIdentity::new("alice", [String::from("tenant-a")]);
763 let header_denial = resolver
764 .resolve(&header_caller, "tenant-b")
765 .err()
766 .map(|error| error.to_wire_error())
767 .ok_or("expected header-sourced caller to be denied")?;
768 assert!(
769 header_denial.message.contains("x-aion-namespaces"),
770 "header-path denial must hint the dev header: {}",
771 header_denial.message
772 );
773 assert!(
774 !header_denial.message.contains("namespace claim"),
775 "header-path denial must not hint the token claim: {}",
776 header_denial.message
777 );
778
779 let token_caller = CallerIdentity::from_token_claims("alice", [String::from("tenant-a")]);
780 let token_denial = resolver
781 .resolve(&token_caller, "tenant-b")
782 .err()
783 .map(|error| error.to_wire_error())
784 .ok_or("expected token-sourced caller to be denied")?;
785 assert!(
786 token_denial.message.contains("namespace claim"),
787 "JWT-path denial must hint the token's namespace claim: {}",
788 token_denial.message
789 );
790 assert!(
791 !token_denial.message.contains("x-aion-namespaces"),
792 "JWT-path denial must not hint the dev header: {}",
793 token_denial.message
794 );
795 Ok(())
796 }
797
798 #[test]
799 fn empty_namespace_is_denied_before_scoping() {
800 let resolver = resolver(NamespaceMode::SharedEngine);
801 let caller = CallerIdentity::new("alice", [String::new()]);
802
803 let denied = resolver.resolve(&caller, "");
804
805 assert!(denied.is_err());
806 }
807
808 #[tokio::test]
809 async fn ownership_misses_are_indistinguishable_not_found()
810 -> Result<(), Box<dyn std::error::Error>> {
811 let ownership = StaticWorkflowNamespaces::default();
812 let owned = WorkflowId::new(uuid::Uuid::from_u128(1));
813 let unknown = WorkflowId::new(uuid::Uuid::from_u128(2));
814 ownership.record(owned.clone(), "tenant-a")?;
815 let resolver = NamespaceResolver::authorization_only(
816 NamespaceMode::SharedEngine,
817 ownership,
818 StaticScheduleNamespaces::default(),
819 );
820
821 resolver
822 .verify_workflow_ownership("tenant-a", &owned)
823 .await?;
824
825 let foreign = resolver
829 .verify_workflow_ownership("tenant-b", &owned)
830 .await
831 .err()
832 .map(|error| error.to_wire_error())
833 .ok_or("expected foreign-owned workflow to be rejected")?;
834 let absent = resolver
835 .verify_workflow_ownership("tenant-b", &unknown)
836 .await
837 .err()
838 .map(|error| error.to_wire_error())
839 .ok_or("expected unknown workflow to be rejected")?;
840
841 assert_eq!(foreign.code, aion_proto::WireErrorCode::NotFound);
842 assert_eq!(foreign, absent);
843 assert_eq!(foreign.message, "workflow not found in namespace tenant-b");
844
845 let absent_in_granted = resolver
846 .verify_workflow_ownership("tenant-a", &unknown)
847 .await
848 .err()
849 .map(|error| error.to_wire_error())
850 .ok_or("expected unknown workflow to be rejected in granted namespace")?;
851 assert_eq!(absent_in_granted.code, aion_proto::WireErrorCode::NotFound);
852 assert_eq!(
853 absent_in_granted.message,
854 "workflow not found in namespace tenant-a"
855 );
856 Ok(())
857 }
858
859 #[tokio::test]
860 async fn schedule_ownership_misses_are_indistinguishable_not_found()
861 -> Result<(), Box<dyn std::error::Error>> {
862 let schedule_ownership = StaticScheduleNamespaces::default();
863 let owned = ScheduleId::new(uuid::Uuid::from_u128(1));
864 let unknown = ScheduleId::new(uuid::Uuid::from_u128(2));
865 schedule_ownership.record(owned.clone(), "tenant-a")?;
866 let resolver = NamespaceResolver::authorization_only(
867 NamespaceMode::SharedEngine,
868 StaticWorkflowNamespaces::default(),
869 schedule_ownership,
870 );
871
872 resolver
873 .verify_schedule_ownership("tenant-a", &owned)
874 .await?;
875
876 let foreign = resolver
880 .verify_schedule_ownership("tenant-b", &owned)
881 .await
882 .err()
883 .map(|error| error.to_wire_error())
884 .ok_or("expected foreign-owned schedule to be rejected")?;
885 let absent = resolver
886 .verify_schedule_ownership("tenant-b", &unknown)
887 .await
888 .err()
889 .map(|error| error.to_wire_error())
890 .ok_or("expected unknown schedule to be rejected")?;
891
892 assert_eq!(foreign.code, aion_proto::WireErrorCode::NotFound);
893 assert_eq!(foreign, absent);
894 assert_eq!(foreign.message, "schedule not found in namespace tenant-b");
895
896 let absent_in_granted = resolver
897 .verify_schedule_ownership("tenant-a", &unknown)
898 .await
899 .err()
900 .map(|error| error.to_wire_error())
901 .ok_or("expected unknown schedule to be rejected in granted namespace")?;
902 assert_eq!(absent_in_granted.code, aion_proto::WireErrorCode::NotFound);
903 assert_eq!(
904 absent_in_granted.message,
905 "schedule not found in namespace tenant-a"
906 );
907 Ok(())
908 }
909
910 #[tokio::test]
911 async fn static_source_reports_recorded_namespace() -> Result<(), Box<dyn std::error::Error>> {
912 let ownership = StaticWorkflowNamespaces::default();
913 let workflow_id = WorkflowId::new(uuid::Uuid::from_u128(3));
914 ownership.record(workflow_id.clone(), "tenant-a")?;
915
916 assert_eq!(
917 ownership.workflow_attribution(&workflow_id).await?,
918 Some(super::WorkflowAttribution {
919 namespace: String::from("tenant-a"),
920 workflow_type: None,
921 })
922 );
923 Ok(())
924 }
925
926 #[tokio::test]
927 async fn static_source_reports_recorded_workflow_type() -> Result<(), Box<dyn std::error::Error>>
928 {
929 let ownership = StaticWorkflowNamespaces::default();
930 let workflow_id = WorkflowId::new(uuid::Uuid::from_u128(4));
931 ownership.record_with_type(workflow_id.clone(), "tenant-a", "checkout")?;
932
933 assert_eq!(
934 ownership.workflow_attribution(&workflow_id).await?,
935 Some(super::WorkflowAttribution {
936 namespace: String::from("tenant-a"),
937 workflow_type: Some(String::from("checkout")),
938 })
939 );
940 Ok(())
941 }
942
943 #[tokio::test]
947 async fn scoped_attribution_hides_foreign_and_unknown_identically()
948 -> Result<(), Box<dyn std::error::Error>> {
949 let ownership = StaticWorkflowNamespaces::default();
950 let owned = WorkflowId::new(uuid::Uuid::from_u128(5));
951 let foreign = WorkflowId::new(uuid::Uuid::from_u128(6));
952 let unknown = WorkflowId::new(uuid::Uuid::from_u128(7));
953 ownership.record_with_type(owned.clone(), "tenant-a", "checkout")?;
954 ownership.record_with_type(foreign.clone(), "tenant-b", "checkout")?;
955 let resolver = NamespaceResolver::authorization_only(
956 NamespaceMode::SharedEngine,
957 ownership,
958 StaticScheduleNamespaces::default(),
959 );
960
961 let visible = resolver
962 .workflow_attribution("tenant-a", &owned)
963 .await?
964 .ok_or("owned workflow attribution must be visible")?;
965 assert_eq!(visible.workflow_type.as_deref(), Some("checkout"));
966 assert_eq!(
967 resolver.workflow_attribution("tenant-a", &foreign).await?,
968 None
969 );
970 assert_eq!(
971 resolver.workflow_attribution("tenant-a", &unknown).await?,
972 None
973 );
974 Ok(())
975 }
976}