use std::collections::HashMap;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use zeroize::Zeroizing;
use crate::executor::PipelineResult;
use crate::extensions::raw_credentials::{DelegationMode, TokenRole};
use crate::extensions::{
DelegationExtension, Extensions, RawCredentialsExtension, RawDelegatedToken,
};
use crate::impl_plugin_payload;
#[non_exhaustive]
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DelegationSubject {
#[default]
User,
Client,
CallerWorkload,
ThisWorkload,
}
impl DelegationSubject {
pub fn inbound_role(&self) -> Option<TokenRole> {
match self {
DelegationSubject::User => Some(TokenRole::User),
DelegationSubject::Client => Some(TokenRole::Client),
DelegationSubject::CallerWorkload => Some(TokenRole::CallerWorkload),
DelegationSubject::ThisWorkload => None,
}
}
pub fn default_mode(&self) -> DelegationMode {
match self {
DelegationSubject::User => DelegationMode::OnBehalfOfUser,
DelegationSubject::Client => DelegationMode::AsClient,
DelegationSubject::CallerWorkload => DelegationMode::AsCallerWorkload,
DelegationSubject::ThisWorkload => DelegationMode::AsThisWorkload,
}
}
pub fn from_config_str(s: &str) -> Option<Self> {
match s {
"user" => Some(DelegationSubject::User),
"client" => Some(DelegationSubject::Client),
"caller_workload" | "workload" => Some(DelegationSubject::CallerWorkload),
"this_workload" | "gateway" => Some(DelegationSubject::ThisWorkload),
_ => None,
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TargetType {
Tool,
Agent,
Resource,
Service,
#[serde(untagged)]
Custom(String),
}
impl Default for TargetType {
fn default() -> Self {
TargetType::Tool
}
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AuthEnforcedBy {
Caller,
Target,
Both,
}
impl Default for AuthEnforcedBy {
fn default() -> Self {
AuthEnforcedBy::Caller
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct AttenuationConfig {
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub capabilities: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub resource_template: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub actions: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ttl_seconds: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DelegationPayload {
#[serde(skip)]
bearer_token: Zeroizing<String>,
#[serde(skip)]
actor_token: Zeroizing<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
actor_role: Option<TokenRole>,
#[serde(default)]
subject: DelegationSubject,
target_name: String,
#[serde(default)]
target_type: TargetType,
#[serde(default, skip_serializing_if = "Option::is_none")]
target_audience: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
required_permissions: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
trust_domain: Option<String>,
#[serde(default)]
auth_enforced_by: AuthEnforcedBy,
#[serde(default, skip_serializing_if = "Option::is_none")]
route_attenuation: Option<AttenuationConfig>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub delegated_token: Option<RawDelegatedToken>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub delegation_update: Option<DelegationExtension>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub delegation_mode: Option<DelegationMode>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub minted_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub metadata: HashMap<String, serde_json::Value>,
}
impl DelegationPayload {
pub fn new(bearer_token: impl Into<String>, target_name: impl Into<String>) -> Self {
Self {
bearer_token: Zeroizing::new(bearer_token.into()),
actor_token: Zeroizing::new(String::new()),
actor_role: None,
subject: DelegationSubject::default(),
target_name: target_name.into(),
target_type: TargetType::Tool,
target_audience: None,
required_permissions: Vec::new(),
trust_domain: None,
auth_enforced_by: AuthEnforcedBy::Caller,
route_attenuation: None,
delegated_token: None,
delegation_update: None,
delegation_mode: None,
minted_at: None,
metadata: HashMap::new(),
}
}
pub fn with_actor(mut self, role: TokenRole, actor_token: impl Into<String>) -> Self {
self.actor_token = Zeroizing::new(actor_token.into());
self.actor_role = Some(role);
self
}
pub fn with_subject(mut self, subject: DelegationSubject) -> Self {
self.subject = subject;
self
}
pub fn with_target_type(mut self, t: TargetType) -> Self {
self.target_type = t;
self
}
pub fn with_target_audience(mut self, aud: impl Into<String>) -> Self {
self.target_audience = Some(aud.into());
self
}
pub fn with_required_permissions(mut self, perms: Vec<String>) -> Self {
self.required_permissions = perms;
self
}
pub fn with_trust_domain(mut self, td: impl Into<String>) -> Self {
self.trust_domain = Some(td.into());
self
}
pub fn with_auth_enforced_by(mut self, who: AuthEnforcedBy) -> Self {
self.auth_enforced_by = who;
self
}
pub fn with_route_attenuation(mut self, cfg: AttenuationConfig) -> Self {
self.route_attenuation = Some(cfg);
self
}
pub fn bearer_token(&self) -> &str {
&self.bearer_token
}
pub fn actor_token(&self) -> &str {
&self.actor_token
}
pub fn subject(&self) -> &DelegationSubject {
&self.subject
}
pub fn actor_role(&self) -> Option<&TokenRole> {
self.actor_role.as_ref()
}
pub fn involves_workload(&self) -> bool {
self.subject == DelegationSubject::CallerWorkload
|| self.actor_role == Some(TokenRole::CallerWorkload)
}
pub fn involves_client(&self) -> bool {
self.subject == DelegationSubject::Client || self.actor_role == Some(TokenRole::Client)
}
pub fn target_name(&self) -> &str {
&self.target_name
}
pub fn target_type(&self) -> &TargetType {
&self.target_type
}
pub fn target_audience(&self) -> Option<&str> {
self.target_audience.as_deref()
}
pub fn required_permissions(&self) -> &[String] {
&self.required_permissions
}
pub fn trust_domain(&self) -> Option<&str> {
self.trust_domain.as_deref()
}
pub fn auth_enforced_by(&self) -> AuthEnforcedBy {
self.auth_enforced_by
}
pub fn route_attenuation(&self) -> Option<&AttenuationConfig> {
self.route_attenuation.as_ref()
}
pub fn merge(&mut self, other: DelegationPayload) {
if other.delegated_token.is_some() {
self.delegated_token = other.delegated_token;
}
if other.delegation_update.is_some() {
self.delegation_update = other.delegation_update;
}
if other.delegation_mode.is_some() {
self.delegation_mode = other.delegation_mode;
}
if other.minted_at.is_some() {
self.minted_at = other.minted_at;
}
for (k, v) in other.metadata {
self.metadata.insert(k, v);
}
}
pub fn from_pipeline_result(result: &PipelineResult) -> Option<Self> {
result
.modified_payload
.as_ref()
.and_then(|p| p.as_any().downcast_ref::<DelegationPayload>())
.cloned()
}
pub fn apply_to_extensions(&self, mut ext: Extensions) -> Extensions {
if let Some(ref token) = self.delegated_token {
use crate::extensions::raw_credentials::DelegationKey;
let subject_id = ext
.security
.as_ref()
.and_then(|s| s.subject.as_ref())
.and_then(|s| s.id.clone())
.unwrap_or_default();
let workload_id = if self.involves_workload() {
ext.security
.as_ref()
.and_then(|s| s.caller_workload.as_ref())
.and_then(|w| w.spiffe_id.clone())
} else {
None
};
let client_id = if self.involves_client() {
ext.security
.as_ref()
.and_then(|s| s.client.as_ref())
.map(|c| c.client_id.clone())
.filter(|id| !id.is_empty())
} else {
None
};
let mode = self
.delegation_mode
.clone()
.unwrap_or(DelegationMode::OnBehalfOfUser);
let key = DelegationKey {
subject_id,
workload_id,
client_id,
audience: token.audience.clone(),
scopes: token.scopes.clone(),
mode,
};
let mut raw = ext
.raw_credentials
.as_ref()
.map(|arc| (**arc).clone())
.unwrap_or_else(RawCredentialsExtension::default);
raw.delegated_tokens.insert(key, token.clone());
ext.raw_credentials = Some(Arc::new(raw));
}
if let Some(ref update) = self.delegation_update {
ext.delegation = Some(Arc::new(update.clone()));
}
ext
}
}
impl_plugin_payload!(DelegationPayload);
#[cfg(test)]
mod tests {
use super::*;
use crate::extensions::raw_credentials::RawDelegatedToken;
#[test]
fn bearer_token_does_not_serialize() {
let p = DelegationPayload::new("eyJ.caller.tok", "get_compensation");
let json = serde_json::to_string(&p).unwrap();
assert!(
!json.contains("eyJ.caller.tok"),
"bearer_token leaked into serialized form: {}",
json,
);
assert!(json.contains("get_compensation"));
}
#[test]
fn deserialize_yields_empty_bearer_token() {
let json = r#"{"target_name":"get_compensation"}"#;
let p: DelegationPayload = serde_json::from_str(json).unwrap();
assert_eq!(p.bearer_token(), "");
assert_eq!(p.target_name(), "get_compensation");
}
#[test]
fn actor_token_defaults_empty_and_builder_sets_it() {
let p = DelegationPayload::new("caller.tok", "get_compensation");
assert_eq!(p.actor_token(), "");
assert_eq!(p.actor_role(), None);
let p = p.with_actor(TokenRole::CallerWorkload, "svid.jwt.bytes");
assert_eq!(p.actor_token(), "svid.jwt.bytes");
assert_eq!(p.actor_role(), Some(&TokenRole::CallerWorkload));
assert_eq!(p.bearer_token(), "caller.tok");
}
#[test]
fn involves_workload_covers_both_subject_and_actor_positions() {
let user_only = DelegationPayload::new("caller.tok", "t");
assert!(!user_only.involves_workload());
let mode_a =
DelegationPayload::new("svid", "t").with_subject(DelegationSubject::CallerWorkload);
assert!(mode_a.involves_workload());
let mode_b = DelegationPayload::new("user.tok", "t")
.with_actor(TokenRole::CallerWorkload, "svid.jwt.bytes");
assert!(mode_b.involves_workload());
let client_actor =
DelegationPayload::new("user.tok", "t").with_actor(TokenRole::Client, "client.tok");
assert!(!client_actor.involves_workload());
}
#[test]
fn subject_defaults_to_user_and_builder_overrides_it() {
let p = DelegationPayload::new("caller.tok", "get_compensation");
assert_eq!(p.subject(), &DelegationSubject::User);
let p = p.with_subject(DelegationSubject::CallerWorkload);
assert_eq!(p.subject(), &DelegationSubject::CallerWorkload);
}
#[test]
fn only_this_workload_has_no_inbound_role() {
assert_eq!(
DelegationSubject::User.inbound_role(),
Some(TokenRole::User),
);
assert_eq!(
DelegationSubject::Client.inbound_role(),
Some(TokenRole::Client),
);
assert_eq!(
DelegationSubject::CallerWorkload.inbound_role(),
Some(TokenRole::CallerWorkload),
);
assert_eq!(DelegationSubject::ThisWorkload.inbound_role(), None);
}
#[test]
fn subject_parses_from_config_including_legacy_spellings() {
assert_eq!(
DelegationSubject::from_config_str("caller_workload"),
Some(DelegationSubject::CallerWorkload),
);
assert_eq!(
DelegationSubject::from_config_str("workload"),
Some(DelegationSubject::CallerWorkload),
);
assert_eq!(
DelegationSubject::from_config_str("this_workload"),
Some(DelegationSubject::ThisWorkload),
);
assert_eq!(
DelegationSubject::from_config_str("gateway"),
Some(DelegationSubject::ThisWorkload),
);
assert_eq!(DelegationSubject::from_config_str("gatewy"), None);
}
#[test]
fn actor_token_does_not_serialize() {
let p = DelegationPayload::new("caller.tok", "get_compensation")
.with_actor(TokenRole::CallerWorkload, "eyJ.workload.svid");
let json = serde_json::to_string(&p).unwrap();
assert!(
!json.contains("eyJ.workload.svid"),
"actor_token leaked into serialized form: {}",
json,
);
}
#[test]
fn input_builders_chain() {
let p = DelegationPayload::new("tok", "get_compensation")
.with_target_type(TargetType::Tool)
.with_target_audience("https://hr.example.com")
.with_required_permissions(vec!["read:compensation".into()])
.with_trust_domain("hr.example.com")
.with_auth_enforced_by(AuthEnforcedBy::Target)
.with_route_attenuation(AttenuationConfig {
capabilities: vec!["read:compensation".into()],
resource_template: Some("hr://employees/{{ args.employee_id }}".into()),
actions: vec!["read".into()],
ttl_seconds: Some(60),
});
assert_eq!(p.bearer_token(), "tok");
assert_eq!(p.target_name(), "get_compensation");
assert_eq!(p.target_audience(), Some("https://hr.example.com"));
assert_eq!(p.required_permissions(), &["read:compensation".to_string()]);
assert_eq!(p.trust_domain(), Some("hr.example.com"));
assert_eq!(p.auth_enforced_by(), AuthEnforcedBy::Target);
let att = p.route_attenuation().unwrap();
assert_eq!(att.ttl_seconds, Some(60));
assert_eq!(att.actions, vec!["read"]);
}
#[test]
fn target_type_custom_round_trips() {
let t = TargetType::Custom("workflow".into());
let json = serde_json::to_string(&t).unwrap();
let back: TargetType = serde_json::from_str(&json).unwrap();
assert_eq!(t, back);
}
#[test]
fn handler_can_populate_output_on_clone() {
let original = DelegationPayload::new("caller-tok", "downstream-tool");
let mut updated = original.clone();
updated.delegated_token = Some(RawDelegatedToken::new(
"minted-bytes",
"Authorization",
"https://api.example.com",
vec!["read".into()],
Utc::now(),
));
assert_eq!(updated.bearer_token(), "caller-tok");
assert_eq!(updated.target_name(), "downstream-tool");
assert!(updated.delegated_token.is_some());
assert!(original.delegated_token.is_none());
}
#[test]
fn merge_overlays_outputs() {
let mut base = DelegationPayload::new("tok", "tool");
base.metadata.insert("attempt".into(), serde_json::json!(1));
let mut overlay = DelegationPayload::new("", "");
overlay.delegated_token = Some(RawDelegatedToken::new(
"x",
"Authorization",
"aud",
vec![],
Utc::now(),
));
overlay
.metadata
.insert("latency_ms".into(), serde_json::json!(42));
base.merge(overlay);
assert!(base.delegated_token.is_some());
assert!(base.metadata.contains_key("attempt"));
assert!(base.metadata.contains_key("latency_ms"));
}
#[test]
fn apply_to_extensions_writes_delegated_token_keyed_by_audience() {
use crate::extensions::raw_credentials::DelegationMode;
use crate::extensions::SubjectExtension;
let mut p = DelegationPayload::new("tok", "get_compensation");
p.delegated_token = Some(RawDelegatedToken::new(
"minted-jwt",
"Authorization",
"https://hr.example.com",
vec!["read:compensation".into()],
Utc::now() + chrono::Duration::seconds(300),
));
let initial_ext = Extensions {
security: Some(Arc::new(crate::extensions::SecurityExtension {
subject: Some(SubjectExtension {
id: Some("alice".into()),
..Default::default()
}),
..Default::default()
})),
..Default::default()
};
let updated = p.apply_to_extensions(initial_ext);
let raw = updated.raw_credentials.as_ref().unwrap();
assert_eq!(raw.delegated_tokens.len(), 1);
let expected_key = crate::extensions::raw_credentials::DelegationKey {
subject_id: "alice".into(),
workload_id: None,
client_id: None,
audience: "https://hr.example.com".into(),
scopes: vec!["read:compensation".into()],
mode: DelegationMode::OnBehalfOfUser,
};
assert!(raw.delegated_tokens.contains_key(&expected_key));
}
#[test]
fn two_agents_do_not_share_one_cache_entry() {
use crate::extensions::security::WorkloadIdentity;
fn workload_exchange(minted: &str) -> DelegationPayload {
let mut p = DelegationPayload::new("svid-bytes", "get_compensation")
.with_subject(DelegationSubject::CallerWorkload);
p.delegated_token = Some(RawDelegatedToken::new(
minted,
"Authorization",
"https://hr.example.com",
vec!["read:compensation".into()],
Utc::now() + chrono::Duration::seconds(300),
));
p.delegation_mode = Some(DelegationMode::AsCallerWorkload);
p
}
fn ext_for(spiffe_id: &str, carry: Option<Arc<RawCredentialsExtension>>) -> Extensions {
Extensions {
security: Some(Arc::new(crate::extensions::SecurityExtension {
caller_workload: Some(WorkloadIdentity {
spiffe_id: Some(spiffe_id.into()),
..Default::default()
}),
..Default::default()
})),
raw_credentials: carry,
..Default::default()
}
}
let after_payroll = workload_exchange("payroll-token")
.apply_to_extensions(ext_for("spiffe://corp/payroll", None));
let carried = after_payroll.raw_credentials.clone();
let after_both = workload_exchange("recruiting-token")
.apply_to_extensions(ext_for("spiffe://corp/recruiting", carried));
let raw = after_both.raw_credentials.as_ref().unwrap();
assert_eq!(
raw.delegated_tokens.len(),
2,
"each calling agent must get its own cache entry; keys: {:?}",
raw.delegated_tokens.keys().collect::<Vec<_>>(),
);
let lookup = |spiffe: &str| {
raw.delegated_tokens
.get(&crate::extensions::raw_credentials::DelegationKey {
subject_id: String::new(),
workload_id: Some(spiffe.into()),
client_id: None,
audience: "https://hr.example.com".into(),
scopes: vec!["read:compensation".into()],
mode: DelegationMode::AsCallerWorkload,
})
.map(|t| (*t.token).clone())
};
assert_eq!(
lookup("spiffe://corp/payroll").as_deref(),
Some("payroll-token"),
);
assert_eq!(
lookup("spiffe://corp/recruiting").as_deref(),
Some("recruiting-token"),
);
}
#[test]
fn two_clients_do_not_share_one_cache_entry() {
use crate::extensions::ClientExtension;
fn client_exchange(minted: &str) -> DelegationPayload {
let mut p = DelegationPayload::new("client-tok", "get_compensation")
.with_subject(DelegationSubject::Client);
p.delegated_token = Some(RawDelegatedToken::new(
minted,
"Authorization",
"https://hr.example.com",
vec!["read:compensation".into()],
Utc::now() + chrono::Duration::seconds(300),
));
p.delegation_mode = Some(DelegationSubject::Client.default_mode());
p
}
fn ext_for(client_id: &str, carry: Option<Arc<RawCredentialsExtension>>) -> Extensions {
Extensions {
security: Some(Arc::new(crate::extensions::SecurityExtension {
client: Some(ClientExtension {
client_id: client_id.into(),
..Default::default()
}),
..Default::default()
})),
raw_credentials: carry,
..Default::default()
}
}
let after_a = client_exchange("token-a").apply_to_extensions(ext_for("client-a", None));
let carried = after_a.raw_credentials.clone();
let after_both =
client_exchange("token-b").apply_to_extensions(ext_for("client-b", carried));
let raw = after_both.raw_credentials.as_ref().unwrap();
assert_eq!(
raw.delegated_tokens.len(),
2,
"each client must get its own cache entry; keys: {:?}",
raw.delegated_tokens.keys().collect::<Vec<_>>(),
);
let lookup = |client_id: &str| {
raw.delegated_tokens
.get(
&crate::extensions::raw_credentials::DelegationKey::new(
DelegationMode::AsClient,
"https://hr.example.com",
vec!["read:compensation".into()],
)
.with_client_id(Some(client_id.into())),
)
.map(|t| (*t.token).clone())
};
assert_eq!(lookup("client-a").as_deref(), Some("token-a"));
assert_eq!(lookup("client-b").as_deref(), Some("token-b"));
}
#[test]
fn apply_to_extensions_respects_explicit_delegation_mode() {
let mut p = DelegationPayload::new("tok", "tool");
p.delegated_token = Some(RawDelegatedToken::new(
"this-workload-token",
"Authorization",
"https://downstream.example.com",
vec!["service:call".into()],
Utc::now(),
));
p.delegation_mode =
Some(crate::extensions::raw_credentials::DelegationMode::AsThisWorkload);
let updated = p.apply_to_extensions(Extensions::default());
let raw = updated.raw_credentials.as_ref().unwrap();
let key = raw.delegated_tokens.keys().next().unwrap();
assert!(matches!(
key.mode,
crate::extensions::raw_credentials::DelegationMode::AsThisWorkload
));
}
#[test]
fn apply_to_extensions_defaults_delegation_mode_when_unset() {
let mut p = DelegationPayload::new("tok", "tool");
p.delegated_token = Some(RawDelegatedToken::new(
"user-token",
"Authorization",
"https://aud.example.com",
vec!["read".into()],
Utc::now(),
));
let updated = p.apply_to_extensions(Extensions::default());
let raw = updated.raw_credentials.as_ref().unwrap();
let key = raw.delegated_tokens.keys().next().unwrap();
assert!(matches!(
key.mode,
crate::extensions::raw_credentials::DelegationMode::OnBehalfOfUser
));
}
#[test]
fn merge_threads_delegation_mode_through_chain() {
let mut base = DelegationPayload::new("tok", "tool");
let mut overlay = DelegationPayload::new("", "");
overlay.delegation_mode =
Some(crate::extensions::raw_credentials::DelegationMode::AsThisWorkload);
base.merge(overlay);
assert!(matches!(
base.delegation_mode,
Some(crate::extensions::raw_credentials::DelegationMode::AsThisWorkload)
));
}
#[test]
fn apply_to_extensions_falls_back_to_empty_subject_id_when_no_subject() {
let mut p = DelegationPayload::new("tok", "tool");
p.delegated_token = Some(RawDelegatedToken::new(
"minted",
"Authorization",
"aud",
vec![],
Utc::now(),
));
let updated = p.apply_to_extensions(Extensions::default());
let raw = updated.raw_credentials.as_ref().unwrap();
let key = raw.delegated_tokens.keys().next().unwrap();
assert_eq!(key.subject_id, "");
}
#[test]
fn auth_enforced_by_defaults_to_caller() {
let p = DelegationPayload::new("tok", "tool");
assert_eq!(p.auth_enforced_by(), AuthEnforcedBy::Caller);
}
#[test]
fn target_type_defaults_to_tool() {
let p = DelegationPayload::new("tok", "tool");
assert_eq!(p.target_type(), &TargetType::Tool);
}
}