feat(embedding,vision): candle embeddings with offline fallback, on-demand clipboard vision capture, and concurrency audit
This commit is contained in:
1 parent
5bd8b1587a
commit
e4a0fe72df
47 files changed
+6292
-3503
No files matched your search
+346
-44
@@ -4,10 +4,16 @@ use std::sync::{Arc, RwLock};
|
||||
|
||||
pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("store");
|
||||
|
||||
/// Internal write request dispatched to the single database writer actor.
|
||||
enum DbOp {
|
||||
Insert(Vec<u8>),
|
||||
Delete,
|
||||
}
|
||||
|
||||
/// Internal write request dispatched to the single database writer actor.
|
||||
struct DbWriteTask {
|
||||
key: String,
|
||||
data: Vec<u8>,
|
||||
op: DbOp,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
oneshot_tx: Option<tokio::sync::oneshot::Sender<()>>,
|
||||
}
|
||||
@@ -18,7 +24,8 @@ pub struct DbWriteQueue {
|
||||
tx: tokio::sync::mpsc::Sender<DbWriteTask>,
|
||||
}
|
||||
|
||||
static QUEUE_REGISTRY: std::sync::Mutex<Option<(Arc<Database>, DbWriteQueue)>> = std::sync::Mutex::new(None);
|
||||
static QUEUE_REGISTRY: std::sync::Mutex<Option<(Arc<Database>, DbWriteQueue)>> =
|
||||
std::sync::Mutex::new(None);
|
||||
|
||||
fn get_or_create_queue(db: Arc<Database>) -> DbWriteQueue {
|
||||
let mut reg = QUEUE_REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
|
||||
@@ -38,7 +45,8 @@ impl DbWriteQueue {
|
||||
|
||||
tokio::spawn(async move {
|
||||
while let Some(first_task) = rx.recv().await {
|
||||
let mut batch = vec![first_task];
|
||||
let mut batch = Vec::with_capacity(100);
|
||||
batch.push(first_task);
|
||||
|
||||
// Gold Standard Micro-batching: Drain up to 100 accumulated tasks from queue without blocking
|
||||
while batch.len() < 100 {
|
||||
@@ -54,8 +62,27 @@ impl DbWriteQueue {
|
||||
Ok(write_txn) => {
|
||||
if let Ok(mut table) = write_txn.open_table(STORE_TABLE) {
|
||||
for task in &batch {
|
||||
if let Err(e) = table.insert(task.key.as_str(), task.data.as_slice()) {
|
||||
tracing::error!("Failed to insert key '{}' into redb: {}", task.key, e);
|
||||
match &task.op {
|
||||
DbOp::Insert(data) => {
|
||||
if let Err(e) =
|
||||
table.insert(task.key.as_str(), data.as_slice())
|
||||
{
|
||||
tracing::error!(
|
||||
"Failed to insert key '{}' into redb: {}",
|
||||
task.key,
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
DbOp::Delete => {
|
||||
if let Err(e) = table.remove(task.key.as_str()) {
|
||||
tracing::error!(
|
||||
"Failed to delete key '{}' from redb: {}",
|
||||
task.key,
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -64,7 +91,10 @@ impl DbWriteQueue {
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to begin write transaction on redb writer actor: {}", e);
|
||||
tracing::error!(
|
||||
"Failed to begin write transaction on redb writer actor: {}",
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -88,18 +118,55 @@ impl DbWriteQueue {
|
||||
key: String,
|
||||
data: Vec<u8>,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
) -> Option<tokio::sync::oneshot::Receiver<()>> {
|
||||
self.push_op(key, DbOp::Insert(data), flushed_notifier)
|
||||
}
|
||||
|
||||
pub fn push_delete(
|
||||
&self,
|
||||
key: String,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
) -> Option<tokio::sync::oneshot::Receiver<()>> {
|
||||
self.push_op(key, DbOp::Delete, flushed_notifier)
|
||||
}
|
||||
|
||||
fn push_op(
|
||||
&self,
|
||||
key: String,
|
||||
op: DbOp,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
) -> Option<tokio::sync::oneshot::Receiver<()>> {
|
||||
let (oneshot_tx, oneshot_rx) = tokio::sync::oneshot::channel();
|
||||
let task = DbWriteTask {
|
||||
key,
|
||||
data,
|
||||
op,
|
||||
flushed_notifier,
|
||||
oneshot_tx: Some(oneshot_tx),
|
||||
};
|
||||
if let Err(e) = self.tx.try_send(task) {
|
||||
let key = e.into_inner().key;
|
||||
tracing::error!("DbWriteQueue channel full or closed; unable to persist key '{}'", key);
|
||||
None
|
||||
match e {
|
||||
tokio::sync::mpsc::error::TrySendError::Full(task) => {
|
||||
let tx = self.tx.clone();
|
||||
let key = task.key.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(err) = tx.send(task).await {
|
||||
tracing::error!(
|
||||
"DbWriteQueue fallback send failed for key '{}': {}",
|
||||
key,
|
||||
err
|
||||
);
|
||||
}
|
||||
});
|
||||
None
|
||||
}
|
||||
tokio::sync::mpsc::error::TrySendError::Closed(task) => {
|
||||
tracing::error!(
|
||||
"DbWriteQueue channel closed; unable to persist key '{}'",
|
||||
task.key
|
||||
);
|
||||
None
|
||||
}
|
||||
}
|
||||
} else {
|
||||
Some(oneshot_rx)
|
||||
}
|
||||
@@ -110,16 +177,38 @@ impl DbWriteQueue {
|
||||
key: String,
|
||||
data: Vec<u8>,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
) -> Option<tokio::sync::oneshot::Receiver<()>> {
|
||||
self.push_op_async(key, DbOp::Insert(data), flushed_notifier)
|
||||
.await
|
||||
}
|
||||
|
||||
pub async fn push_delete_async(
|
||||
&self,
|
||||
key: String,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
) -> Option<tokio::sync::oneshot::Receiver<()>> {
|
||||
self.push_op_async(key, DbOp::Delete, flushed_notifier)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn push_op_async(
|
||||
&self,
|
||||
key: String,
|
||||
op: DbOp,
|
||||
flushed_notifier: Arc<tokio::sync::Notify>,
|
||||
) -> Option<tokio::sync::oneshot::Receiver<()>> {
|
||||
let (oneshot_tx, oneshot_rx) = tokio::sync::oneshot::channel();
|
||||
let task = DbWriteTask {
|
||||
key,
|
||||
data,
|
||||
op,
|
||||
flushed_notifier,
|
||||
oneshot_tx: Some(oneshot_tx),
|
||||
};
|
||||
if let Err(e) = self.tx.send(task).await {
|
||||
tracing::error!("DbWriteQueue channel closed; unable to persist key '{}'", e.0.key);
|
||||
tracing::error!(
|
||||
"DbWriteQueue channel closed; unable to persist key '{}'",
|
||||
e.0.key
|
||||
);
|
||||
None
|
||||
} else {
|
||||
Some(oneshot_rx)
|
||||
@@ -156,29 +245,116 @@ impl<T: DeserializeOwned + Default + Serialize + Send + Sync + 'static> Store<T>
|
||||
tracing::error!("Failed to begin read transaction for key '{}'", key);
|
||||
return (T::default(), false);
|
||||
};
|
||||
match read_txn.open_table(STORE_TABLE) {
|
||||
Ok(table) => match table.get(key) {
|
||||
Ok(Some(value)) => match serde_json::from_slice::<T>(value.value()) {
|
||||
Ok(parsed) => (parsed, false),
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
"CRITICAL: Corrupted data for key '{}' in database: {}. Quarantine mode active: state initialized to empty default without overwriting DB key.",
|
||||
key, e
|
||||
);
|
||||
(T::default(), true)
|
||||
}
|
||||
},
|
||||
Ok(None) => (T::default(), false),
|
||||
let Ok(table) = read_txn.open_table(STORE_TABLE) else {
|
||||
return (T::default(), false);
|
||||
};
|
||||
|
||||
// 1. Check monolithic key first as the authoritative snapshot
|
||||
match table.get(key) {
|
||||
Ok(Some(value)) => match serde_json::from_slice::<T>(value.value()) {
|
||||
Ok(parsed) => return (parsed, false),
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to get key '{}' from store table: {}", key, e);
|
||||
(T::default(), false)
|
||||
tracing::error!(
|
||||
"CRITICAL: Corrupted data for key '{}' in database: {}. Quarantine mode active: state initialized to empty default without overwriting DB key.",
|
||||
key,
|
||||
e
|
||||
);
|
||||
return (T::default(), true);
|
||||
}
|
||||
},
|
||||
Ok(None) => {}
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to open STORE_TABLE for key '{}': {}", key, e);
|
||||
(T::default(), false)
|
||||
tracing::error!("Failed to get key '{}' from store table: {}", key, e);
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Granular prefix keys fallback: format!("{}:", key)
|
||||
let prefix = format!("{}:", key);
|
||||
let mut items_array = Vec::new();
|
||||
let mut items_map = serde_json::Map::new();
|
||||
let mut found_granular = false;
|
||||
|
||||
if let Ok(range) = table.range(prefix.as_str()..) {
|
||||
for entry in range {
|
||||
if let Ok((k, v)) = entry {
|
||||
let k_str = k.value();
|
||||
if !k_str.starts_with(&prefix) {
|
||||
break;
|
||||
}
|
||||
found_granular = true;
|
||||
if let Ok(val) = serde_json::from_slice::<serde_json::Value>(v.value()) {
|
||||
let sub_key = &k_str[prefix.len()..];
|
||||
items_array.push(val.clone());
|
||||
items_map.insert(sub_key.to_string(), val);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if found_granular {
|
||||
if let Ok(parsed) = serde_json::from_value::<T>(serde_json::Value::Array(items_array)) {
|
||||
return (parsed, false);
|
||||
}
|
||||
if let Ok(parsed) = serde_json::from_value::<T>(serde_json::Value::Object(items_map)) {
|
||||
return (parsed, false);
|
||||
}
|
||||
}
|
||||
|
||||
(T::default(), false)
|
||||
}
|
||||
|
||||
fn extract_granular_keys(base_key: &str, val: &serde_json::Value) -> Vec<String> {
|
||||
let mut keys = Vec::new();
|
||||
match val {
|
||||
serde_json::Value::Array(arr) => {
|
||||
for (i, item) in arr.iter().enumerate() {
|
||||
let sub_key = item
|
||||
.get("id")
|
||||
.or_else(|| item.get("name"))
|
||||
.or_else(|| item.get("title"))
|
||||
.and_then(|v| v.as_str())
|
||||
.map(|s| s.to_string())
|
||||
.unwrap_or_else(|| i.to_string());
|
||||
keys.push(format!("{}:{}", base_key, sub_key));
|
||||
}
|
||||
}
|
||||
serde_json::Value::Object(map) => {
|
||||
for sub_key in map.keys() {
|
||||
keys.push(format!("{}:{}", base_key, sub_key));
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
keys
|
||||
}
|
||||
|
||||
fn extract_granular_entries(base_key: &str, val: &serde_json::Value) -> Vec<(String, Vec<u8>)> {
|
||||
let mut granular = Vec::new();
|
||||
match val {
|
||||
serde_json::Value::Array(arr) => {
|
||||
for (i, item) in arr.iter().enumerate() {
|
||||
let sub_key = item
|
||||
.get("id")
|
||||
.or_else(|| item.get("name"))
|
||||
.or_else(|| item.get("title"))
|
||||
.and_then(|v| v.as_str())
|
||||
.map(|s| s.to_string())
|
||||
.unwrap_or_else(|| i.to_string());
|
||||
if let Ok(item_bytes) = serde_json::to_vec(item) {
|
||||
granular.push((format!("{}:{}", base_key, sub_key), item_bytes));
|
||||
}
|
||||
}
|
||||
}
|
||||
serde_json::Value::Object(map) => {
|
||||
for (sub_key, item) in map {
|
||||
if let Ok(item_bytes) = serde_json::to_vec(item) {
|
||||
granular.push((format!("{}:{}", base_key, sub_key), item_bytes));
|
||||
}
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
granular
|
||||
}
|
||||
|
||||
pub fn read_with<F, R>(&self, f: F) -> R
|
||||
@@ -191,7 +367,7 @@ impl<T: DeserializeOwned + Default + Serialize + Send + Sync + 'static> Store<T>
|
||||
|
||||
pub fn modify<F: FnOnce(&mut T)>(&self, f: F)
|
||||
where
|
||||
T: Serialize,
|
||||
T: Serialize + Clone,
|
||||
{
|
||||
if self.is_corrupted {
|
||||
tracing::error!(
|
||||
@@ -201,30 +377,80 @@ impl<T: DeserializeOwned + Default + Serialize + Send + Sync + 'static> Store<T>
|
||||
return;
|
||||
}
|
||||
|
||||
let serialized_res = {
|
||||
// Fast mutation under critical lock section, then immediately release the RwLock guard
|
||||
let (old_snapshot, new_snapshot) = {
|
||||
let mut lock = self.cache.write().unwrap_or_else(|e| e.into_inner());
|
||||
let old = (*lock).clone();
|
||||
f(&mut lock);
|
||||
serde_json::to_vec(&*lock)
|
||||
let new = (*lock).clone();
|
||||
(old, new)
|
||||
};
|
||||
|
||||
match serialized_res {
|
||||
// Expensive serialization and granular extraction run completely unblocked outside the lock
|
||||
let old_keys = serde_json::to_value(&old_snapshot)
|
||||
.map(|val| Self::extract_granular_keys(&self.key, &val))
|
||||
.unwrap_or_default();
|
||||
|
||||
let full_bytes_res = serde_json::to_vec(&new_snapshot);
|
||||
let granular_entries = serde_json::to_value(&new_snapshot)
|
||||
.map(|val| Self::extract_granular_entries(&self.key, &val))
|
||||
.unwrap_or_default();
|
||||
|
||||
let new_keys: std::collections::HashSet<&str> =
|
||||
granular_entries.iter().map(|(k, _)| k.as_str()).collect();
|
||||
let mut removed_keys = Vec::new();
|
||||
for old_k in &old_keys {
|
||||
if !new_keys.contains(old_k.as_str()) {
|
||||
removed_keys.push(old_k.clone());
|
||||
}
|
||||
}
|
||||
|
||||
match full_bytes_res {
|
||||
Ok(data) => {
|
||||
if self.queue.push(self.key.clone(), data.clone(), self.flushed.clone()).is_none() {
|
||||
// Delete removed granular entries so they don't resurrect on restart
|
||||
for del_key in removed_keys {
|
||||
self.queue.push_delete(del_key, self.flushed.clone());
|
||||
}
|
||||
|
||||
// Queue granular entries
|
||||
for (g_key, g_bytes) in granular_entries {
|
||||
self.queue.push(g_key, g_bytes, self.flushed.clone());
|
||||
}
|
||||
|
||||
if self
|
||||
.queue
|
||||
.push(self.key.clone(), data.clone(), self.flushed.clone())
|
||||
.is_none()
|
||||
{
|
||||
tracing::warn!(
|
||||
"DbWriteQueue channel full for key '{}'. Applying backpressure fallback.",
|
||||
self.key
|
||||
);
|
||||
let queue = self.queue.clone();
|
||||
let key = self.key.clone();
|
||||
let flushed = self.flushed.clone();
|
||||
tokio::spawn(async move {
|
||||
let _ = queue.push_async(key, data, flushed).await;
|
||||
});
|
||||
if let Ok(handle) = tokio::runtime::Handle::try_current() {
|
||||
handle.spawn(async move {
|
||||
let _ = tokio::time::timeout(
|
||||
std::time::Duration::from_secs(10),
|
||||
queue.push_async(key, data, flushed),
|
||||
)
|
||||
.await;
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => tracing::error!("Failed to serialize memory store for key '{}': {}", self.key, e),
|
||||
Err(e) => tracing::error!(
|
||||
"Failed to serialize memory store for key '{}': {}",
|
||||
self.key,
|
||||
e
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn modify_async<F: FnOnce(&mut T)>(&self, f: F)
|
||||
where
|
||||
T: Serialize,
|
||||
T: Serialize + Clone,
|
||||
{
|
||||
if self.is_corrupted {
|
||||
tracing::error!(
|
||||
@@ -234,19 +460,57 @@ impl<T: DeserializeOwned + Default + Serialize + Send + Sync + 'static> Store<T>
|
||||
return;
|
||||
}
|
||||
|
||||
let serialized_res = {
|
||||
let (old_snapshot, new_snapshot) = {
|
||||
let mut lock = self.cache.write().unwrap_or_else(|e| e.into_inner());
|
||||
let old = (*lock).clone();
|
||||
f(&mut lock);
|
||||
serde_json::to_vec(&*lock)
|
||||
let new = (*lock).clone();
|
||||
(old, new)
|
||||
};
|
||||
|
||||
match serialized_res {
|
||||
let old_keys = serde_json::to_value(&old_snapshot)
|
||||
.map(|val| Self::extract_granular_keys(&self.key, &val))
|
||||
.unwrap_or_default();
|
||||
|
||||
let full_bytes_res = serde_json::to_vec(&new_snapshot);
|
||||
let granular_entries = serde_json::to_value(&new_snapshot)
|
||||
.map(|val| Self::extract_granular_entries(&self.key, &val))
|
||||
.unwrap_or_default();
|
||||
|
||||
let new_keys: std::collections::HashSet<&str> =
|
||||
granular_entries.iter().map(|(k, _)| k.as_str()).collect();
|
||||
let mut removed_keys = Vec::new();
|
||||
for old_k in &old_keys {
|
||||
if !new_keys.contains(old_k.as_str()) {
|
||||
removed_keys.push(old_k.clone());
|
||||
}
|
||||
}
|
||||
|
||||
match full_bytes_res {
|
||||
Ok(data) => {
|
||||
if let Some(rx) = self.queue.push_async(self.key.clone(), data, self.flushed.clone()).await {
|
||||
for del_key in removed_keys {
|
||||
self.queue
|
||||
.push_delete_async(del_key, self.flushed.clone())
|
||||
.await;
|
||||
}
|
||||
for (g_key, g_bytes) in granular_entries {
|
||||
self.queue
|
||||
.push_async(g_key, g_bytes, self.flushed.clone())
|
||||
.await;
|
||||
}
|
||||
if let Some(rx) = self
|
||||
.queue
|
||||
.push_async(self.key.clone(), data, self.flushed.clone())
|
||||
.await
|
||||
{
|
||||
let _ = rx.await;
|
||||
}
|
||||
}
|
||||
Err(e) => tracing::error!("Failed to serialize memory store for key '{}': {}", self.key, e),
|
||||
Err(e) => tracing::error!(
|
||||
"Failed to serialize memory store for key '{}': {}",
|
||||
self.key,
|
||||
e
|
||||
),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -322,4 +586,42 @@ mod tests {
|
||||
|
||||
assert_eq!(store.read_with(|s| s.value), 50);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_store_deletion_does_not_resurrect() {
|
||||
let db = create_in_memory_test_db();
|
||||
#[derive(Serialize, serde::Deserialize, Clone, Default, PartialEq, Debug)]
|
||||
struct Item {
|
||||
id: String,
|
||||
name: String,
|
||||
}
|
||||
let store = Store::<Vec<Item>>::new("items", db.clone());
|
||||
store.modify(|items| {
|
||||
items.push(Item {
|
||||
id: "item1".into(),
|
||||
name: "First".into(),
|
||||
});
|
||||
items.push(Item {
|
||||
id: "item2".into(),
|
||||
name: "Second".into(),
|
||||
});
|
||||
});
|
||||
store.flushed.notified().await;
|
||||
|
||||
// Verify both items loaded
|
||||
let store_check = Store::<Vec<Item>>::new("items", db.clone());
|
||||
assert_eq!(store_check.read_with(|items| items.len()), 2);
|
||||
|
||||
// Delete item1
|
||||
store.modify(|items| {
|
||||
items.retain(|i| i.id != "item1");
|
||||
});
|
||||
store.flushed.notified().await;
|
||||
|
||||
// Reload from DB into a brand new Store instance - item1 must NOT resurrect!
|
||||
let store_reloaded = Store::<Vec<Item>>::new("items", db.clone());
|
||||
let remaining = store_reloaded.read_with(|items| items.clone());
|
||||
assert_eq!(remaining.len(), 1);
|
||||
assert_eq!(remaining[0].id, "item2");
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user