use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use async_trait::async_trait;
use axum::Router;
use axum::body::Body;
use axum::http::{HeaderMap, Request, StatusCode};
use axum::routing::post;
use serde_json::json;
use tokio::sync::{Semaphore, mpsc};
use super::sweep::LeakSweep;
use crate::backends::control_plane::ControlPlaneStore;
use crate::backends::fakes::InMemorySecrets;
use crate::backends::secrets::{SecretMaterial, SecretResolver as _, SecretStore};
use crate::budget::NoBudget;
use crate::config::{Config, Model, Target};
use crate::convergence::compile::{CandidateCompiler, CompileError, RevisionProjection};
use crate::convergence::credentials::RuntimeProjection;
use crate::convergence::secrets::{MaterialLedger, SecretMaterialization};
use crate::convergence::status::testing::ManualClock;
use crate::convergence::{BackoffPolicy, ConvergenceSettings, Outcome, Reconciler};
use crate::desired_state::credentials::ProviderCredentialBody;
use crate::desired_state::models::WireFamily;
use crate::desired_state::oracle::InMemoryControlPlane;
use crate::desired_state::providers::ProviderBody;
use crate::desired_state::{
CanonicalValue, DesiredState, ExpectedRevision, LoadedRevision, ResourceBody, ResourceKind,
ResourceRef, ResourceScope, ResourceVersion, ResourceVersionNumber, RevisionId,
SecretLifecycle, SecretOwner, SecretRef, Slug, fixtures,
};
use crate::state::{AppState, ConfigSnapshot};
use crate::telemetry;
use crate::usage::{UsageFanout, UsageRecord, UsageSink};
pub(crate) const PROVIDER_MATERIAL: &str = "sk-axond-sentinel-provider-6f21a9d0c7b4";
pub(crate) const ROTATED_MATERIAL: &str = "sk-axond-sentinel-rotated-b48c37e1590a";
pub(crate) const INBOUND_MATERIAL: &str = "axond-sentinel-inbound-2c8a11f5e304";
pub(crate) fn sweep() -> LeakSweep {
LeakSweep::of([
("provider", PROVIDER_MATERIAL),
("rotated", ROTATED_MATERIAL),
("inbound", INBOUND_MATERIAL),
])
}
pub(crate) fn bootstrap(base_url: &str) -> Config {
Config::from_toml_str(&format!(
r#"
[[namespace]]
id = "platform"
default = true
[[provider]]
id = "openai"
kind = "openai"
base_url = "{base_url}"
[[gateway_key]]
env = "AXOND_SENTINEL_INBOUND"
namespace = "platform"
"#
))
.expect("a valid bootstrap config")
}
pub(crate) fn bootstrap_env() -> HashMap<String, String> {
HashMap::from([(
"AXOND_SENTINEL_INBOUND".to_owned(),
INBOUND_MATERIAL.to_owned(),
)])
}
pub(crate) const SERVING_NAMESPACE: &str = "acme/core";
pub(crate) struct SecretResolvingCompiler {
bootstrap: Config,
env: HashMap<String, String>,
materialization: Arc<SecretMaterialization>,
provider: &'static str,
resolutions: AtomicUsize,
}
impl SecretResolvingCompiler {
pub(crate) fn new(bootstrap: Config, secrets: Arc<dyn SecretStore>) -> Self {
Self {
bootstrap,
env: bootstrap_env(),
materialization: Arc::new(SecretMaterialization::new(secrets, MaterialLedger::new())),
provider: "openai",
resolutions: AtomicUsize::new(0),
}
}
pub(crate) fn ledger(&self) -> &Arc<MaterialLedger> {
self.materialization.ledger()
}
pub(crate) fn resolutions(&self) -> usize {
self.resolutions.load(Ordering::Relaxed)
}
}
#[async_trait]
impl CandidateCompiler for SecretResolvingCompiler {
async fn compile(
&self,
revision: &LoadedRevision,
generation: u64,
) -> Result<ConfigSnapshot, CompileError> {
let id = revision.id();
let projection = |source| CompileError::Projection {
revision: id,
source,
};
let mut config = RuntimeProjection
.project(&self.bootstrap, revision.state(), id)
.map_err(projection)?;
let resolved = self
.materialization
.resolve(revision.state())
.await
.map_err(projection)?;
self.resolutions
.fetch_add(resolved.len(), Ordering::Relaxed);
assert!(
config
.namespace
.iter()
.any(|namespace| namespace.id == SERVING_NAMESPACE),
"the fixture project must project as {SERVING_NAMESPACE}"
);
for key in &mut config.gateway_key {
key.namespace = SERVING_NAMESPACE.to_owned();
}
for resource in revision.state().resources() {
if resource.reference.kind != ResourceKind::Alias {
continue;
}
config.model.push(Model {
name: resource.slug.as_str().to_owned(),
namespace: None,
targets: vec![Target {
provider: self.provider.to_owned(),
model: "gpt-4o".to_owned(),
price: gateway_core::catalog::ModelPrice {
input_microdollars_per_million: 1_000_000,
output_microdollars_per_million: 2_000_000,
reasoning_microdollars_per_million: None,
cache_read_microdollars_per_million: None,
cache_write_microdollars_per_million: None,
},
}],
});
}
config
.validate_compiled()
.map_err(|source| CompileError::Validation {
revision: id,
source,
})?;
ConfigSnapshot::build_with(config, &self.env, generation, resolved).map_err(|source| {
CompileError::Snapshot {
revision: id,
source,
}
})
}
}
pub(crate) fn material(plaintext: &str) -> SecretMaterial {
SecretMaterial::new(plaintext.to_owned())
}
#[derive(Clone, Default)]
pub(crate) struct CapturingSink(Arc<Mutex<Vec<UsageRecord>>>);
impl CapturingSink {
pub(crate) fn records(&self) -> Vec<UsageRecord> {
self.0.lock().expect("not poisoned").clone()
}
}
#[async_trait]
impl UsageSink for CapturingSink {
fn name(&self) -> &'static str {
"capture"
}
async fn record(&self, record: &UsageRecord) {
self.0.lock().expect("not poisoned").push(record.clone());
}
}
pub(crate) struct FakeProvider {
pub(crate) base_url: String,
presented: Arc<Mutex<Vec<String>>>,
release: Option<Arc<Semaphore>>,
arrivals: tokio::sync::Mutex<mpsc::UnboundedReceiver<()>>,
}
impl FakeProvider {
pub(crate) async fn serving() -> Self {
Self::spawn(None).await
}
pub(crate) async fn gated() -> Self {
Self::spawn(Some(Arc::new(Semaphore::new(0)))).await
}
pub(crate) fn unreachable() -> Self {
let (_, arrivals) = mpsc::unbounded_channel();
Self {
base_url: "http://127.0.0.1:1".to_owned(),
presented: Arc::new(Mutex::new(Vec::new())),
release: None,
arrivals: tokio::sync::Mutex::new(arrivals),
}
}
async fn spawn(release: Option<Arc<Semaphore>>) -> Self {
let presented: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
let (arrived, arrivals) = mpsc::unbounded_channel();
let handler = {
let presented = Arc::clone(&presented);
let release = release.clone();
move |headers: HeaderMap, _: axum::body::Bytes| {
let presented = Arc::clone(&presented);
let release = release.clone();
let arrived = arrived.clone();
async move {
let credential = headers
.get(axum::http::header::AUTHORIZATION)
.and_then(|value| value.to_str().ok())
.unwrap_or_default()
.to_owned();
presented.lock().expect("not poisoned").push(credential);
let _ = arrived.send(());
if let Some(release) = release {
release
.acquire()
.await
.expect("the gate outlives the request")
.forget();
}
(
StatusCode::OK,
axum::Json(json!({
"id": "chatcmpl-sentinel",
"choices": [],
"usage": { "prompt_tokens": 7, "completion_tokens": 3 }
})),
)
}
}
};
let app = Router::new()
.route("/chat/completions", post(handler.clone()))
.route("/responses", post(handler));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("a loopback port");
let addr = listener.local_addr().expect("a bound address");
tokio::spawn(async move {
let _ = axum::serve(listener, app).await;
});
Self {
base_url: format!("http://{addr}"),
presented,
release,
arrivals: tokio::sync::Mutex::new(arrivals),
}
}
pub(crate) async fn await_arrival(&self) {
let mut arrivals = self.arrivals.lock().await;
tokio::time::timeout(std::time::Duration::from_secs(5), arrivals.recv())
.await
.expect("a request reaches the provider")
.expect("the provider outlives the test");
}
pub(crate) fn release(&self, count: usize) {
self.release
.as_ref()
.expect("a gated provider")
.add_permits(count);
}
pub(crate) fn presented(&self) -> Vec<String> {
self.presented.lock().expect("not poisoned").clone()
}
}
pub(crate) fn owner() -> SecretOwner {
SecretOwner::tenant(fixtures::tenant_id(1))
}
pub(crate) fn first() -> SecretRef {
fixtures::secret_ref(3)
}
pub(crate) fn credential(secret: SecretRef, version: ResourceVersionNumber) -> ResourceVersion {
ProviderCredentialBody::staged(
fixtures::resource_id(3),
owner(),
fixtures::provider_id(3),
fixtures::display_name("Primary"),
secret,
)
.transitioned(SecretLifecycle::Active)
.expect("staged material may be activated")
.version_at(Slug::parse("primary").expect("fixture slug"), version)
}
pub(crate) fn provider_connection() -> ResourceVersion {
ProviderBody::for_tenant(
fixtures::provider_id(3),
fixtures::tenant_id(1),
fixtures::display_name("OpenAI"),
WireFamily::OpenaiChat,
"https://api.openai.com/v1",
)
.version(Slug::parse("openai").expect("fixture slug"))
}
pub(crate) fn state_pinning(secret: SecretRef, version: ResourceVersionNumber) -> DesiredState {
let credential = credential(secret, version);
let alias = ResourceVersion::new(
ResourceRef::new(ResourceKind::Alias, fixtures::resource_id(4), version),
ResourceScope::Tenant(fixtures::tenant_id(1)),
Slug::parse("fast").expect("fixture slug"),
ResourceBody::Inline(CanonicalValue::map([(
"wire_family",
CanonicalValue::string("openai-chat"),
)])),
)
.depending_on([credential.reference]);
let mut state = DesiredState::new();
state
.insert(fixtures::tenant(1, "acme"))
.and_then(|state| state.insert(fixtures::project(&fixtures::tenant_id(1), 2, "core")))
.and_then(|state| state.insert(provider_connection()))
.and_then(|state| state.insert(credential))
.and_then(|state| state.insert(alias))
.expect("a valid revision");
state
}
pub(crate) fn state_sharing(secret: SecretRef, version: ResourceVersionNumber) -> DesiredState {
let primary = credential(secret, version);
let secondary = ProviderCredentialBody::staged(
fixtures::resource_id(5),
owner(),
fixtures::provider_id(3),
fixtures::display_name("Secondary"),
secret,
)
.transitioned(SecretLifecycle::Active)
.expect("staged material may be activated")
.version_at(Slug::parse("secondary").expect("fixture slug"), version);
let alias = ResourceVersion::new(
ResourceRef::new(ResourceKind::Alias, fixtures::resource_id(4), version),
ResourceScope::Tenant(fixtures::tenant_id(1)),
Slug::parse("fast").expect("fixture slug"),
ResourceBody::Inline(CanonicalValue::map([(
"wire_family",
CanonicalValue::string("openai-chat"),
)])),
)
.depending_on([primary.reference, secondary.reference]);
let mut state = DesiredState::new();
state
.insert(fixtures::tenant(1, "acme"))
.and_then(|state| state.insert(fixtures::project(&fixtures::tenant_id(1), 2, "core")))
.and_then(|state| state.insert(provider_connection()))
.and_then(|state| state.insert(primary))
.and_then(|state| state.insert(secondary))
.and_then(|state| state.insert(alias))
.expect("a valid revision");
state
}
pub(crate) struct Replica<S = InMemorySecrets> {
pub(crate) store: Arc<InMemoryControlPlane>,
pub(crate) secrets: Arc<S>,
pub(crate) compiler: Arc<SecretResolvingCompiler>,
pub(crate) state: AppState,
pub(crate) reconciler: Arc<Reconciler>,
}
impl Replica<InMemorySecrets> {
pub(crate) fn new(provider: &FakeProvider) -> Self {
Self::with_sinks(provider, Vec::new())
}
pub(crate) fn with_sinks(provider: &FakeProvider, sinks: Vec<Box<dyn UsageSink>>) -> Self {
Self::backed_by(provider, Arc::new(InMemorySecrets::new()), sinks)
}
}
impl<S: SecretStore + 'static> Replica<S> {
pub(crate) fn backed_by(
provider: &FakeProvider,
secrets: Arc<S>,
sinks: Vec<Box<dyn UsageSink>>,
) -> Self {
let store = Arc::new(InMemoryControlPlane::new());
let compiler = Arc::new(SecretResolvingCompiler::new(
bootstrap(&provider.base_url),
Arc::clone(&secrets) as Arc<dyn SecretStore>,
));
let state = AppState::new(
bootstrap(&provider.base_url),
&bootstrap_env(),
UsageFanout::new(sinks),
Box::new(NoBudget),
)
.expect("the bootstrap config is servable");
let reconciler = Arc::new(Reconciler::new(
Arc::clone(&store) as Arc<dyn ControlPlaneStore>,
Arc::clone(&compiler) as Arc<dyn CandidateCompiler>,
Arc::new(state.clone()),
settings(),
None,
Arc::new(ManualClock::new()),
));
Self {
store,
secrets,
compiler,
state,
reconciler,
}
}
pub(crate) async fn publish(&self, key: &str, state: DesiredState) -> RevisionId {
let expected = self
.store
.desired_revision()
.await
.expect("the control plane is readable")
.map_or(ExpectedRevision::Empty, ExpectedRevision::Exactly);
self.store
.publish_revision(fixtures::candidate(expected, key, state))
.await
.expect("the candidate is valid")
.id
}
pub(crate) async fn converge(&self) -> Outcome {
self.reconciler
.converge_once(telemetry::CONVERGENCE_POLLED)
.await
}
pub(crate) fn generation(&self) -> u64 {
self.state.config().generation
}
}
fn settings() -> ConvergenceSettings {
ConvergenceSettings {
poll_interval: Duration::from_millis(100),
target: Duration::from_secs(1),
backoff: BackoffPolicy {
initial: Duration::from_millis(100),
max: Duration::from_secs(4),
multiplier: 2,
},
}
}
pub(crate) fn chat_request() -> Request<Body> {
Request::post("/v1/chat/completions")
.header("content-type", "application/json")
.header(
axum::http::header::AUTHORIZATION,
format!("Bearer {INBOUND_MATERIAL}"),
)
.body(Body::from(r#"{"model":"fast","messages":[]}"#))
.expect("a valid request")
}
pub(crate) async fn live_material(pairs: &[(SecretRef, &'static str)]) -> Vec<String> {
let secrets = InMemorySecrets::new();
let mut resolved = Vec::with_capacity(pairs.len());
for (reference, plaintext) in pairs {
secrets.seed(owner(), *reference, plaintext, SecretLifecycle::Active);
resolved.push(
secrets
.resolve(owner(), reference)
.await
.expect("active material resolves")
.expose()
.to_owned(),
);
}
resolved
}