use std::collections::HashMap;
use anyhow::Result;
use parking_lot::Mutex;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::kvs::{KVValue, Val};
use crate::lq::event::LiveEvents;
use crate::val::{RecordId, TableName, Value};
type PreparedWrite = (NamespaceId, DatabaseId, TableName, Val);
#[derive(Hash, Eq, PartialEq, Debug)]
struct BufferKey {
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
}
pub(crate) struct LiveEventBuffer {
buffer: Mutex<HashMap<BufferKey, LiveEvents>>,
}
impl LiveEventBuffer {
pub(crate) fn new() -> Self {
Self {
buffer: Mutex::new(HashMap::new()),
}
}
pub(crate) fn buffer_record_change(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
id: RecordId,
before: Value,
after: Value,
) {
let mut buffer = self.buffer.lock();
buffer
.entry(BufferKey {
ns,
db,
tb: tb.clone(),
})
.or_insert_with(LiveEvents::new)
.push_record_change(id, before, after);
}
pub(crate) fn changes(&self) -> Result<Vec<PreparedWrite>> {
let buffer = self.buffer.lock();
if buffer.is_empty() {
return Ok(Vec::new());
}
let mut res = Vec::with_capacity(buffer.len());
for (key, events) in buffer.iter() {
let value = events.kv_encode_value()?;
res.push((key.ns, key.db, key.tb.clone(), value));
}
Ok(res)
}
pub(crate) fn clear(&self) {
self.buffer.lock().clear();
}
}