use std::sync::Arc;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use crate::prelude::*;
use crate::scheduler::{Task, TaskId};
const SHARED_TN: TnId = TnId(0);
const DEFAULT_CRON: &str = "20 4 * * *";
const DEFAULT_MIN_FREE_PCT: i64 = 20;
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct DbMaintenanceTask {
#[serde(default, skip_serializing)]
pub notify_tn: Option<TnId>,
}
#[async_trait]
impl Task<App> for DbMaintenanceTask {
fn kind() -> &'static str {
"core.db_maintenance"
}
fn kind_of(&self) -> &'static str {
Self::kind()
}
fn build(_id: TaskId, ctx: &str) -> ClResult<Arc<dyn Task<App>>> {
if ctx.is_empty() {
Ok(Arc::new(Self::default()))
} else {
Ok(Arc::new(serde_json::from_str::<Self>(ctx)?))
}
}
fn serialize(&self) -> String {
serde_json::to_string(self).unwrap_or_default()
}
async fn run(&self, app: &App) -> ClResult<()> {
let min_free_pct = app
.settings
.get_int_opt(SHARED_TN, "core.vacuum_min_free_pct")
.await
.ok()
.flatten()
.unwrap_or(DEFAULT_MIN_FREE_PCT);
if let Err(e) = app.meta_adapter.optimize_search_index(false).await {
warn!(error = %e, "db_maintenance: FTS merge failed");
}
let report = match app.meta_adapter.reclaim_space(min_free_pct).await {
Ok(report) => {
info!(
page_size = report.page_size,
page_count = report.page_count,
freelist_count = report.freelist_count,
vacuumed = report.vacuumed,
min_free_pct,
"db_maintenance: sweep complete"
);
Some(report)
}
Err(e) => {
warn!(error = %e, "db_maintenance: space reclaim failed");
None
}
};
for (what, result) in [
("rtdb", app.rtdb_adapter.compact_storage().await),
("crdt", app.crdt_adapter.compact_storage().await),
] {
match result {
Ok(r) => info!(
store = what,
files = r.files,
bytes_before = r.bytes_before,
bytes_after = r.bytes_after,
"db_maintenance: compacted"
),
Err(e) => warn!(store = what, error = %e, "db_maintenance: compaction failed"),
}
}
if let Some(tn_id) = self.notify_tn {
let data = match report {
Some(r) => serde_json::json!({
"ok": true,
"vacuumed": r.vacuumed,
"pageSize": r.page_size,
"pageCount": r.page_count,
"freelistCount": r.freelist_count,
}),
None => serde_json::json!({ "ok": false }),
};
let msg =
crate::ws_broadcast::BroadcastMessage::new("DB_MAINTENANCE_DONE", data, "system");
let delivered = app.broadcast.send_to_tenant(tn_id, msg).await;
debug!(tn_id = %tn_id, delivered, "db_maintenance outcome broadcast");
}
Ok(())
}
}
pub fn init(app: &App) -> ClResult<()> {
app.scheduler.register::<DbMaintenanceTask>()?;
Ok(())
}
pub async fn schedule(app: &App) -> ClResult<()> {
let cron = app
.settings
.get_string_opt(SHARED_TN, "core.db_maintenance_cron")
.await
.ok()
.flatten()
.unwrap_or_else(|| DEFAULT_CRON.to_string());
let task: Arc<dyn Task<App>> = Arc::new(DbMaintenanceTask::default());
app.scheduler
.task(task)
.key("core.db_maintenance")
.cron(cron)
.run_on_startup()
.schedule()
.await?;
Ok(())
}