diff --git a/src/modules/blob/manager.rs b/src/modules/blob/manager.rs index 038d36e..f042ac9 100644 --- a/src/modules/blob/manager.rs +++ b/src/modules/blob/manager.rs @@ -84,13 +84,13 @@ impl EnvelopeIndexManager { Some(doc) => { buffer.push(doc); if buffer.len() >= ENVELOPE_BATCH_SIZE { - ENVELOPE_INDEX_MANAGER.drain_and_commit(&mut buffer).await; + ENVELOPE_INDEX_MANAGER.flush(&mut buffer).await; } } None => { if !buffer.is_empty() { tracing::info!("Channel closed, flushing remaining {} items", buffer.len()); - ENVELOPE_INDEX_MANAGER.drain_and_commit(&mut buffer).await; + ENVELOPE_INDEX_MANAGER.flush(&mut buffer).await; } break; }, @@ -98,11 +98,11 @@ impl EnvelopeIndexManager { } _ = interval.tick() => { if !buffer.is_empty() { - ENVELOPE_INDEX_MANAGER.drain_and_commit(&mut buffer).await; + ENVELOPE_INDEX_MANAGER.flush(&mut buffer).await; } } _ = shutdown.recv() => { - ENVELOPE_INDEX_MANAGER.drain_and_commit(&mut buffer).await; + ENVELOPE_INDEX_MANAGER.flush(&mut buffer).await; break; } } @@ -114,11 +114,11 @@ impl EnvelopeIndexManager { } } - pub async fn add_document(&self, doc: (Envelope, Vec)) { + pub async fn queue(&self, doc: (Envelope, Vec)) { let _ = self.sender.send(doc).await; } - async fn drain_and_commit(&self, buffer: &mut Vec<(Envelope, Vec)>) { + async fn flush(&self, buffer: &mut Vec<(Envelope, Vec)>) { if buffer.is_empty() { return; } diff --git a/src/modules/blob/storage.rs b/src/modules/blob/storage.rs index 64617a9..f943cb7 100644 --- a/src/modules/blob/storage.rs +++ b/src/modules/blob/storage.rs @@ -37,6 +37,45 @@ impl BlobManager { } } + fn process_detached_email( + eml: DetachedEmail, + store: &Database, + email_ks: &Keyspace, + attach_ks: &Keyspace, + ) { + let (email_hash, email_data) = eml.email; + let mut batch = store.batch(); + let mut needs_commit = false; + + match email_ks.contains_key(&email_hash) { + Ok(false) => { + batch.insert(email_ks, email_hash.as_bytes(), email_data); + needs_commit = true; + } + Err(e) => tracing::error!("Fjall email_ks error: {:?}", e), + _ => {} + } + + if let Some(attachments) = eml.attachments { + for (a_hash, a_data) in attachments { + match attach_ks.contains_key(&a_hash) { + Ok(false) => { + batch.insert(attach_ks, a_hash.as_bytes(), a_data); + needs_commit = true; + } + Err(e) => tracing::error!("Fjall attach_ks error: {:?}", e), + _ => {} + } + } + } + + if needs_commit { + if let Err(e) = batch.commit() { + tracing::error!("Fjall Batch Commit Error: {:?}", e); + } + } + } + pub fn new() -> Self { let db = Database::builder(&DATA_DIR_MANAGER.eml_dir) .open() @@ -74,27 +113,16 @@ impl BlobManager { let (sender, mut receiver) = mpsc::channel::(100); let store = db.clone(); - let email_keyspace_clone = email_keyspace.clone(); - let attachments_keyspace_clone = attachments_keyspace.clone(); + let email_ks = email_keyspace.clone(); + let attach_ks = attachments_keyspace.clone(); let handler = task::spawn(async move { let mut shutdown = SIGNAL_MANAGER.subscribe(); loop { tokio::select! { res = receiver.recv() => { match res { - Some(email) => { - let mut batch = store.batch(); - batch.insert(&email_keyspace_clone, email.email.0, email.email.1); - if let Some(attachments) = email.attachments { - for a in attachments { - batch.insert(&attachments_keyspace_clone,a.0, a.1); - } - } - if let Err(e) = batch.commit() { - tracing::error!("Fjall Put Error {:?}", e); - } else { - tracing::info!("Fjall Put Success"); - } + Some(eml) => { + Self::process_detached_email(eml, &store, &email_ks, &attach_ks); } None => { tracing::info!("BlobManager: All senders dropped, closing storage."); @@ -110,19 +138,8 @@ impl BlobManager { remaining ); - while let Some(email) = receiver.recv().await { - let mut batch = store.batch(); - batch.insert(&email_keyspace_clone, email.email.0, email.email.1); - if let Some(attachments) = email.attachments { - for a in attachments { - batch.insert(&attachments_keyspace_clone,a.0, a.1); - } - } - if let Err(e) = batch.commit() { - tracing::error!("Fjall Put Error {:?}", e); - } else { - tracing::info!("Fjall Put Success"); - } + while let Some(eml) = receiver.recv().await { + Self::process_detached_email(eml, &store, &email_ks, &attach_ks); } tracing::info!("BlobManager: All remaining tasks processed. Closing Fjall."); diff --git a/src/modules/dashboard/mod.rs b/src/modules/dashboard/mod.rs index 64e7a93..80bc746 100644 --- a/src/modules/dashboard/mod.rs +++ b/src/modules/dashboard/mod.rs @@ -83,8 +83,9 @@ impl DashboardStats { stat.storage_usage_bytes = get_total_size(&DATA_DIR_MANAGER.eml_dir) .map_err(|e| raise_error!(format!("{:#?}", e), ErrorCode::InternalError))?; - stat.index_usage_bytes = get_total_size(&DATA_DIR_MANAGER.envelope_dir) - .map_err(|e| raise_error!(format!("{:#?}", e), ErrorCode::InternalError))?; + stat.index_usage_bytes = + get_total_size(&DATA_DIR_MANAGER.envelope_dir.join("envelopes.db")) + .map_err(|e| raise_error!(format!("{:#?}", e), ErrorCode::InternalError))?; } else { stat.storage_usage_bytes = 0; stat.index_usage_bytes = 0; diff --git a/src/modules/duckdb/init.rs b/src/modules/duckdb/init.rs index 479bd1f..a75ad34 100644 --- a/src/modules/duckdb/init.rs +++ b/src/modules/duckdb/init.rs @@ -1487,38 +1487,17 @@ impl DuckDBManager { } if let Some(to) = filter.to { - base_sql.push_str( - " - AND EXISTS ( - SELECT 1 FROM UNNEST(e.recipients) r - WHERE r::VARCHAR ILIKE ? - ) - ", - ); + base_sql.push_str(" AND array_to_string(e.recipients, ',') ILIKE ?"); args.push(format!("%{}%", to).into()); } if let Some(cc) = filter.cc { - base_sql.push_str( - " - AND EXISTS ( - SELECT 1 FROM UNNEST(e.cc) r - WHERE r::VARCHAR ILIKE ? - ) - ", - ); + base_sql.push_str(" AND array_to_string(e.cc, ',') ILIKE ?"); args.push(format!("%{}%", cc).into()); } if let Some(bcc) = filter.bcc { - base_sql.push_str( - " - AND EXISTS ( - SELECT 1 FROM UNNEST(e.bcc) r - WHERE r::VARCHAR ILIKE ? - ) - ", - ); + base_sql.push_str(" AND array_to_string(e.bcc, ',') ILIKE ?"); args.push(format!("%{}%", bcc).into()); } diff --git a/src/modules/envelope/extractor.rs b/src/modules/envelope/extractor.rs index 64aaf10..b0435ef 100644 --- a/src/modules/envelope/extractor.rs +++ b/src/modules/envelope/extractor.rs @@ -153,7 +153,7 @@ async fn extract_envelope_core( mailbox_id, uid, subject, - text, + text: String::new(),//for test from, to, cc, @@ -169,9 +169,7 @@ async fn extract_envelope_core( mailbox_name: None, content_hash: email_content_hash, }; - ENVELOPE_INDEX_MANAGER - .add_document((envelope, attachments)) - .await; + ENVELOPE_INDEX_MANAGER.queue((envelope, attachments)).await; Ok(()) } diff --git a/web/src/features/dashboard/index.tsx b/web/src/features/dashboard/index.tsx index 2fab104..21fa42a 100644 --- a/web/src/features/dashboard/index.tsx +++ b/web/src/features/dashboard/index.tsx @@ -196,13 +196,6 @@ export default function MailArchiveDashboard() {
-
-
-

{t('dashboard.title')}

-
-
- -