JMAP Registry API implementation - part 7
This commit is contained in:
@@ -10,7 +10,7 @@ use common::{
|
||||
ipc::{BroadcastEvent, HousekeeperEvent, PurgeType},
|
||||
};
|
||||
use email::message::delete::EmailDeletion;
|
||||
use smtp::reporting::SmtpReporting;
|
||||
use smtp::reporting::send::MtaReportSend;
|
||||
use spam_filter::modules::classifier::SpamClassifier;
|
||||
use std::{
|
||||
collections::BinaryHeap,
|
||||
|
||||
@@ -621,7 +621,7 @@ async fn delete_email_metadata(
|
||||
use email::message::metadata::MESSAGE_RECEIVED_MASK;
|
||||
use registry::{
|
||||
pickle::Pickle,
|
||||
schema::structs::{DeletedEmail, DeletedItem},
|
||||
schema::structs::{ArchivedEmail, ArchivedItem},
|
||||
types::{datetime::UTCDateTime, id::ObjectId},
|
||||
};
|
||||
use store::{
|
||||
@@ -660,11 +660,11 @@ async fn delete_email_metadata(
|
||||
let until = now + undelete_retention.as_secs();
|
||||
let blob_hash = BlobHash::from(&metadata.blob_hash);
|
||||
|
||||
let item = DeletedItem::Email(DeletedEmail {
|
||||
let item = ArchivedItem::Email(ArchivedEmail {
|
||||
account_id: account_id.into(),
|
||||
blob_id: BlobId::new(blob_hash.clone(), Default::default()),
|
||||
cleanup_at: UTCDateTime::from_timestamp(until as i64),
|
||||
deleted_at: UTCDateTime::now(),
|
||||
archived_until: UTCDateTime::from_timestamp(until as i64),
|
||||
archived_at: UTCDateTime::now(),
|
||||
from: from.unwrap_or_default(),
|
||||
received_at: UTCDateTime::from_timestamp(
|
||||
(metadata.rcvd_attach.to_native() & MESSAGE_RECEIVED_MASK) as i64,
|
||||
@@ -673,7 +673,7 @@ async fn delete_email_metadata(
|
||||
size: root_part.offset_end.to_native() as u64,
|
||||
})
|
||||
.to_pickled_vec();
|
||||
let object_id = ObjectType::DeletedItem.to_id();
|
||||
let object_id = ObjectType::ArchivedItem.to_id();
|
||||
let item_id = SnowflakeIdGenerator::from_sequence_id(
|
||||
xxhash_rust::xxh3::xxh3_64(item.as_slice()),
|
||||
)
|
||||
@@ -685,7 +685,7 @@ async fn delete_email_metadata(
|
||||
hash: blob_hash,
|
||||
to: BlobLink::Temporary { until },
|
||||
},
|
||||
ObjectId::new(ObjectType::DeletedItem, item_id.into()).serialize(),
|
||||
ObjectId::new(ObjectType::ArchivedItem, item_id.into()).serialize(),
|
||||
)
|
||||
.set(
|
||||
ValueClass::Registry(RegistryClass::Index {
|
||||
|
||||
@@ -9,6 +9,7 @@ use crate::task_manager::index::SearchIndexTask;
|
||||
use crate::task_manager::lock::TaskLockManager;
|
||||
use crate::task_manager::merge_threads::MergeThreadsTask;
|
||||
use crate::task_manager::report::SubmitReportTask;
|
||||
use crate::task_manager::restore_item::RestoreItemTask;
|
||||
use alarm::SendAlarmTask;
|
||||
use common::config::server::ServerProtocol;
|
||||
use common::network::limiter::ConcurrencyLimiter;
|
||||
@@ -45,6 +46,7 @@ pub mod index;
|
||||
pub mod lock;
|
||||
pub mod merge_threads;
|
||||
pub mod report;
|
||||
pub mod restore_item;
|
||||
|
||||
const QUEUE_REFRESH_INTERVAL: u64 = 60 * 5; // 5 minutes
|
||||
const DEFAULT_LOCK_EXPIRY: u64 = 60 * 5; // 5 minutes
|
||||
@@ -219,6 +221,7 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
|
||||
.submit_report(report::ReportId::Tls(task.report_id.id()))
|
||||
.await
|
||||
}
|
||||
Task::RestoreArchivedItem(task) => server.restore_item(task).await,
|
||||
Task::IndexDocument(_)
|
||||
| Task::UnindexDocument(_)
|
||||
| Task::IndexTrace(_) => unreachable!(),
|
||||
@@ -397,7 +400,7 @@ impl TaskQueueManager for Server {
|
||||
TaskType::MergeThreads => roles
|
||||
.merge_threads
|
||||
.is_enabled_for_integer(task_job.id as u32),
|
||||
TaskType::DmarcReport | TaskType::TlsReport => true,
|
||||
TaskType::DmarcReport | TaskType::TlsReport | TaskType::RestoreArchivedItem => true,
|
||||
};
|
||||
|
||||
if enabled {
|
||||
@@ -584,6 +587,7 @@ impl TaskInfo for Task {
|
||||
Task::MergeThreads(_) => "MergeThreads",
|
||||
Task::DmarcReport(_) => "DmarcReport",
|
||||
Task::TlsReport(_) => "TlsReport",
|
||||
Task::RestoreArchivedItem(_) => "RestoreArchivedItem",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
105
crates/services/src/task_manager/restore_item.rs
Normal file
105
crates/services/src/task_manager/restore_item.rs
Normal file
@@ -0,0 +1,105 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||
*/
|
||||
|
||||
use common::{Server, auth::BuildAccessToken};
|
||||
use email::{
|
||||
mailbox::INBOX_ID,
|
||||
message::ingest::{EmailIngest, IngestEmail, IngestSource},
|
||||
};
|
||||
use mail_parser::MessageParser;
|
||||
use registry::schema::{enums::ArchivedItemType, structs::TaskRestoreArchivedItem};
|
||||
use store::write::{BatchBuilder, BlobLink, BlobOp};
|
||||
use trc::AddContext;
|
||||
|
||||
use crate::task_manager::TaskResult;
|
||||
|
||||
pub(crate) trait RestoreItemTask: Sync + Send {
|
||||
fn restore_item(
|
||||
&self,
|
||||
task: &TaskRestoreArchivedItem,
|
||||
) -> impl Future<Output = TaskResult> + Send;
|
||||
}
|
||||
|
||||
impl RestoreItemTask for Server {
|
||||
async fn restore_item(&self, task: &TaskRestoreArchivedItem) -> TaskResult {
|
||||
match restore_item(self, task).await {
|
||||
Ok(result) => result,
|
||||
Err(err) => {
|
||||
let result = TaskResult::temporary(err.to_string());
|
||||
trc::error!(
|
||||
err.account_id(task.account_id.document_id())
|
||||
.details("Failed to restore item")
|
||||
);
|
||||
result
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn restore_item(server: &Server, task: &TaskRestoreArchivedItem) -> trc::Result<TaskResult> {
|
||||
match task.archived_item_type {
|
||||
ArchivedItemType::Email => {
|
||||
let account_id = task.account_id.document_id();
|
||||
let access_token = server
|
||||
.access_token(account_id)
|
||||
.await
|
||||
.caused_by(trc::location!())?;
|
||||
|
||||
let Some(bytes) = server
|
||||
.blob_store()
|
||||
.get_blob(task.blob_id.hash.as_slice(), 0..usize::MAX)
|
||||
.await?
|
||||
else {
|
||||
return Ok(TaskResult::permanent("Blob not found"));
|
||||
};
|
||||
|
||||
match server
|
||||
.email_ingest(IngestEmail {
|
||||
raw_message: &bytes,
|
||||
message: MessageParser::new().parse(&bytes),
|
||||
blob_hash: Some(&task.blob_id.hash),
|
||||
access_token: &access_token.build(),
|
||||
mailbox_ids: vec![INBOX_ID],
|
||||
keywords: vec![],
|
||||
received_at: (task.created_at.timestamp() as u64).into(),
|
||||
source: IngestSource::Restore,
|
||||
session_id: 0,
|
||||
})
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
let mut batch = BatchBuilder::new();
|
||||
batch.with_account_id(account_id).clear(BlobOp::Link {
|
||||
hash: task.blob_id.hash.clone(),
|
||||
to: BlobLink::Temporary {
|
||||
until: task.archived_until.timestamp() as u64,
|
||||
},
|
||||
});
|
||||
server.store().write(batch.build_all()).await?;
|
||||
|
||||
Ok(TaskResult::Success)
|
||||
}
|
||||
Err(mut err)
|
||||
if err.matches(trc::EventType::MessageIngest(
|
||||
trc::MessageIngestEvent::Error,
|
||||
)) =>
|
||||
{
|
||||
Ok(TaskResult::permanent(
|
||||
err.take_value(trc::Key::Reason)
|
||||
.and_then(|v| v.into_string())
|
||||
.unwrap()
|
||||
.to_string(),
|
||||
))
|
||||
}
|
||||
Err(err) => Err(err.caused_by(trc::location!())),
|
||||
}
|
||||
}
|
||||
ArchivedItemType::FileNode
|
||||
| ArchivedItemType::CalendarEvent
|
||||
| ArchivedItemType::ContactCard
|
||||
| ArchivedItemType::SieveScript => Ok(TaskResult::permanent("Not implemented")),
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user