use crate::{CapacityDimension, DiagnosticSubmission, EmergencyDiagnosticHandle, Observer};
use saddle_admission::{
AdmissionError, ProfuseGwAdmissionObservationReceipt, ProfuseGwCapacityBottleneck, ProfuseGwCapacityDecision,
ProfuseGwCapacityRejectReason, ProfuseGwCapacitySnapshot,
};
use saddle_core::{ContextFact, ContextLabel, RequestExecutionView, RequestViewPhase};
use serde::{Serialize, Serializer, ser::SerializeStruct};
pub struct AdmissionCapacityFacts {
snapshot: ProfuseGwCapacitySnapshot,
used: usize,
decision: ProfuseGwCapacityDecision,
reason: Option<ProfuseGwCapacityRejectReason>,
elapsed_ms: u64,
resource: Option<saddle_admission::ResourceDecisionSnapshot>,
}
impl AdmissionCapacityFacts {
pub fn from_receipt(receipt: &ProfuseGwAdmissionObservationReceipt) -> Self {
Self {
snapshot: receipt.snapshot(),
used: receipt.used(),
decision: receipt.decision(),
reason: receipt.reason(),
elapsed_ms: receipt.elapsed_ms(),
resource: receipt.resource(),
}
}
}
#[derive(Clone, Copy)]
pub enum AdmissionEventContext<'a> {
Rooted(&'a RequestExecutionView),
Unrooted {
application: &'a ContextLabel,
lifecycle: RequestViewPhase,
},
}
#[doc(hidden)]
#[derive(Clone, Copy, Serialize)]
pub struct AdmissionRecognizedIngress<'a> {
pub request: &'a str,
pub route: &'a str,
}
impl Serialize for AdmissionEventContext<'_> {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
match self {
Self::Rooted(view) => view.serialize(serializer),
Self::Unrooted {
application,
lifecycle,
} => {
let mut out = serializer.serialize_struct("UnrootedRequestContext", 20)?;
out.serialize_field("schema_version", &2u8)?;
out.serialize_field("application", &ContextFact::Present(application))?;
out.serialize_field("lifecycle", &ContextFact::Present(lifecycle))?;
for field in [
"local_request",
"publication",
"call_application",
"module",
"service",
"operation",
"trace_id",
"request",
"span_id",
"route",
"attempt",
"rpc_id",
"zone",
"db_operation",
"scope",
"task",
"target",
] {
out.serialize_field(field, &ContextFact::<()>::NotEstablished)?;
}
out.end()
}
}
}
}
#[derive(Clone, Copy, Debug, Serialize)]
#[serde(tag = "status", content = "failure", rename_all = "snake_case")]
pub enum AdmissionConstruction {
NotAttempted,
QueuedForRetry,
Ready,
Cancelled,
Failed(AdmissionConstructionFailure),
}
#[derive(Clone, Copy, Debug, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum AdmissionConstructionFailure {
Storage,
Deadline,
RequestIdentity,
Context,
RejectionTaskBytes,
RejectionTaskSlot,
}
#[must_use]
pub struct AdmissionCapacitySubmission {
pub cpu: DiagnosticSubmission,
pub memory: DiagnosticSubmission,
pub database: DiagnosticSubmission,
pub profuse_contract: DiagnosticSubmission,
}
#[derive(Serialize)]
struct CapacityRecord<'a> {
schema_version: u8,
timestamp_unix_ms: u128,
level: &'static str,
event: &'static str,
stage: &'static str,
context: AdmissionEventContext<'a>,
recognized_ingress: Option<AdmissionRecognizedIngress<'a>>,
capacity_dimension: &'static str,
budget: usize,
limit: usize,
used: usize,
bottleneck: &'static str,
decision: &'static str,
outcome: &'static str,
reject_reason: Option<&'static str>,
elapsed_ms: u64,
resource_source: &'static str,
resource: Option<saddle_admission::ResourceDecisionSnapshot>,
construction: AdmissionConstruction,
}
#[derive(Serialize)]
struct IngressResource<'a> {
domain: &'a str,
requested_bytes: usize,
available_bytes: usize,
}
#[derive(Serialize)]
struct IngressResourceDecisionRecord<'a> {
schema_version: u8,
timestamp_unix_ms: u128,
timestamp_role: &'static str,
decision_monotonic_ms: Option<u64>,
decision_time_valid: bool,
decision_time_source: &'static str,
resource_valid: bool,
level: &'static str,
event: &'static str,
stage: &'static str,
context: AdmissionEventContext<'a>,
decision: &'static str,
resource_source: &'static str,
resource: IngressResource<'a>,
}
impl EmergencyDiagnosticHandle {
pub fn record_ingress_resource_failure(
&self,
application: &ContextLabel,
error: &AdmissionError,
) -> Option<DiagnosticSubmission> {
let (domain, requested_bytes, available_bytes) = match error {
AdmissionError::FrameworkReserveExceeded { requested, available } =>
("framework", *requested, *available),
AdmissionError::ProcessCapacityExceeded { requested, available } =>
("managed", *requested, *available),
_ => return None,
};
let timestamp_unix_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_millis();
Some(self.submit_fixed_record(&IngressResourceDecisionRecord {
schema_version: 1,
timestamp_unix_ms,
timestamp_role: "output_enqueued",
decision_monotonic_ms: None,
decision_time_valid: false,
decision_time_source: "not_captured_by_ledger_transaction",
resource_valid: true,
level: "warn",
event: "framework.ingress_resource_decision",
stage: "admission",
context: AdmissionEventContext::Unrooted {
application,
lifecycle: RequestViewPhase::Admitted,
},
decision: "capacity_rejected",
resource_source: "admission_ledger_same_lock",
resource: IngressResource { domain, requested_bytes, available_bytes },
}))
}
}
impl Observer {
pub fn record_admission_capacity(
&self,
output: Option<&EmergencyDiagnosticHandle>,
context: AdmissionEventContext<'_>,
facts: &AdmissionCapacityFacts,
construction: AdmissionConstruction,
) -> AdmissionCapacitySubmission {
self.record_admission_capacity_with_recognized(output, context, facts, construction, None)
}
#[doc(hidden)]
pub fn record_admission_capacity_with_recognized<'a>(
&self,
output: Option<&EmergencyDiagnosticHandle>,
context: AdmissionEventContext<'a>,
facts: &AdmissionCapacityFacts,
construction: AdmissionConstruction,
recognized_ingress: Option<AdmissionRecognizedIngress<'a>>,
) -> AdmissionCapacitySubmission {
let snapshot = facts.snapshot;
let dimensions = [
(CapacityDimension::Cpu, "cpu", snapshot.cpu_budget()),
(
CapacityDimension::Memory,
"memory",
snapshot.memory_budget_bytes(),
),
(
CapacityDimension::Database,
"database",
snapshot.database_budget(),
),
(
CapacityDimension::ProfuseContract,
"profuse_contract",
snapshot.profusecontract_budget(),
),
];
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis();
let rejected = facts.decision == ProfuseGwCapacityDecision::CapacityRejected;
let submissions = dimensions.map(|(dimension, name, budget)| {
self.inner.metrics.capacity(
Some(dimension),
u64::try_from(facts.used).unwrap_or(u64::MAX),
rejected,
);
let Some(output) = output else {
return DiagnosticSubmission::OutputUnavailable;
};
output.submit_fixed_record(&CapacityRecord {
schema_version: 1,
timestamp_unix_ms: timestamp,
level: "info",
event: "framework.capacity",
stage: "admission",
context,
recognized_ingress,
capacity_dimension: name,
budget,
limit: snapshot.active_limit(),
used: facts.used,
bottleneck: match snapshot.bottleneck() {
ProfuseGwCapacityBottleneck::Cpu => "cpu",
ProfuseGwCapacityBottleneck::Memory => "memory",
ProfuseGwCapacityBottleneck::Database => "database",
ProfuseGwCapacityBottleneck::ProfuseContract => "profuse_contract",
},
decision: if rejected {
"capacity_rejected"
} else {
"accepted"
},
outcome: if rejected { "rejected" } else { "accepted" },
reject_reason: facts.reason.map(|reason| match reason {
ProfuseGwCapacityRejectReason::AtLimit => "at_limit",
}),
elapsed_ms: facts.elapsed_ms,
resource_source: if facts.resource.is_some() {
"admission_ledger_same_lock"
} else {
"not_observed_before_resource_check"
},
resource: facts.resource,
construction,
})
});
let [cpu, memory, database, profuse_contract] = submissions;
AdmissionCapacitySubmission {
cpu,
memory,
database,
profuse_contract,
}
}
}
pub fn admission_capacity_layouts() -> [std::alloc::Layout; 4] {
[
std::alloc::Layout::new::<AdmissionCapacityFacts>(),
std::alloc::Layout::new::<AdmissionEventContext<'static>>(),
std::alloc::Layout::new::<CapacityRecord<'static>>(),
std::alloc::Layout::new::<AdmissionCapacitySubmission>(),
]
}