use std::sync::Arc;
use async_trait::async_trait;
use chrono::SecondsFormat;
use tokio::sync::Mutex;
use cpex_core::delegation::{
payload::{AuthEnforcedBy, TargetType},
DelegationPayload, DelegationSubject, TokenDelegateHook,
};
use cpex_core::extensions::raw_credentials::TokenRole;
use cpex_core::hooks::payload::Extensions;
use cpex_core::manager::PluginManager;
use apl_core::evaluator::Decision;
use apl_core::step::{DelegateStep, DelegationError, DelegationInvoker, DelegationOutcome};
use crate::dispatch_plan::RouteDispatchPlan;
pub struct DelegationPluginInvoker {
manager: Arc<PluginManager>,
extensions: Arc<Mutex<Extensions>>,
plan: Arc<RouteDispatchPlan>,
}
impl DelegationPluginInvoker {
pub fn new(
manager: Arc<PluginManager>,
extensions: Arc<Mutex<Extensions>>,
plan: Arc<RouteDispatchPlan>,
) -> Self {
Self {
manager,
extensions,
plan,
}
}
}
#[async_trait]
impl DelegationInvoker for DelegationPluginInvoker {
async fn delegate(&self, step: &DelegateStep) -> Result<DelegationOutcome, DelegationError> {
let entry = self
.plan
.delegation_entries
.get(&step.plugin_name)
.ok_or_else(|| DelegationError::NotFound(step.plugin_name.clone()))?
.clone();
let current_extensions = self.extensions.lock().await.clone();
let cfg = step.config_override.as_ref().and_then(|v| v.as_mapping());
let subject = subject_from_cfg(cfg)?;
let bearer_token = subject
.inbound_role()
.and_then(|role| {
current_extensions
.raw_credentials
.as_ref()
.and_then(|rc| rc.inbound_tokens.get(&role))
.map(|tok| (*tok.token).clone())
})
.unwrap_or_default();
let target_name: String = cfg
.and_then(|m| m.get(serde_yaml::Value::String("target".into())))
.and_then(|v| v.as_str())
.unwrap_or(&step.plugin_name)
.to_string();
let actor_role = role_from_cfg(cfg, "actor")?;
reject_unsupported_actor_combo(&subject, actor_role.as_ref())?;
let mut payload = DelegationPayload::new(bearer_token, target_name).with_subject(subject);
if let Some(actor_role) = actor_role {
let actor_token = current_extensions
.raw_credentials
.as_ref()
.and_then(|rc| rc.inbound_tokens.get(&actor_role))
.map(|tok| (*tok.token).clone())
.unwrap_or_default();
if !actor_token.is_empty() {
payload = payload.with_actor(actor_role, actor_token);
}
}
if let Some(audience) = cfg
.and_then(|m| m.get(serde_yaml::Value::String("audience".into())))
.and_then(|v| v.as_str())
{
payload = payload.with_target_audience(audience);
}
if let Some(perms) = cfg
.and_then(|m| m.get(serde_yaml::Value::String("permissions".into())))
.and_then(|v| v.as_sequence())
{
let list: Vec<String> = perms
.iter()
.filter_map(|v| v.as_str().map(str::to_string))
.collect();
if !list.is_empty() {
payload = payload.with_required_permissions(list);
}
}
if let Some(t_kind) = cfg
.and_then(|m| m.get(serde_yaml::Value::String("target_type".into())))
.and_then(|v| v.as_str())
{
payload = payload.with_target_type(target_type_from_str(t_kind));
}
if let Some(enforcer) = cfg
.and_then(|m| m.get(serde_yaml::Value::String("auth_enforced_by".into())))
.and_then(|v| v.as_str())
{
payload = payload.with_auth_enforced_by(auth_enforced_by_from_str(enforcer));
}
let (result, _bg) = self
.manager
.invoke_entries::<TokenDelegateHook>(
std::slice::from_ref(&entry),
payload,
current_extensions,
None,
)
.await;
if !result.continue_processing {
let decision = match result.violation {
Some(v) => Decision::Deny {
reason: Some(v.reason),
rule_source: v.code,
},
None => Decision::Deny {
reason: Some(format!(
"delegate `{}` denied without violation detail",
step.plugin_name
)),
rule_source: step.source.clone(),
},
};
return Ok(DelegationOutcome::deny(decision));
}
let resolved = DelegationPayload::from_pipeline_result(&result).ok_or_else(|| {
DelegationError::Dispatch(format!(
"plugin `{}` returned allow but no DelegationPayload",
step.plugin_name,
))
})?;
{
let mut ext_lock = self.extensions.lock().await;
let merged = resolved.clone().apply_to_extensions(ext_lock.clone());
*ext_lock = merged;
}
let (granted_permissions, granted_audience, granted_expires_at) =
match resolved.delegated_token {
Some(tok) => (
tok.scopes,
Some(tok.audience),
Some(tok.expires_at.to_rfc3339_opts(SecondsFormat::Secs, true)),
),
None => (Vec::new(), None, None),
};
Ok(DelegationOutcome {
decision: Decision::Allow,
granted_permissions,
granted_audience,
granted_expires_at,
})
}
}
fn subject_from_cfg(
cfg: Option<&serde_yaml::Mapping>,
) -> Result<DelegationSubject, DelegationError> {
let Some(v) = cfg.and_then(|m| m.get(serde_yaml::Value::String("subject".into()))) else {
return Ok(DelegationSubject::default());
};
let s = v
.as_str()
.ok_or_else(|| DelegationError::InvalidConfig("`subject:` must be a string".into()))?;
DelegationSubject::from_config_str(s).ok_or_else(|| {
DelegationError::InvalidConfig(format!(
"unknown `subject: {s}` (expected user | client | caller_workload | this_workload)"
))
})
}
fn role_from_cfg(
cfg: Option<&serde_yaml::Mapping>,
key: &str,
) -> Result<Option<TokenRole>, DelegationError> {
let Some(v) = cfg.and_then(|m| m.get(serde_yaml::Value::String(key.into()))) else {
return Ok(None);
};
let s = v
.as_str()
.ok_or_else(|| DelegationError::InvalidConfig(format!("`{key}:` must be a string")))?;
match s {
"user" => Ok(Some(TokenRole::User)),
"client" => Ok(Some(TokenRole::Client)),
"caller_workload" | "workload" => Ok(Some(TokenRole::CallerWorkload)),
other => Err(DelegationError::InvalidConfig(format!(
"unknown `{key}: {other}` (expected user | client | caller_workload)"
))),
}
}
fn reject_unsupported_actor_combo(
subject: &DelegationSubject,
actor: Option<&TokenRole>,
) -> Result<(), DelegationError> {
if actor.is_some()
&& matches!(
subject,
DelegationSubject::CallerWorkload | DelegationSubject::ThisWorkload
)
{
return Err(DelegationError::InvalidConfig(
"`actor:` is not supported with `subject: caller_workload` or \
`subject: this_workload` — the actor would be silently ignored"
.into(),
));
}
Ok(())
}
fn target_type_from_str(s: &str) -> TargetType {
match s.to_ascii_lowercase().as_str() {
"tool" => TargetType::Tool,
"agent" => TargetType::Agent,
"resource" => TargetType::Resource,
"service" => TargetType::Service,
other => TargetType::Custom(other.to_string()),
}
}
fn auth_enforced_by_from_str(s: &str) -> AuthEnforcedBy {
match s.to_ascii_lowercase().as_str() {
"caller" => AuthEnforcedBy::Caller,
"target" => AuthEnforcedBy::Target,
_ => AuthEnforcedBy::Caller,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn cfg(yaml: &str) -> serde_yaml::Mapping {
serde_yaml::from_str::<serde_yaml::Value>(yaml)
.expect("valid yaml")
.as_mapping()
.expect("yaml is a mapping")
.clone()
}
#[test]
fn subject_variants_parse() {
let s = |y| subject_from_cfg(Some(&cfg(y))).unwrap();
assert_eq!(s("subject: user"), DelegationSubject::User);
assert_eq!(s("subject: client"), DelegationSubject::Client);
assert_eq!(
s("subject: caller_workload"),
DelegationSubject::CallerWorkload
);
assert_eq!(s("subject: workload"), DelegationSubject::CallerWorkload); assert_eq!(s("subject: this_workload"), DelegationSubject::ThisWorkload);
}
#[test]
fn subject_absent_defaults_to_user() {
assert_eq!(
subject_from_cfg(Some(&cfg("target: hr-service"))).unwrap(),
DelegationSubject::User
);
assert_eq!(subject_from_cfg(None).unwrap(), DelegationSubject::User);
}
#[test]
fn subject_typo_is_rejected_not_defaulted_to_user() {
let err = subject_from_cfg(Some(&cfg("subject: workloadd"))).unwrap_err();
assert!(matches!(err, DelegationError::InvalidConfig(_)), "{err:?}");
}
#[test]
fn subject_non_string_is_rejected() {
let err = subject_from_cfg(Some(&cfg("subject: [a, b]"))).unwrap_err();
assert!(matches!(err, DelegationError::InvalidConfig(_)), "{err:?}");
}
#[test]
fn actor_variants_parse() {
let a = |y| role_from_cfg(Some(&cfg(y)), "actor").unwrap();
assert_eq!(a("actor: user"), Some(TokenRole::User));
assert_eq!(a("actor: client"), Some(TokenRole::Client));
assert_eq!(a("actor: workload"), Some(TokenRole::CallerWorkload));
}
#[test]
fn actor_absent_is_ok_none() {
assert_eq!(
role_from_cfg(Some(&cfg("subject: user")), "actor").unwrap(),
None
);
assert_eq!(role_from_cfg(None, "actor").unwrap(), None);
}
#[test]
fn actor_typo_is_rejected_not_silently_dropped() {
let err = role_from_cfg(Some(&cfg("actor: workloadd")), "actor").unwrap_err();
assert!(matches!(err, DelegationError::InvalidConfig(_)), "{err:?}");
}
#[test]
fn subject_and_actor_resolve_independently() {
let m = cfg("subject: user\nactor: workload");
assert_eq!(subject_from_cfg(Some(&m)).unwrap(), DelegationSubject::User);
assert_eq!(
role_from_cfg(Some(&m), "actor").unwrap(),
Some(TokenRole::CallerWorkload)
);
}
#[test]
fn actor_with_user_or_client_subject_is_allowed() {
for s in [DelegationSubject::User, DelegationSubject::Client] {
assert!(reject_unsupported_actor_combo(&s, Some(&TokenRole::Client)).is_ok());
}
assert!(reject_unsupported_actor_combo(&DelegationSubject::CallerWorkload, None).is_ok());
}
#[test]
fn actor_with_workload_or_this_workload_subject_is_rejected() {
for s in [
DelegationSubject::CallerWorkload,
DelegationSubject::ThisWorkload,
] {
let err =
reject_unsupported_actor_combo(&s, Some(&TokenRole::CallerWorkload)).unwrap_err();
assert!(matches!(err, DelegationError::InvalidConfig(_)), "{err:?}");
}
}
}