use core_storage::Value;
use serde::Serialize;
use std::collections::VecDeque;
use std::sync::{Arc, Condvar, Mutex, Weak};
use std::time::Duration;
pub const DEFAULT_SUB_CAPACITY: usize = 65_536;
#[derive(Debug, Clone, Serialize, PartialEq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum DbEvent {
EdgeFired {
rule: String,
src_key: String,
dst_key: String,
edge_type: String,
#[serde(skip_serializing_if = "Option::is_none")]
weight: Option<f64>,
commit_seq: u64,
},
EdgeRetracted {
rule: String,
src_key: String,
dst_key: String,
edge_type: String,
commit_seq: u64,
},
NodeInserted {
label: String,
key: String,
commit_seq: u64,
},
NodeDeleted { key: String, commit_seq: u64 },
EdgeInserted {
edge_type: String,
src: String,
dst: String,
commit_seq: u64,
},
EdgeDeleted {
edge_type: String,
src: String,
dst: String,
commit_seq: u64,
},
PropSet {
key: String,
field: String,
commit_seq: u64,
},
PropRemoved {
key: String,
field: String,
commit_seq: u64,
},
Lagged { missed: u64 },
QueryRowAdded {
columns: Vec<String>,
row: Vec<Option<Value>>,
},
QueryRowRemoved {
columns: Vec<String>,
row: Vec<Option<Value>>,
},
}
pub(crate) struct SubInner {
mu: Mutex<SubQueue>,
condvar: Condvar,
}
impl std::fmt::Debug for SubInner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SubInner").finish_non_exhaustive()
}
}
struct SubQueue {
items: VecDeque<DbEvent>,
missed: u64,
capacity: usize,
}
impl SubInner {
pub(crate) fn new(capacity: usize) -> Arc<Self> {
Arc::new(SubInner {
mu: Mutex::new(SubQueue {
items: VecDeque::new(),
missed: 0,
capacity,
}),
condvar: Condvar::new(),
})
}
pub(crate) fn push(&self, event: DbEvent) {
let mut q = self.mu.lock().unwrap();
if q.items.len() >= q.capacity {
q.missed += 1;
return;
}
q.items.push_back(event);
drop(q);
self.condvar.notify_one();
}
fn pop_one(q: &mut SubQueue) -> Option<DbEvent> {
if let Some(item) = q.items.pop_front() {
return Some(item);
}
if q.missed > 0 {
let missed = std::mem::take(&mut q.missed);
return Some(DbEvent::Lagged { missed });
}
None
}
pub(crate) fn try_recv(&self) -> Option<DbEvent> {
let mut q = self.mu.lock().unwrap();
Self::pop_one(&mut q)
}
pub(crate) fn recv_timeout(&self, timeout: Duration) -> Option<DbEvent> {
let mut q = self.mu.lock().unwrap();
let deadline = std::time::Instant::now() + timeout;
loop {
if let Some(item) = Self::pop_one(&mut q) {
return Some(item);
}
let now = std::time::Instant::now();
if now >= deadline {
return None;
}
let remaining = deadline - now;
let (q2, timed_out) = self.condvar.wait_timeout(q, remaining).unwrap();
q = q2;
if timed_out.timed_out() {
return Self::pop_one(&mut q);
}
}
}
}
pub(crate) enum SubFilter {
Rule(String),
AllRules,
Writes,
}
pub(crate) fn event_matches(event: &DbEvent, filter: &SubFilter) -> bool {
match filter {
SubFilter::Rule(name) => match event {
DbEvent::EdgeFired { rule, .. } | DbEvent::EdgeRetracted { rule, .. } => rule == name,
_ => false,
},
SubFilter::AllRules => {
matches!(
event,
DbEvent::EdgeFired { .. } | DbEvent::EdgeRetracted { .. }
)
}
SubFilter::Writes => matches!(
event,
DbEvent::NodeInserted { .. }
| DbEvent::NodeDeleted { .. }
| DbEvent::EdgeInserted { .. }
| DbEvent::EdgeDeleted { .. }
| DbEvent::PropSet { .. }
| DbEvent::PropRemoved { .. }
),
}
}
pub(crate) struct SubEntry {
pub(crate) filter: SubFilter,
pub(crate) inner: Weak<SubInner>,
}
#[derive(Clone, Debug)]
pub struct Subscription(pub(crate) Arc<SubInner>);
impl Subscription {
pub fn try_recv(&self) -> Option<DbEvent> {
self.0.try_recv()
}
pub fn recv_timeout(&self, timeout: Duration) -> Option<DbEvent> {
self.0.recv_timeout(timeout)
}
}