From 2e2d1b3506eba0cf471ba4968b9528cd336d520e Mon Sep 17 00:00:00 2001 From: YoursFunny Date: Thu, 13 Aug 2026 22:25:13 +0800 Subject: [PATCH] style: rustfmt the DbPool call sites from the connection-pool change Formatting-only; the pool commit landed before cargo fmt was run. --- crates/xmedia-bot/src/link_cache.rs | 90 ++++++++++++++++------------- crates/xmedia-bot/src/queue.rs | 37 ++++++------ crates/xmedia-bot/src/state.rs | 54 +++++++++-------- 3 files changed, 99 insertions(+), 82 deletions(-) diff --git a/crates/xmedia-bot/src/link_cache.rs b/crates/xmedia-bot/src/link_cache.rs index a5eb43c..839d24e 100644 --- a/crates/xmedia-bot/src/link_cache.rs +++ b/crates/xmedia-bot/src/link_cache.rs @@ -70,24 +70,26 @@ impl LinkCache { pub async fn get(&self, key: &str, ttl: Duration) -> Option { let key = key.to_string(); let ttl = ttl.as_secs_f64(); - let result = self.pool.with_conn(move |conn| { - let mut stmt = - conn.prepare("SELECT payload, created_at FROM link_cache WHERE url = ?1")?; - let mut rows = stmt.query(params![key])?; - let Some(row) = rows.next()? else { - return Ok(None); - }; - let payload: String = row.get(0)?; - let created_at: f64 = row.get(1)?; - if now_f64() - created_at > ttl { - conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key])?; - return Ok(None); - } - Ok(Some(serde_json::from_str::(&payload).map_err( - |e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)), - )?)) - }) - .await; + let result = self + .pool + .with_conn(move |conn| { + let mut stmt = + conn.prepare("SELECT payload, created_at FROM link_cache WHERE url = ?1")?; + let mut rows = stmt.query(params![key])?; + let Some(row) = rows.next()? else { + return Ok(None); + }; + let payload: String = row.get(0)?; + let created_at: f64 = row.get(1)?; + if now_f64() - created_at > ttl { + conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key])?; + return Ok(None); + } + Ok(Some(serde_json::from_str::(&payload).map_err( + |e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)), + )?)) + }) + .await; match result { Ok(v) => v, Err(e) => { @@ -100,14 +102,16 @@ impl LinkCache { pub async fn put(&self, key: &str, post: &CachedPost) { let key = key.to_string(); let payload = serde_json::to_string(post).expect("cached post serializes"); - let result = self.pool.with_conn(move |conn| { - conn.execute( + let result = self + .pool + .with_conn(move |conn| { + conn.execute( "INSERT OR REPLACE INTO link_cache (url, payload, created_at) VALUES (?1, ?2, ?3)", params![key, payload, now_f64()], )?; - Ok(()) - }) - .await; + Ok(()) + }) + .await; if let Err(e) = result { log::error!("link cache write failed: {e}"); } @@ -116,11 +120,13 @@ impl LinkCache { /// Drops an entry (e.g. a cached file id that turned out invalid). pub async fn remove(&self, key: &str) { let key = key.to_string(); - let result = self.pool.with_conn(move |conn| { - conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key])?; - Ok(()) - }) - .await; + let result = self + .pool + .with_conn(move |conn| { + conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key])?; + Ok(()) + }) + .await; if let Err(e) = result { log::error!("link cache delete failed: {e}"); } @@ -129,13 +135,15 @@ impl LinkCache { /// Removes expired entries; returns how many were deleted. pub async fn prune(&self, ttl: Duration) -> usize { let cutoff = now_f64() - ttl.as_secs_f64(); - let result = self.pool.with_conn(move |conn| { - conn.execute( - "DELETE FROM link_cache WHERE created_at < ?1", - params![cutoff], - ) - }) - .await; + let result = self + .pool + .with_conn(move |conn| { + conn.execute( + "DELETE FROM link_cache WHERE created_at < ?1", + params![cutoff], + ) + }) + .await; match result { Ok(n) => n, Err(e) => { @@ -149,11 +157,13 @@ impl LinkCache { /// `key` is `None`. Returns how many rows were removed. pub async fn clear(&self, key: Option<&str>) -> usize { let key = key.map(str::to_string); - let result = self.pool.with_conn(move |conn| match &key { - Some(key) => conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key]), - None => conn.execute("DELETE FROM link_cache", []), - }) - .await; + let result = self + .pool + .with_conn(move |conn| match &key { + Some(key) => conn.execute("DELETE FROM link_cache WHERE url = ?1", params![key]), + None => conn.execute("DELETE FROM link_cache", []), + }) + .await; match result { Ok(n) => n, Err(e) => { diff --git a/crates/xmedia-bot/src/queue.rs b/crates/xmedia-bot/src/queue.rs index 8606f72..c9e7591 100644 --- a/crates/xmedia-bot/src/queue.rs +++ b/crates/xmedia-bot/src/queue.rs @@ -160,8 +160,7 @@ impl PersistentTaskQueue { if sweep_stop.load(Ordering::Relaxed) { break; } - let result = - sweep_pool.with_conn(move |conn| recover_update(conn)).await; + let result = sweep_pool.with_conn(move |conn| recover_update(conn)).await; if let Err(e) = result { log::error!("queue sweep failed: {e}"); } @@ -309,16 +308,18 @@ impl QueueWorker { } async fn earliest_run_after(&self) -> Option { - let result = self.pool.with_conn(|conn| { - let mut stmt = - conn.prepare("SELECT MIN(run_after) FROM tasks WHERE status='pending'")?; - let mut rows = stmt.query([])?; - match rows.next()? { - Some(row) => Ok(row.get::<_, Option>(0)?), - None => Ok(None), - } - }) - .await; + let result = self + .pool + .with_conn(|conn| { + let mut stmt = + conn.prepare("SELECT MIN(run_after) FROM tasks WHERE status='pending'")?; + let mut rows = stmt.query([])?; + match rows.next()? { + Some(row) => Ok(row.get::<_, Option>(0)?), + None => Ok(None), + } + }) + .await; match result { Ok(v) => v, Err(e) => { @@ -374,11 +375,13 @@ impl QueueWorker { async fn delete_row(&self, id: &str) { let id = id.to_string(); - let result = self.pool.with_conn(move |conn| { - conn.execute("DELETE FROM tasks WHERE id = ?1", params![id])?; - Ok(()) - }) - .await; + let result = self + .pool + .with_conn(move |conn| { + conn.execute("DELETE FROM tasks WHERE id = ?1", params![id])?; + Ok(()) + }) + .await; if let Err(e) = result { log::error!("queue delete failed: {e}"); } diff --git a/crates/xmedia-bot/src/state.rs b/crates/xmedia-bot/src/state.rs index d26dcde..bbc8543 100644 --- a/crates/xmedia-bot/src/state.rs +++ b/crates/xmedia-bot/src/state.rs @@ -76,23 +76,25 @@ impl ChatStore { return data.clone(); } let chat_key = chat_id.to_string(); - let payload = self.pool.with_conn(move |conn| { - // Concurrent handler tasks (batch-forwards) may write chat_state - // while this read runs; the shared busy timeout handles the - // write-lock collision instead of failing the query. - let mut stmt = conn.prepare("SELECT payload FROM chat_state WHERE chat_id = ?1")?; - let mut rows = stmt.query(params![chat_key])?; - match rows.next()? { - Some(row) => Ok(Some(row.get::<_, String>(0)?)), - None => Ok(None), - } - }) - .await - .unwrap_or_else(|e| { - log::error!("chat_state read failed: {e}"); - None - }) - .unwrap_or_default(); + let payload = self + .pool + .with_conn(move |conn| { + // Concurrent handler tasks (batch-forwards) may write chat_state + // while this read runs; the shared busy timeout handles the + // write-lock collision instead of failing the query. + let mut stmt = conn.prepare("SELECT payload FROM chat_state WHERE chat_id = ?1")?; + let mut rows = stmt.query(params![chat_key])?; + match rows.next()? { + Some(row) => Ok(Some(row.get::<_, String>(0)?)), + None => Ok(None), + } + }) + .await + .unwrap_or_else(|e| { + log::error!("chat_state read failed: {e}"); + None + }) + .unwrap_or_default(); let data: ChatData = serde_json::from_str(&payload).unwrap_or_default(); self.cache.lock().insert(chat_id, data.clone()); data @@ -103,14 +105,16 @@ impl ChatStore { self.cache.lock().insert(chat_id, data.clone()); let payload = serde_json::to_string(data).expect("chat state serializes"); let chat_id = chat_id.to_string(); - let result = self.pool.with_conn(move |conn| { - conn.execute( - "INSERT OR REPLACE INTO chat_state (chat_id, payload) VALUES (?1, ?2)", - params![chat_id, payload], - )?; - Ok(()) - }) - .await; + let result = self + .pool + .with_conn(move |conn| { + conn.execute( + "INSERT OR REPLACE INTO chat_state (chat_id, payload) VALUES (?1, ?2)", + params![chat_id, payload], + )?; + Ok(()) + }) + .await; if let Err(e) = result { log::error!("chat_state write failed: {e}"); }