diff --git a/src/common/infra/db/mod.rs b/src/common/infra/db/mod.rs index cdcb720b53..98a2348c89 100644 --- a/src/common/infra/db/mod.rs +++ b/src/common/infra/db/mod.rs @@ -219,6 +219,7 @@ impl std::fmt::Debug for DbEventFileList { pub enum DbEventStreamStats { Set(String, Vec<(String, StreamStats)>), ResetMinTS(String, i64), + ResetAll, } impl std::fmt::Debug for DbEventStreamStats { @@ -226,6 +227,7 @@ impl std::fmt::Debug for DbEventStreamStats { match self { DbEventStreamStats::Set(key, _) => write!(f, "Set({})", key), DbEventStreamStats::ResetMinTS(key, _) => write!(f, "ResetMinTS({})", key), + DbEventStreamStats::ResetAll => write!(f, "ResetAll"), } } } diff --git a/src/common/infra/db/sqlite.rs b/src/common/infra/db/sqlite.rs index d7fc90f078..51fac14b34 100644 --- a/src/common/infra/db/sqlite.rs +++ b/src/common/infra/db/sqlite.rs @@ -309,6 +309,24 @@ impl SqliteDbChannel { log::error!("[SQLITE] reset stream stats min_ts error: {}", e); } } + DbEvent::StreamStats(DbEventStreamStats::ResetAll) => { + let mut err: Option = None; + for _ in 0..DB_RETRY_TIMES { + match sqlite_file_list::reset_stream_stats(&client).await { + Ok(_) => { + err = None; + break; + } + Err(e) => { + err = Some(e.to_string()); + } + } + time::sleep(time::Duration::from_secs(1)).await; + } + if let Some(e) = err { + log::error!("[SQLITE] reset stream stats error: {}", e); + } + } DbEvent::CreateTableMeta => { let mut err: Option = None; for _ in 0..DB_RETRY_TIMES { diff --git a/src/common/infra/file_list/sqlite.rs b/src/common/infra/file_list/sqlite.rs index ec2f645e01..0040b88457 100644 --- a/src/common/infra/file_list/sqlite.rs +++ b/src/common/infra/file_list/sqlite.rs @@ -326,10 +326,10 @@ SELECT stream, MIN(min_ts) as min_ts, MAX(max_ts) as max_ts, COUNT(*) as file_nu } async fn reset_stream_stats(&self) -> Result<()> { - let pool = CLIENT.clone(); - sqlx::query(r#"UPDATE stream_stats SET file_num = 0, min_ts = 0, max_ts = 0, records = 0, original_size = 0, compressed_size = 0;"#) - .execute(&pool) - .await?; + let tx = CHANNEL.db_tx.clone(); + tx.send(DbEvent::StreamStats(DbEventStreamStats::ResetAll)) + .await + .map_err(|e| Error::Message(e.to_string()))?; Ok(()) } @@ -580,6 +580,13 @@ pub async fn reset_stream_stats_min_ts( Ok(()) } +pub async fn reset_stream_stats(client: &Pool) -> Result<()> { + sqlx::query(r#"UPDATE stream_stats SET file_num = 0, min_ts = 0, max_ts = 0, records = 0, original_size = 0, compressed_size = 0;"#) + .execute(client) + .await?; + Ok(()) +} + pub async fn create_table(client: &Pool) -> Result<()> { sqlx::query( r#" diff --git a/src/main.rs b/src/main.rs index daf904e31b..b06ecb73bf 100644 --- a/src/main.rs +++ b/src/main.rs @@ -520,7 +520,11 @@ async fn cli() -> Result { } } - println!("command {name} execute succeeded"); + // flush db + if let Err(e) = infra::db::DEFAULT.close().await { + log::error!("waiting for db close failed, error: {}", e); + } + println!("command {name} execute succeeded"); Ok(true) }