use std::collections::btree_map::Entry;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::SystemTime;
use gateway_core::{CircuitState, FailoverTarget};
use crate::backends::catalog::CatalogContent;
use crate::convergence::ResolvedSecrets;
use crate::desired_state::credentials::{CredentialError, Credentials};
use crate::desired_state::models::{
CatalogOffering, ModelEnablement, ModelError, ModelOwner, Models, OfferingId,
};
use crate::desired_state::policy::{PolicyError, PolicyScope, PolicySet};
use crate::desired_state::providers::{ProviderError, Providers};
use crate::desired_state::secrets::{SecretLifecycle, SecretRef};
use crate::desired_state::{Checksum, DesiredState};
use super::dimensions::{
CataloguePresence, Enablement, Entitlement, PolicyDecision, RuntimeHealth,
};
use super::discovery::DiscoveryObservation;
use super::index::{AvailabilityIndex, AvailabilityIndexBuilder, AvailabilityRecord};
use super::refs::{AvailabilityKey, CredentialRef, ScopeRef, TargetRef};
use super::store::{self, EvidenceClear, EvidenceWrite, StoredObservation};
use super::verdict::Availability;
#[derive(Debug, thiserror::Error)]
pub enum AvailabilityProjectionError {
#[error("the revision's model contracts could not be read: {0}")]
Models(#[from] ModelError),
#[error("the revision's provider connections could not be read: {0}")]
Providers(#[from] ProviderError),
#[error("the revision's credentials could not be read: {0}")]
Credentials(#[from] CredentialError),
#[error("the revision's policy documents could not be read: {0}")]
Policy(#[from] PolicyError),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CatalogueListing {
snapshot: Checksum,
offerings: BTreeMap<OfferingId, TargetRef>,
unnamed: usize,
}
impl CatalogueListing {
pub fn of(snapshot: Checksum, content: &CatalogContent) -> Self {
let mut offerings = BTreeMap::new();
let mut unnamed = 0;
for model in content.models() {
for offering in &model.offerings {
let provider = offering.provider.as_str();
let published = offering.published_model_id.as_str();
let (Ok(identity), Ok(target)) = (
OfferingId::of(provider, published),
TargetRef::parse(provider, published),
) else {
unnamed += 1;
continue;
};
offerings.insert(identity, target);
}
}
Self {
snapshot,
offerings,
unnamed,
}
}
pub const fn snapshot(&self) -> Checksum {
self.snapshot
}
pub fn target(&self, offering: OfferingId) -> Option<&TargetRef> {
self.offerings.get(&offering)
}
pub fn len(&self) -> usize {
self.offerings.len()
}
pub fn is_empty(&self) -> bool {
self.offerings.is_empty()
}
pub const fn unnamed(&self) -> usize {
self.unnamed
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Catalogue {
active: CatalogueListing,
superseded: BTreeMap<Checksum, CatalogueListing>,
}
impl Catalogue {
pub fn active(active: CatalogueListing) -> Self {
Self {
active,
superseded: BTreeMap::new(),
}
}
#[must_use]
pub fn with_superseded(mut self, listing: CatalogueListing) -> Self {
self.superseded.insert(listing.snapshot(), listing);
self
}
fn presence(&self, pinned: CatalogOffering) -> Option<(TargetRef, CataloguePresence)> {
if let Some(target) = self.active.target(pinned.offering) {
return Some((target.clone(), CataloguePresence::Present));
}
let named = self
.superseded
.get(&pinned.snapshot)
.and_then(|listing| listing.target(pinned.offering))?;
Some((named.clone(), CataloguePresence::Withdrawn))
}
fn is_current(&self, pinned: CatalogOffering) -> bool {
pinned.is_pinned_to(self.active.snapshot())
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct CredentialReadiness {
resolved: BTreeSet<SecretRef>,
}
impl CredentialReadiness {
pub fn none() -> Self {
Self::default()
}
pub fn of(secrets: &ResolvedSecrets) -> Self {
Self {
resolved: secrets.references().into_iter().collect(),
}
}
#[must_use]
pub fn holding(mut self, secret: SecretRef) -> Self {
self.resolved.insert(secret);
self
}
fn holds(&self, secret: SecretRef) -> bool {
self.resolved.contains(&secret)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RuntimeObservations {
health: BTreeMap<String, RuntimeHealth>,
}
impl RuntimeObservations {
pub fn none() -> Self {
Self::default()
}
pub fn of_circuits(circuits: impl IntoIterator<Item = (String, CircuitState)>) -> Self {
Self {
health: circuits
.into_iter()
.map(|(target, state)| {
let health = match state {
CircuitState::Closed => RuntimeHealth::Healthy,
CircuitState::HalfOpen => RuntimeHealth::Impaired,
CircuitState::Open => RuntimeHealth::Unavailable,
};
(target, health)
})
.collect(),
}
}
fn health(&self, target: &TargetRef) -> RuntimeHealth {
self.health
.get(&Self::circuit_key(target))
.copied()
.unwrap_or(RuntimeHealth::Unobserved)
}
pub(crate) fn circuit_key(target: &TargetRef) -> String {
FailoverTarget::new(target.provider.as_str(), target.model.as_str()).qualified_model()
}
}
pub const REPORTED_KEYS: usize = 32;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProjectedAvailability {
derivation: u64,
index: AvailabilityIndex,
unnameable: usize,
undescribed: usize,
orphaned: Vec<AvailabilityKey>,
skewed: usize,
superseded: usize,
misfiled: usize,
undescribed_looks: usize,
undescribed_look_keys: Vec<AvailabilityKey>,
conflicting: usize,
conflicted: Vec<AvailabilityKey>,
}
impl ProjectedAvailability {
pub const fn derivation(&self) -> u64 {
self.derivation
}
pub const fn index(&self) -> &AvailabilityIndex {
&self.index
}
pub fn into_index(self) -> AvailabilityIndex {
self.index
}
pub const fn unnameable(&self) -> usize {
self.unnameable
}
pub const fn undescribed(&self) -> usize {
self.undescribed
}
pub fn orphaned(&self) -> &[AvailabilityKey] {
&self.orphaned
}
pub const fn skewed(&self) -> usize {
self.skewed
}
pub const fn superseded(&self) -> usize {
self.superseded
}
pub const fn misfiled(&self) -> usize {
self.misfiled
}
pub const fn undescribed_looks(&self) -> usize {
self.undescribed_looks
}
pub fn undescribed_look_keys(&self) -> &[AvailabilityKey] {
&self.undescribed_look_keys
}
pub const fn conflicting(&self) -> usize {
self.conflicting
}
pub fn conflicted(&self) -> &[AvailabilityKey] {
&self.conflicted
}
}
pub struct AvailabilityProjection<'a> {
catalogue: &'a Catalogue,
readiness: &'a CredentialReadiness,
}
impl<'a> AvailabilityProjection<'a> {
pub const fn new(catalogue: &'a Catalogue, readiness: &'a CredentialReadiness) -> Self {
Self {
catalogue,
readiness,
}
}
pub fn project(
&self,
state: &DesiredState,
previous: &AvailabilityIndex,
observations: impl IntoIterator<Item = DiscoveryObservation>,
) -> Result<ProjectedAvailability, AvailabilityProjectionError> {
let models = Models::of(state)?;
let providers = Providers::of(state)?;
let credentials = Credentials::of(state)?;
let policies = PolicySet::of(state)?;
let mut described = BTreeSet::new();
let mut declarations: BTreeMap<AvailabilityKey, AvailabilityRecord> = BTreeMap::new();
let mut unnameable = 0;
let mut skewed = 0;
let mut conflicting = 0;
let mut conflicted: BTreeSet<AvailabilityKey> = BTreeSet::new();
for enablement in models.enablements() {
let pinned = enablement.body.offering();
let Some((target, presence)) = self.catalogue.presence(pinned) else {
unnameable += 1;
continue;
};
if !self.catalogue.is_current(pinned) {
skewed += 1;
}
let owner = enablement.body.owner();
let scope = scope_of(owner);
let (entitlement, credential) =
self.entitlement(&providers, &credentials, owner, &target);
let record = AvailabilityRecord {
presence,
enablement: enablement_of(enablement),
entitlement,
policy: policy_of(&policies, owner),
credential,
..AvailabilityRecord::default()
};
let key = AvailabilityKey::new(scope, target);
described.insert(key.clone());
match declarations.entry(key) {
Entry::Vacant(slot) => {
slot.insert(record);
}
Entry::Occupied(mut held) => {
conflicting += 1;
note(&mut conflicted, held.key().clone());
let combined = least_permissive(held.get(), record);
held.insert(combined);
}
}
}
let mut builder = AvailabilityIndexBuilder::carrying_evidence_for(previous, &described);
for (key, record) in declarations {
builder = builder.record(key, record);
}
let mut undescribed_looks = 0;
let mut undescribed_look_keys: BTreeSet<AvailabilityKey> = BTreeSet::new();
for observation in observations {
if described.contains(&observation.key()) {
builder = builder.observe(observation);
} else {
undescribed_looks += 1;
note(&mut undescribed_look_keys, observation.key());
}
}
let orphaned: Vec<AvailabilityKey> = previous
.records()
.filter(|(key, record)| record.holds_evidence() && !described.contains(*key))
.map(|(key, _)| key.clone())
.collect();
let undescribed = orphaned.len();
Ok(ProjectedAvailability {
derivation: 0,
unnameable,
undescribed,
orphaned,
skewed,
superseded: builder.superseded(),
misfiled: builder.misfiled(),
undescribed_looks,
undescribed_look_keys: reported(undescribed_look_keys),
conflicting,
conflicted: reported(conflicted),
index: builder.build(),
})
}
fn entitlement(
&self,
providers: &Providers,
credentials: &Credentials,
owner: ModelOwner,
target: &TargetRef,
) -> (Entitlement, Option<CredentialRef>) {
let connections: BTreeSet<_> = providers
.all()
.filter(|provider| provider.slug.as_str() == target.provider.as_str())
.filter(|provider| owner.reaches(owner_of_provider(provider)))
.map(|provider| provider.body.provider())
.collect();
if connections.is_empty() {
return (Entitlement::Missing, None);
}
let mut best: Option<(Entitlement, Option<CredentialRef>)> = None;
for credential in credentials.all() {
let body = &credential.body;
if !connections.contains(&body.provider()) {
continue;
}
let holder = ModelOwner {
tenant: body.owner().tenant,
project: body.owner().project,
};
if !owner.reaches(holder) {
continue;
}
let entitlement = match body.lifecycle() {
SecretLifecycle::Active if self.readiness.holds(body.secret()) => {
Entitlement::Granted
}
SecretLifecycle::Active | SecretLifecycle::Staged => Entitlement::Unknown,
SecretLifecycle::Disabled
| SecretLifecycle::Revoked
| SecretLifecycle::Tombstoned => Entitlement::Revoked,
};
let reference = CredentialRef::parse(credential.slug.as_str()).ok();
best = Some(match best {
Some(held) if rank(held.0) >= rank(entitlement) => held,
_ => (entitlement, reference),
});
}
best.unwrap_or((Entitlement::Missing, None))
}
}
#[derive(Debug)]
pub struct AvailabilityEvidence {
catalogue: Mutex<Arc<Catalogue>>,
index: Mutex<Arc<AvailabilityIndex>>,
pending: Mutex<Vec<DiscoveryObservation>>,
orphaned: Mutex<BTreeMap<AvailabilityKey, SystemTime>>,
derived_from: Mutex<Option<(Arc<DesiredState>, CredentialReadiness)>>,
deriving: Mutex<()>,
replaced: Mutex<Option<Superseded>>,
derivations: Mutex<u64>,
}
#[derive(Debug)]
struct Superseded {
derivation: u64,
index: Arc<AvailabilityIndex>,
orphaned: BTreeMap<AvailabilityKey, SystemTime>,
derived_from: Option<(Arc<DesiredState>, CredentialReadiness)>,
looks: Vec<DiscoveryObservation>,
}
fn latest_evidence_at(record: &AvailabilityRecord) -> Option<SystemTime> {
record
.discovery
.iter()
.chain(record.last_known_good.iter())
.map(|observation| observation.observed_at)
.chain(record.definitive_at)
.max()
}
impl AvailabilityEvidence {
pub fn new(catalogue: Catalogue) -> Self {
Self {
catalogue: Mutex::new(Arc::new(catalogue)),
index: Mutex::new(Arc::new(AvailabilityIndex::empty())),
pending: Mutex::new(Vec::new()),
orphaned: Mutex::new(BTreeMap::new()),
derived_from: Mutex::new(None),
deriving: Mutex::new(()),
replaced: Mutex::new(None),
derivations: Mutex::new(0),
}
}
pub fn refresh(&self, catalogue: Catalogue) {
*self.lock(&self.catalogue) = Arc::new(catalogue);
}
pub fn observe(&self, observation: DiscoveryObservation) {
self.lock(&self.pending).push(observation);
}
pub fn index(&self) -> Arc<AvailabilityIndex> {
Arc::clone(&self.lock(&self.index))
}
pub fn restore(&self, rows: impl IntoIterator<Item = StoredObservation>) -> usize {
let _deriving = self.lock(&self.deriving);
let mut held = self.lock(&self.index);
let mut builder = AvailabilityIndexBuilder::from_index(&held);
for (key, record) in store::restored_records(rows) {
builder = builder.record(key, record);
}
let refused = builder.superseded() + builder.misfiled();
*held = Arc::new(builder.build());
refused
}
pub fn persistable(&self) -> EvidenceWrite {
let orphaned = self
.lock(&self.orphaned)
.iter()
.map(|(key, before)| EvidenceClear::new(key.clone(), *before))
.collect::<Vec<_>>();
EvidenceWrite::of_index(&self.index()).clearing(orphaned)
}
pub fn acknowledge_persisted(&self, write: &EvidenceWrite) {
let mut orphaned = self.lock(&self.orphaned);
for clear in write.cleared() {
if orphaned
.get(&clear.key)
.is_some_and(|before| *before <= clear.before)
{
orphaned.remove(&clear.key);
}
}
}
pub fn derive(
&self,
state: &DesiredState,
readiness: &CredentialReadiness,
) -> Result<ProjectedAvailability, AvailabilityProjectionError> {
let _deriving = self.lock(&self.deriving);
let catalogue = Arc::clone(&self.lock(&self.catalogue));
let previous = self.index();
let pending: Vec<DiscoveryObservation> = self.lock(&self.pending).drain(..).collect();
let projected = match AvailabilityProjection::new(&catalogue, readiness).project(
state,
&previous,
pending.clone(),
) {
Ok(projected) => projected,
Err(error) => {
let mut queued = self.lock(&self.pending);
let since: Vec<DiscoveryObservation> = queued.drain(..).collect();
queued.extend(pending);
queued.extend(since);
return Err(error);
}
};
let derivation = {
let mut derivations = self.lock(&self.derivations);
*derivations += 1;
*derivations
};
let orphaned = projected
.orphaned()
.iter()
.filter_map(|key| {
previous
.record(key)
.and_then(latest_evidence_at)
.map(|before| (key.clone(), before))
})
.collect::<BTreeMap<_, _>>();
let mut pending_orphaned = self.lock(&self.orphaned);
let previous_orphaned = pending_orphaned.clone();
pending_orphaned.retain(|key, _| {
projected
.index()
.record(key)
.is_none_or(|record| !record.holds_evidence())
});
pending_orphaned.extend(orphaned);
drop(pending_orphaned);
*self.lock(&self.replaced) = Some(Superseded {
derivation,
index: previous,
orphaned: previous_orphaned,
derived_from: self.lock(&self.derived_from).clone(),
looks: pending,
});
*self.lock(&self.index) = Arc::new(projected.index().clone());
*self.lock(&self.derived_from) = Some((Arc::new(state.clone()), readiness.clone()));
Ok(ProjectedAvailability {
derivation,
..projected
})
}
pub fn abandon(&self, derivation: u64) -> bool {
let _deriving = self.lock(&self.deriving);
let mut held = self.lock(&self.replaced);
if held
.as_ref()
.is_none_or(|superseded| superseded.derivation != derivation)
{
return false;
}
let Some(replaced) = held.take() else {
return false;
};
drop(held);
*self.lock(&self.index) = replaced.index;
*self.lock(&self.orphaned) = replaced.orphaned;
*self.lock(&self.derived_from) = replaced.derived_from;
let mut queued = self.lock(&self.pending);
let since: Vec<DiscoveryObservation> = queued.drain(..).collect();
queued.extend(replaced.looks);
queued.extend(since);
true
}
pub fn reproject(&self) -> Option<Result<ProjectedAvailability, AvailabilityProjectionError>> {
let (state, readiness) = self.lock(&self.derived_from).clone()?;
Some(self.derive(&state, &readiness))
}
fn lock<'a, T>(&self, guarded: &'a Mutex<T>) -> MutexGuard<'a, T> {
guarded.lock().unwrap_or_else(PoisonError::into_inner)
}
}
pub trait AvailabilityReader: Send + Sync {
fn read(&self) -> Option<(Arc<AvailabilityIndex>, RuntimeObservations)>;
}
pub struct AvailabilityView<'a> {
index: &'a AvailabilityIndex,
runtime: &'a RuntimeObservations,
}
impl<'a> AvailabilityView<'a> {
pub const fn new(index: &'a AvailabilityIndex, runtime: &'a RuntimeObservations) -> Self {
Self { index, runtime }
}
pub fn evaluate(&self, key: &AvailabilityKey, now: SystemTime) -> Availability {
self.index
.evaluate_with(key, now, self.runtime.health(&key.target))
}
pub fn evaluate_effective(
&self,
scope: ScopeRef,
target: &TargetRef,
now: SystemTime,
) -> Availability {
let own = AvailabilityKey::new(scope, target.clone());
if scope.is_tenant_wide() || self.index.record(&own).is_some() {
return self.evaluate(&own, now);
}
self.evaluate(
&AvailabilityKey::new(ScopeRef::tenant(scope.tenant), target.clone()),
now,
)
}
pub fn evaluate_inherited_scope(
&self,
scope: ScopeRef,
now: SystemTime,
) -> Vec<(TargetRef, Availability)> {
if scope.is_tenant_wide() {
return self.evaluate_scope(scope, now);
}
let inherited = ScopeRef::tenant(scope.tenant);
let targets: BTreeSet<TargetRef> = self
.index
.evaluate_scope(&inherited, now)
.into_iter()
.chain(self.index.evaluate_scope(&scope, now))
.map(|(target, _)| target)
.collect();
targets
.into_iter()
.map(|target| {
let verdict = self.evaluate_effective(scope, &target, now);
(target, verdict)
})
.collect()
}
pub fn evaluate_scope(
&self,
scope: ScopeRef,
now: SystemTime,
) -> Vec<(TargetRef, Availability)> {
self.index
.evaluate_scope(&scope, now)
.into_iter()
.map(|(target, _)| {
let verdict = self.evaluate(&AvailabilityKey::new(scope, target.clone()), now);
(target, verdict)
})
.collect()
}
}
const fn rank(entitlement: Entitlement) -> u8 {
match entitlement {
Entitlement::Granted => 3,
Entitlement::Unknown => 2,
Entitlement::Revoked => 1,
Entitlement::Missing => 0,
}
}
const fn presence_rank(presence: CataloguePresence) -> u8 {
match presence {
CataloguePresence::Present => 2,
CataloguePresence::Withdrawn => 1,
CataloguePresence::Absent => 0,
}
}
const fn policy_rank(policy: PolicyDecision) -> u8 {
match policy {
PolicyDecision::Permitted => 2,
PolicyDecision::Indeterminate => 1,
PolicyDecision::Denied => 0,
}
}
fn least_permissive(held: &AvailabilityRecord, other: AvailabilityRecord) -> AvailabilityRecord {
let (entitlement, credential) = if rank(other.entitlement) < rank(held.entitlement) {
(other.entitlement, other.credential.clone())
} else {
(held.entitlement, held.credential.clone())
};
AvailabilityRecord {
presence: if presence_rank(other.presence) < presence_rank(held.presence) {
other.presence
} else {
held.presence
},
enablement: if held.enablement.is_enabled() && other.enablement.is_enabled() {
Enablement::Enabled
} else {
Enablement::NotEnabled
},
entitlement,
policy: if policy_rank(other.policy) < policy_rank(held.policy) {
other.policy
} else {
held.policy
},
credential,
..AvailabilityRecord::default()
}
}
fn note(named: &mut BTreeSet<AvailabilityKey>, key: AvailabilityKey) {
named.insert(key);
if named.len() > REPORTED_KEYS {
named.pop_last();
}
}
fn reported(named: BTreeSet<AvailabilityKey>) -> Vec<AvailabilityKey> {
named.into_iter().collect()
}
fn scope_of(owner: ModelOwner) -> ScopeRef {
ScopeRef {
tenant: owner.tenant,
project: owner.project,
}
}
fn owner_of_provider(provider: &crate::desired_state::providers::Provider) -> ModelOwner {
ModelOwner {
tenant: provider.body.tenant(),
project: provider.body.project(),
}
}
fn enablement_of(enablement: &ModelEnablement) -> Enablement {
if enablement.body.state().is_enabled() {
Enablement::Enabled
} else {
Enablement::NotEnabled
}
}
fn policy_of(policies: &PolicySet, owner: ModelOwner) -> PolicyDecision {
let scope = match owner.project {
None => PolicyScope::Tenant(owner.tenant),
Some(project) => PolicyScope::Project {
tenant: owner.tenant,
project,
},
};
match policies.effective(scope) {
None => PolicyDecision::Indeterminate,
Some(document) if document.body.budget().subject_limit_microdollars() == 0 => {
PolicyDecision::Denied
}
Some(_) => PolicyDecision::Permitted,
}
}
#[cfg(test)]
pub(crate) mod testing {
use super::{Catalogue, CatalogueListing};
use crate::backends::catalog::{
CatalogContent, CatalogModelEntry, CatalogProvider, JsonPointer, ModelFacts, ModelId,
ProviderEndpoint, ProviderId, ProviderOffering,
};
use crate::desired_state::Checksum;
pub(crate) fn listing(snapshot: Checksum, provider: &str, model: &str) -> CatalogueListing {
let id = ProviderId::parse(provider).expect("a well-formed provider id");
let content = CatalogContent::new(
vec![CatalogProvider {
id: id.clone(),
display_name: None,
doc_url: None,
endpoint: ProviderEndpoint::default(),
env_vars: Vec::new(),
pointer: JsonPointer::new("").child("providers").child(provider),
}],
vec![CatalogModelEntry {
id: ModelId::parse(model).expect("a well-formed model id"),
neutral: None,
offerings: vec![ProviderOffering {
provider: id,
model: ModelId::parse(model).expect("a well-formed model id"),
published_model_id: model.to_owned(),
facts: ModelFacts::default(),
overrides: Vec::new(),
price: None,
endpoint: ProviderEndpoint::default(),
pointer: JsonPointer::new("").child("models").child(model),
}],
}],
)
.expect("a catalogue with one offering");
CatalogueListing::of(snapshot, &content)
}
pub(crate) fn catalogue(snapshot: Checksum, provider: &str, model: &str) -> Catalogue {
Catalogue::active(listing(snapshot, provider, model))
}
}