use std::collections::BTreeMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, Weak};
use crate::CounterSnapshot;
#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
pub enum Phase {
Started,
Slow,
Heartbeat,
Finished,
Failed,
}
#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
pub enum EventSource {
Engine,
SqliteInternal,
}
#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
pub enum EventCategory {
Writer,
Search,
Admin,
Error,
Corruption,
Recovery,
Io,
}
#[derive(Debug, Clone)]
pub struct Event {
pub phase: Phase,
pub source: EventSource,
pub category: EventCategory,
pub code: Option<&'static str>,
}
pub trait Subscriber: Send + Sync {
fn on_event(&self, event: &Event);
fn on_profile(&self, _record: &ProfileRecord) {}
fn on_slow_statement(&self, _signal: &SlowStatement) {}
fn on_stress_failure(&self, _context: &StressFailureContext) {}
}
#[derive(Debug, Clone)]
pub struct SlowStatement {
pub statement: String,
pub wall_clock_ms: u64,
}
#[derive(Default)]
pub(crate) struct SubscriberRegistry {
next_id: AtomicU64,
entries: Mutex<Vec<(u64, Arc<dyn Subscriber>)>>,
}
impl std::fmt::Debug for SubscriberRegistry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let count = self.entries.lock().map(|e| e.len()).unwrap_or(0);
f.debug_struct("SubscriberRegistry").field("subscribers", &count).finish()
}
}
impl SubscriberRegistry {
pub(crate) fn new() -> Self {
Self::default()
}
pub(crate) fn attach(self: &Arc<Self>, subscriber: Arc<dyn Subscriber>) -> Subscription {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
if let Ok(mut entries) = self.entries.lock() {
entries.push((id, subscriber));
}
Subscription { id, registry: Arc::downgrade(self) }
}
pub(crate) fn attach_persistent(&self, subscriber: Arc<dyn Subscriber>) {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
if let Ok(mut entries) = self.entries.lock() {
entries.push((id, subscriber));
}
}
fn detach(&self, id: u64) {
if let Ok(mut entries) = self.entries.lock() {
entries.retain(|(eid, _)| *eid != id);
}
}
pub(crate) fn dispatch(&self, event: &Event) {
for sub in self.snapshot() {
sub.on_event(event);
}
}
pub(crate) fn dispatch_profile(&self, record: &ProfileRecord) {
for sub in self.snapshot() {
sub.on_profile(record);
}
}
pub(crate) fn dispatch_slow_statement(&self, signal: &SlowStatement) {
for sub in self.snapshot() {
sub.on_slow_statement(signal);
}
}
pub(crate) fn dispatch_stress_failure(&self, context: &StressFailureContext) {
for sub in self.snapshot() {
sub.on_stress_failure(context);
}
}
fn snapshot(&self) -> Vec<Arc<dyn Subscriber>> {
match self.entries.lock() {
Ok(entries) => entries.iter().map(|(_, s)| Arc::clone(s)).collect(),
Err(_) => Vec::new(),
}
}
}
#[derive(Debug)]
pub struct Subscription {
id: u64,
registry: Weak<SubscriberRegistry>,
}
impl Drop for Subscription {
fn drop(&mut self) {
if let Some(registry) = self.registry.upgrade() {
registry.detach(self.id);
}
}
}
#[derive(Debug, Clone, Copy, Eq, PartialEq)]
pub struct ProfileRecord {
pub wall_clock_ms: u64,
pub step_count: u64,
pub cache_delta: i64,
}
#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
pub enum ProjectionStatus {
Pending,
Failed,
UpToDate,
}
#[derive(Debug, Clone)]
pub struct StressFailureContext {
pub thread_group_id: u64,
pub op_kind: String,
pub last_error_chain: Vec<String>,
pub projection_state: String,
}
#[derive(Debug)]
pub(crate) struct Counters {
queries: AtomicU64,
writes: AtomicU64,
write_rows: AtomicU64,
admin_ops: AtomicU64,
cache_hit: AtomicU64,
cache_miss: AtomicU64,
errors_by_code: Mutex<BTreeMap<String, u64>>,
}
impl Counters {
pub(crate) fn new() -> Self {
Self {
queries: AtomicU64::new(0),
writes: AtomicU64::new(0),
write_rows: AtomicU64::new(0),
admin_ops: AtomicU64::new(0),
cache_hit: AtomicU64::new(0),
cache_miss: AtomicU64::new(0),
errors_by_code: Mutex::new(BTreeMap::new()),
}
}
pub(crate) fn record_write(&self, rows: u64) {
self.writes.fetch_add(1, Ordering::Relaxed);
self.write_rows.fetch_add(rows, Ordering::Relaxed);
}
pub(crate) fn record_query(&self) {
self.queries.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_admin(&self) {
self.admin_ops.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_error(&self, code: &str) {
if let Ok(mut map) = self.errors_by_code.lock() {
*map.entry(code.to_string()).or_insert(0) += 1;
}
}
#[allow(dead_code)]
pub(crate) fn record_cache_hit(&self) {
self.cache_hit.fetch_add(1, Ordering::Relaxed);
}
#[allow(dead_code)]
pub(crate) fn record_cache_miss(&self) {
self.cache_miss.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn snapshot(&self) -> CounterSnapshot {
let errors_by_code = self.errors_by_code.lock().map(|map| map.clone()).unwrap_or_default();
CounterSnapshot {
queries: self.queries.load(Ordering::Relaxed),
writes: self.writes.load(Ordering::Relaxed),
write_rows: self.write_rows.load(Ordering::Relaxed),
errors_by_code,
admin_ops: self.admin_ops.load(Ordering::Relaxed),
cache_hit: self.cache_hit.load(Ordering::Relaxed),
cache_miss: self.cache_miss.load(Ordering::Relaxed),
}
}
}