use std::sync::Arc;
use base64::Engine as _;
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
use polyc_crypto::session;
use polyc_crypto::signing_role::{RoleTrustSet, SigningRole as _, TurnReadRole};
use polyc_persona::ScopeResolution;
use crate::principal::{
AdminFleetPrincipal, AdminPrincipal, ControlFleetPrincipal, ConversationGrantPrincipal,
PersonaAccess, PersonaPrincipal, Principal, PrincipalError, SEARCH_SCOPE_CAP, Scoping,
SearchScope, SearchScopeError, canonical_search_scope,
};
use crate::session::{MemorySources, QueryScope};
use polyc_query_model::{GRANT_KIND, GrantClaims, GrantSubject};
#[async_trait::async_trait]
pub trait SessionVerification: Send + Sync {
async fn verify_bearer(
&self,
token: &str,
now_ms: u64,
) -> Result<
polyc_crypto::session::AuthorizedSessionClaims,
polyc_session_family::authority::SessionAuthorityError,
>;
}
struct FullAuthority(Arc<dyn polyc_session_family::authority::BrowserSessionAuthority>);
#[async_trait::async_trait]
impl SessionVerification for FullAuthority {
async fn verify_bearer(
&self,
token: &str,
now_ms: u64,
) -> Result<
polyc_crypto::session::AuthorizedSessionClaims,
polyc_session_family::authority::SessionAuthorityError,
> {
self.0.verify_bearer(token, now_ms).await
}
}
pub struct CredentialAuthority {
persona: PersonaAccess,
bearer_authority: Option<Arc<dyn SessionVerification>>,
turn_read_trust: RoleTrustSet<TurnReadRole>,
#[cfg(any(test, feature = "test-util"))]
legacy_revoked: Option<Arc<polyc_crypto::session::RevokedTokens>>,
#[cfg(any(test, feature = "test-util"))]
legacy_session_trust: Option<RoleTrustSet<polyc_crypto::signing_role::SessionRole>>,
#[cfg(any(test, feature = "test-util"))]
test_search_scope_cap: std::sync::Mutex<Option<usize>>,
}
impl std::fmt::Debug for CredentialAuthority {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("CredentialAuthority")
.field("bearer_authority", &self.bearer_authority.is_some())
.finish_non_exhaustive()
}
}
impl CredentialAuthority {
async fn verified_memory_participants(
&self,
conversation_id: &str,
proposed_sources: &[String],
) -> Result<Vec<String>, PrincipalError> {
if proposed_sources.is_empty() {
return Ok(Vec::new());
}
let Some(persona_host) = self.persona.load_full() else {
return Err(PrincipalError::StoreUnavailable);
};
let mut participants = Vec::with_capacity(proposed_sources.len());
for proposed in proposed_sources {
let active = persona_host
.active_persona(proposed.clone())
.await
.map_err(|_store_error| PrincipalError::StoreUnavailable)?
.ok_or(PrincipalError::SourceOutsideScope)?;
let participations = persona_host
.participations(active.persona_id.clone())
.await
.map_err(|_store_error| PrincipalError::StoreUnavailable)?;
if !participations
.iter()
.any(|participation| participation.conversation_id == conversation_id)
{
return Err(PrincipalError::SourceOutsideScope);
}
participants.push(active.persona_id);
}
participants.sort();
participants.dedup();
Ok(participants)
}
pub async fn composite_trace_memory_scope(
&self,
scoping: &Scoping,
proposed_sources: &[String],
) -> Result<QueryScope, PrincipalError> {
if scoping.composite_trace().is_none() {
return Err(PrincipalError::SourceOutsideScope);
}
let conversation_id = scoping
.conversation_id()
.ok_or(PrincipalError::SourceOutsideScope)?;
let owner = scoping
.caller_identity()
.ok_or(PrincipalError::SourceOutsideScope)?;
let proposed = proposed_sources
.iter()
.filter(|source| source.as_str() != owner)
.cloned()
.collect::<Vec<_>>();
let participants = self
.verified_memory_participants(conversation_id, &proposed)
.await?;
let QueryScope::Conversations { conversations, .. } = scoping.scope() else {
return Err(PrincipalError::SourceOutsideScope);
};
Ok(QueryScope::Conversations {
conversations: conversations.clone(),
memory: MemorySources {
owner: Some(owner.to_owned()),
participants,
},
})
}
#[must_use]
pub fn current(
persona: Arc<dyn crate::principal::PersonaSource>,
turn_read_trust: RoleTrustSet<TurnReadRole>,
authority: Arc<dyn polyc_session_family::authority::BrowserSessionAuthority>,
) -> Self {
Self::verifying(persona, turn_read_trust, Arc::new(FullAuthority(authority)))
}
#[must_use]
pub fn verifying(
persona: Arc<dyn crate::principal::PersonaSource>,
turn_read_trust: RoleTrustSet<TurnReadRole>,
sessions: Arc<dyn SessionVerification>,
) -> Self {
Self {
persona: PersonaAccess::current(persona),
bearer_authority: Some(sessions),
turn_read_trust,
#[cfg(any(test, feature = "test-util"))]
legacy_revoked: None,
#[cfg(any(test, feature = "test-util"))]
legacy_session_trust: None,
#[cfg(any(test, feature = "test-util"))]
test_search_scope_cap: std::sync::Mutex::new(None),
}
}
#[cfg(any(test, feature = "test-util"))]
#[must_use]
pub fn legacy(
persona: crate::principal::PersonaSourceHandle,
revoked: Arc<polyc_crypto::session::RevokedTokens>,
turn_read_trust: RoleTrustSet<TurnReadRole>,
legacy_session_trust: RoleTrustSet<polyc_crypto::signing_role::SessionRole>,
) -> Self {
Self {
persona,
bearer_authority: None,
turn_read_trust,
legacy_revoked: Some(revoked),
legacy_session_trust: Some(legacy_session_trust),
#[cfg(any(test, feature = "test-util"))]
test_search_scope_cap: std::sync::Mutex::new(None),
}
}
#[cfg(any(test, feature = "test-util"))]
pub fn replace_turn_read_trust_for_test(&mut self, trust: RoleTrustSet<TurnReadRole>) {
self.turn_read_trust = trust;
}
#[cfg(any(test, feature = "test-util"))]
pub fn set_test_search_scope_cap(&self, cap: usize) {
*self.test_search_scope_cap.lock().expect("poison") = Some(cap);
}
pub async fn verify_admin_session(
&self,
token: &str,
now_ms: u64,
) -> Result<Principal, PrincipalError> {
let claims = if let Some(authority) = &self.bearer_authority {
authority
.verify_bearer(token, now_ms)
.await
.map_err(|error| match error {
polyc_session_family::authority::SessionAuthorityError::Invalid => {
PrincipalError::InvalidSession
}
_ => PrincipalError::StoreUnavailable,
})?
} else {
#[cfg(not(any(test, feature = "test-util")))]
return Err(PrincipalError::StoreUnavailable);
#[cfg(any(test, feature = "test-util"))]
{
let legacy = session::verify_session_with_trust(
self.legacy_session_trust
.as_ref()
.ok_or(PrincipalError::StoreUnavailable)?,
token,
now_ms,
self.legacy_revoked
.as_ref()
.ok_or(PrincipalError::StoreUnavailable)?,
)
.ok_or(PrincipalError::InvalidSession)?;
polyc_crypto::session::AuthorizedSessionClaims {
issuer: legacy.issuer,
key_id: legacy.key_id,
session_id: "legacy-test-session".to_owned(),
authorization_epoch: 0,
subject: legacy.subject,
scopes: legacy.scopes,
issued_ms: legacy.issued_ms,
expires_ms: legacy.expires_ms,
}
}
};
if !claims.has_scope(session::SessionScope::ExplorerRead) {
return Err(PrincipalError::InvalidSession);
}
let persona_id = claims
.subject
.persona_id()
.ok_or(PrincipalError::InvalidSession)?
.to_owned();
let Some(persona) = self.persona.load_full() else {
return Err(PrincipalError::StoreUnavailable);
};
match persona.active_persona(persona_id).await {
Ok(Some(active)) => {
let persona_id = active.persona_id;
if active.admin {
Ok(Principal::Admin(AdminPrincipal::minted(persona_id)))
} else {
Ok(Principal::Persona(PersonaPrincipal::minted(persona_id)))
}
}
Ok(None) => Err(PrincipalError::NotAuthorizedForFleet),
Err(_store_error) => Err(PrincipalError::StoreUnavailable),
}
}
pub fn verify_conversation_grant(
&self,
token: &str,
now_unix_ms: u64,
) -> Result<Principal, PrincipalError> {
let (claims_b64, sig_b64) = token.split_once('.').ok_or_else(|| {
tracing::warn!("conversation grant token malformed: no `.` separator");
PrincipalError::InvalidGrant
})?;
let canonical = URL_SAFE_NO_PAD.decode(claims_b64).map_err(|_| {
tracing::warn!("conversation grant token malformed: claims segment not base64");
PrincipalError::InvalidGrant
})?;
let signature = URL_SAFE_NO_PAD.decode(sig_b64).map_err(|_| {
tracing::warn!("conversation grant token malformed: signature segment not base64");
PrincipalError::InvalidGrant
})?;
let claims: GrantClaims = serde_json::from_slice(&canonical).map_err(|_| {
tracing::warn!("conversation grant token malformed: claims did not decode as JSON");
PrincipalError::InvalidGrant
})?;
if claims.issuer() != TurnReadRole::ISSUER
|| !self.turn_read_trust.verify_turn_read_capability(
claims.key_id(),
&canonical,
&signature,
)
{
tracing::warn!("conversation grant token signature invalid");
return Err(PrincipalError::InvalidGrant);
}
if claims.kind() != GRANT_KIND {
tracing::warn!(kind = %claims.kind(), "conversation grant token kind tag mismatch");
return Err(PrincipalError::InvalidGrant);
}
if now_unix_ms > claims.expires_at_ms() {
tracing::warn!("conversation grant token expired");
return Err(PrincipalError::GrantExpired);
}
let (conversation_id, subject, memory_sources) = claims.into_parts();
if let Some((purpose, statement_digest)) = subject.control_fleet() {
let digest = polyc_crypto::hex::decode(statement_digest)
.and_then(|bytes| <[u8; 32]>::try_from(bytes).ok())
.ok_or_else(|| {
tracing::warn!("control-fleet grant statement digest was not 32 hex bytes");
PrincipalError::InvalidGrant
})?;
return Ok(Principal::ControlFleet(ControlFleetPrincipal::minted(
purpose.to_owned(),
digest,
)));
}
if let Some((admin_persona, session)) = subject.admin_fleet() {
return Ok(Principal::AdminFleet(AdminFleetPrincipal::minted(
admin_persona.to_owned(),
session.to_owned(),
)));
}
if conversation_id.is_empty()
&& let GrantSubject::Persona(persona_id) = &subject
{
return Ok(Principal::Persona(PersonaPrincipal::minted(
persona_id.clone(),
)));
}
let composite_digests = subject
.composite_trace()
.map(|(_persona_id, trace, memory, routine)| vec![trace, memory, routine])
.or_else(|| {
subject
.admin_composite_trace()
.map(|(trace, routine, addresses)| vec![trace, routine, addresses])
});
if let Some(digests) = composite_digests {
for digest in digests {
let valid = polyc_crypto::hex::decode(digest)
.and_then(|bytes| <[u8; 32]>::try_from(bytes).ok())
.is_some();
if !valid {
tracing::warn!("composite-trace statement digest was not 32 hex bytes");
return Err(PrincipalError::InvalidGrant);
}
}
}
if subject
.composite_trace()
.is_some_and(|(persona_id, ..)| persona_id.is_empty())
{
tracing::warn!("a composite-trace grant named an empty persona");
return Err(PrincipalError::InvalidGrant);
}
if subject.admin_composite_trace().is_some() && !memory_sources.is_empty() {
tracing::warn!("an admin composite-trace grant proposed persona-memory sources");
return Err(PrincipalError::InvalidGrant);
}
Ok(Principal::ConversationGrant(
ConversationGrantPrincipal::minted(conversation_id, subject, memory_sources),
))
}
#[allow(
clippy::too_many_lines,
reason = "one match arm per Principal variant, each spelling out its own Scoping literal in full"
)]
pub async fn scoping_for(&self, principal: &Principal) -> Result<Scoping, PrincipalError> {
let scoping = match principal {
Principal::Admin(admin) => Scoping::admin(admin.persona_id()),
Principal::ConversationGrant(grant) => {
let owner = grant.subject().owner_persona_id().map(str::to_owned);
let participants = self
.verified_memory_participants(grant.conversation_id(), grant.memory_sources())
.await?;
let mut scoping =
Scoping::for_conversation(grant.conversation_id(), grant.subject())?;
let memory = MemorySources {
owner,
participants,
};
let scope = match scoping.scope() {
QueryScope::Conversations { conversations, .. } => QueryScope::Conversations {
conversations: conversations.clone(),
memory,
},
QueryScope::Fleet => QueryScope::Fleet,
QueryScope::FleetConversation { conversation } => {
QueryScope::FleetConversation {
conversation: conversation.clone(),
}
}
};
scoping.set_scope(scope);
scoping
}
Principal::Persona(persona) => {
let Some(persona_host) = self.persona.load_full() else {
return Err(PrincipalError::StoreUnavailable);
};
let participations = persona_host
.participations(persona.persona_id().to_owned())
.await
.map_err(|_store_error| PrincipalError::StoreUnavailable)?;
let conversation_ids = participations
.into_iter()
.map(|participation| participation.conversation_id)
.collect();
Scoping::persona(persona.persona_id(), conversation_ids)
}
Principal::ControlFleet(control_fleet) => {
Scoping::from_control_fleet(control_fleet.clone())
}
Principal::AdminFleet(admin_fleet) => {
let Some(persona_host) = self.persona.load_full() else {
return Err(PrincipalError::StoreUnavailable);
};
match persona_host
.active_persona(admin_fleet.admin_persona().to_owned())
.await
{
Ok(Some(active)) if active.admin => {
Scoping::from_admin_fleet(active.persona_id, admin_fleet.clone())
}
Ok(Some(_) | None) => return Err(PrincipalError::NotAuthorizedForFleet),
Err(_store_error) => return Err(PrincipalError::StoreUnavailable),
}
}
};
Ok(scoping)
}
pub async fn own_rows_scoping(&self, persona_id: &str) -> Result<Scoping, PrincipalError> {
let Some(persona_host) = self.persona.load_full() else {
return Err(PrincipalError::StoreUnavailable);
};
let participations = persona_host
.participations(persona_id.to_owned())
.await
.map_err(|_store_error| PrincipalError::StoreUnavailable)?;
let conversation_ids = participations
.into_iter()
.map(|participation| participation.conversation_id)
.collect();
Ok(Scoping::own_rows(persona_id, conversation_ids))
}
pub async fn authorize_conversation_read(
&self,
principal: &Principal,
) -> Result<(), PrincipalError> {
if matches!(principal, Principal::ControlFleet(_)) {
return Err(PrincipalError::ControlFleetNotUsableEmbedded);
}
if matches!(principal, Principal::AdminFleet(_)) {
return Err(PrincipalError::AdminFleetNotUsableEmbedded);
}
self.scoping_for(principal).await?;
Ok(())
}
pub async fn resolve_search_scope(
&self,
principal_ref: &str,
conversation_id: &str,
turn_id: &str,
) -> Result<SearchScope, SearchScopeError> {
let Some(persona_host) = self.persona.load_full() else {
return Err(SearchScopeError::StoreUnavailable);
};
match persona_host.active_persona(principal_ref.to_owned()).await {
Ok(Some(_active)) => {}
Ok(None) => {
tracing::info!(
persona_id = %principal_ref,
conversation_id = %conversation_id,
turn_id = %turn_id,
"refusing search-scope resolution: persona is not active"
);
return Err(SearchScopeError::PersonaNotActive);
}
Err(_store_error) => {
return Err(SearchScopeError::StoreUnavailable);
}
}
#[cfg(any(test, feature = "test-util"))]
let cap = self
.test_search_scope_cap
.lock()
.expect("poison")
.unwrap_or(SEARCH_SCOPE_CAP);
#[cfg(not(any(test, feature = "test-util")))]
let cap = SEARCH_SCOPE_CAP;
let resolution = persona_host
.participation_scope(principal_ref.to_owned(), cap)
.await
.map_err(|_store_error| SearchScopeError::StoreUnavailable)?;
let mut conversation_ids = match resolution {
ScopeResolution::RefusedOverCap { count } => {
tracing::warn!(
persona_id = %principal_ref,
conversation_id = %conversation_id,
turn_id = %turn_id,
count,
cap,
"refusing search-scope resolution: participation count exceeds the search cap"
);
return Err(SearchScopeError::OverCap { count });
}
ScopeResolution::Resolved { conversation_ids } => conversation_ids,
};
conversation_ids.retain(|id| id != conversation_id);
let (conversation_ids, hash) = canonical_search_scope(conversation_ids);
Ok(SearchScope::minted(conversation_ids, hash))
}
}
pub trait UnixClock: Send + Sync {
fn now_unix_ms(&self) -> u64;
}
pub struct SystemUnixClock;
impl UnixClock for SystemUnixClock {
fn now_unix_ms(&self) -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |elapsed| {
u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX)
})
}
}
pub enum PresentedCredential {
Bearer(String),
ConversationGrant(String),
}
impl std::fmt::Debug for PresentedCredential {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(match self {
Self::Bearer(_) => "bearer",
Self::ConversationGrant(_) => "conversation-grant",
})
}
}
pub struct CredentialWitness {
credential: PresentedCredential,
authority: Arc<CredentialAuthority>,
clock: Arc<dyn UnixClock>,
composite_trace_memory_sources: Option<Vec<String>>,
composite_trace_routine: bool,
composite_trace_addresses: bool,
}
impl std::fmt::Debug for CredentialWitness {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("CredentialWitness")
.field("credential", &self.credential)
.finish_non_exhaustive()
}
}
impl CredentialWitness {
pub async fn admit_bearer(
token: String,
authority: Arc<CredentialAuthority>,
clock: Arc<dyn UnixClock>,
) -> Result<(Self, Scoping), PrincipalError> {
let now = clock.now_unix_ms();
match authority.verify_conversation_grant(&token, now) {
Ok(_grant) => {
Self::admit(
PresentedCredential::ConversationGrant(token),
authority,
clock,
)
.await
}
Err(PrincipalError::GrantExpired) => Err(PrincipalError::GrantExpired),
Err(_not_a_grant) => {
Self::admit(PresentedCredential::Bearer(token), authority, clock).await
}
}
}
pub async fn admit_grant_only(
token: String,
authority: Arc<CredentialAuthority>,
clock: Arc<dyn UnixClock>,
) -> Result<(Self, Scoping), PrincipalError> {
Self::admit(
PresentedCredential::ConversationGrant(token),
authority,
clock,
)
.await
}
pub async fn admit(
credential: PresentedCredential,
authority: Arc<CredentialAuthority>,
clock: Arc<dyn UnixClock>,
) -> Result<(Self, Scoping), PrincipalError> {
let scoping = Self::verify(&credential, &authority, clock.as_ref()).await?;
Ok((
Self {
credential,
authority,
clock,
composite_trace_memory_sources: None,
composite_trace_routine: false,
composite_trace_addresses: false,
},
scoping,
))
}
pub fn retain_composite_trace_memory_sources(&mut self, sources: Vec<String>) {
self.composite_trace_memory_sources = Some(sources);
}
pub const fn retain_composite_trace_routine(&mut self) {
self.composite_trace_routine = true;
}
pub const fn retain_composite_trace_addresses(&mut self) {
self.composite_trace_addresses = true;
}
async fn verify(
credential: &PresentedCredential,
authority: &CredentialAuthority,
clock: &dyn UnixClock,
) -> Result<Scoping, PrincipalError> {
let now = clock.now_unix_ms();
let principal = match credential {
PresentedCredential::Bearer(token) => {
authority.verify_admin_session(token, now).await?
}
PresentedCredential::ConversationGrant(token) => {
authority.verify_conversation_grant(token, now)?
}
};
authority.scoping_for(&principal).await
}
pub async fn current_scope(&self) -> Result<QueryScope, PrincipalError> {
let scoping = Self::verify(&self.credential, &self.authority, self.clock.as_ref()).await?;
if let Some(sources) = &self.composite_trace_memory_sources {
self.authority
.composite_trace_memory_scope(&scoping, sources)
.await
} else if self.composite_trace_routine {
if scoping.composite_trace().is_none() {
Err(PrincipalError::SourceOutsideScope)
} else {
Ok(QueryScope::Fleet)
}
} else if self.composite_trace_addresses {
let addresses = scoping
.composite_trace()
.and_then(crate::principal::CompositeTracePrincipal::address_statement_digest);
match (addresses, scoping.conversation_id()) {
(Some(_), Some(conversation)) => Ok(QueryScope::FleetConversation {
conversation: conversation.to_owned(),
}),
_ => Err(PrincipalError::SourceOutsideScope),
}
} else {
Ok(scoping.into_parts().scope)
}
}
}