use super::ProtectionLayer;
use super::RUNTIME;
#[derive(Clone, Debug, Eq, PartialEq, serde::Serialize)]
pub(crate) enum ProductionTraceObservation {
PersistOneEntry { node_did: String },
YieldActor { node_did: String },
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, serde::Serialize)]
pub(crate) struct DeadlineMissWitness {
pub(crate) observed_virtual_ms: u64,
pub(crate) deadline_virtual_ms: u64,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, serde::Serialize)]
struct CapacityMetric {
peak: usize,
limit: usize,
}
impl CapacityMetric {
fn observe(&mut self, current: usize, limit: usize) {
self.peak = self.peak.max(current);
self.limit = limit;
}
}
#[derive(Clone, Debug, Default, Eq, PartialEq, serde::Serialize)]
pub(crate) struct ProductionCapacityObservations {
outbound_peer_transfers: CapacityMetric,
outbound_peer_bytes: CapacityMetric,
outbound_global_bytes: CapacityMetric,
inbound_peer_transfers: CapacityMetric,
inbound_peer_bytes: CapacityMetric,
inbound_node_transfers: CapacityMetric,
inbound_node_bytes: CapacityMetric,
reassembly_node_bytes: CapacityMetric,
reassembly_peer_bytes: CapacityMetric,
reassembly_pending_messages: CapacityMetric,
}
impl ProductionCapacityObservations {
pub(crate) fn validate(&self) -> Result<(), String> {
let metrics = [
("outbound_peer_transfers", self.outbound_peer_transfers),
("outbound_peer_bytes", self.outbound_peer_bytes),
("outbound_global_bytes", self.outbound_global_bytes),
("inbound_peer_transfers", self.inbound_peer_transfers),
("inbound_peer_bytes", self.inbound_peer_bytes),
("inbound_node_transfers", self.inbound_node_transfers),
("inbound_node_bytes", self.inbound_node_bytes),
("reassembly_node_bytes", self.reassembly_node_bytes),
("reassembly_peer_bytes", self.reassembly_peer_bytes),
(
"reassembly_pending_messages",
self.reassembly_pending_messages,
),
];
for (name, metric) in metrics {
if metric.peak == 0 || metric.limit == 0 {
return Err(format!(
"production capacity metric {name} was not observed"
));
}
if metric.peak > metric.limit {
return Err(format!(
"production capacity metric {name} exceeded its limit: {} > {}",
metric.peak, metric.limit
));
}
}
Ok(())
}
}
pub(crate) fn record_protection_violation(layer: ProtectionLayer) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime.observations.record(layer);
}
});
}
pub(crate) fn record_barrier_control_blocked() {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime.observations.record_barrier_blocked();
}
});
}
pub(crate) fn record_barrier_control_deadline_miss(
observed_virtual_ms: u64,
deadline_virtual_ms: u64,
) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime
.observations
.record_barrier_deadline_miss(DeadlineMissWitness {
observed_virtual_ms,
deadline_virtual_ms,
});
}
});
}
pub(crate) fn record_repair_entries(entries: usize) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime.repair_entries_observed = runtime
.repair_entries_observed
.saturating_add(entries as u64);
}
});
}
pub(crate) fn record_reassembly_advance(transaction_id: uuid::Uuid) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
let advances = runtime
.reassembly_advances
.entry(transaction_id)
.or_default();
*advances = advances.saturating_add(1);
}
});
}
pub(crate) fn record_storage_persisted(node: crate::dht::Did) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime.production_trace_observations.push(
ProductionTraceObservation::PersistOneEntry {
node_did: node.to_string(),
},
);
}
});
}
pub(crate) fn record_storage_actor_yield(node: crate::dht::Did) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime
.production_trace_observations
.push(ProductionTraceObservation::YieldActor {
node_did: node.to_string(),
});
}
});
}
pub(crate) fn signal_storage_progress_probe() -> Option<u64> {
RUNTIME.with(|runtime| {
let runtime = runtime.borrow();
let runtime = runtime.as_ref()?;
runtime.storage_progress_notify.as_ref()?.notify_one();
Some(runtime.storage_progress_epoch)
})
}
pub(crate) fn record_storage_progress() {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime.storage_progress_epoch = runtime.storage_progress_epoch.saturating_add(1);
}
});
}
pub(crate) fn record_storage_progress_between_entries() {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime
.observations
.record_storage_progress_between_entries();
}
});
}
pub(crate) fn storage_progress_epoch() -> Option<u64> {
RUNTIME.with(|runtime| {
runtime
.borrow()
.as_ref()
.map(|runtime| runtime.storage_progress_epoch)
})
}
pub(crate) fn record_outbound_submission(transaction_id: uuid::Uuid) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
let elapsed_ms = u64::try_from(runtime.elapsed_ms).unwrap_or(u64::MAX);
runtime
.outbound_submission_ms
.entry(transaction_id)
.or_insert(elapsed_ms);
}
});
}
pub(crate) fn observe_outbound_peer_capacity(
transfers: usize,
transfer_limit: usize,
bytes: usize,
byte_limit: usize,
) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime
.capacity_observations
.outbound_peer_transfers
.observe(transfers, transfer_limit);
runtime
.capacity_observations
.outbound_peer_bytes
.observe(bytes, byte_limit);
}
});
}
pub(crate) fn observe_outbound_global_capacity(bytes: usize, limit: usize) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
runtime
.capacity_observations
.outbound_global_bytes
.observe(bytes, limit);
}
});
}
pub(crate) fn observe_inbound_capacity(
peer: (usize, usize, usize, usize),
node: (usize, usize, usize, usize),
) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
let observations = &mut runtime.capacity_observations;
observations.inbound_peer_transfers.observe(peer.0, peer.2);
observations.inbound_peer_bytes.observe(peer.1, peer.3);
observations.inbound_node_transfers.observe(node.0, node.2);
observations.inbound_node_bytes.observe(node.1, node.3);
}
});
}
pub(crate) fn observe_reassembly_capacity(
node_bytes: usize,
node_limit: usize,
peer_bytes: usize,
peer_limit: usize,
pending_messages: usize,
pending_limit: usize,
) {
RUNTIME.with(|runtime| {
if let Some(runtime) = runtime.borrow_mut().as_mut() {
let observations = &mut runtime.capacity_observations;
observations
.reassembly_node_bytes
.observe(node_bytes, node_limit);
observations
.reassembly_peer_bytes
.observe(peer_bytes, peer_limit);
observations
.reassembly_pending_messages
.observe(pending_messages, pending_limit);
}
});
}