use super::BearerToken;
use crate::capability::error::CapabilityError;
use futures::Stream;
use otel_arrow_dfe_engine_macros::capability;
use serde_json::{Map, Value};
use std::pin::Pin;
use std::sync::Arc;
#[derive(Clone, Debug)]
pub struct AgentFedCredentialSnapshot {
token: BearerToken,
attributes: Arc<Map<String, Value>>,
}
impl AgentFedCredentialSnapshot {
#[must_use]
pub fn new(token: BearerToken, attributes: Arc<Map<String, Value>>) -> Self {
Self { token, attributes }
}
#[must_use]
pub const fn token(&self) -> &BearerToken {
&self.token
}
#[must_use]
pub fn attributes(&self) -> &Map<String, Value> {
&self.attributes
}
}
pub type AgentFedCredentialSnapshotStream =
Pin<Box<dyn Stream<Item = AgentFedCredentialSnapshot> + 'static>>;
#[capability(
name = "agent_fed_credential_provider",
description = "Provides one atomic bearer-token and vendor-attribute snapshot"
)]
pub trait AgentFedCredentialProvider {
async fn get_credential(&self) -> Result<Arc<AgentFedCredentialSnapshot>, CapabilityError>;
fn credential_stream(&self) -> AgentFedCredentialSnapshotStream;
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn snapshot_keeps_token_and_attributes_together() {
let attributes = Arc::new(
json!({"endpoint": "https://ingest.example"})
.as_object()
.cloned()
.expect("object"),
);
let snapshot = AgentFedCredentialSnapshot::new(
BearerToken::without_expiry("secret-token".to_owned()),
Arc::clone(&attributes),
);
assert_eq!(snapshot.token().expose_token(), "secret-token");
assert_eq!(snapshot.attributes(), attributes.as_ref());
assert!(!format!("{snapshot:?}").contains("secret-token"));
}
}