From c7ee6bb519071cbc56a72a895fe6e122840d6d00 Mon Sep 17 00:00:00 2001 From: oudecheng <13802883547@139.com> Date: Mon, 3 Aug 2026 22:15:13 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E5=B9=B6=E5=8F=91=20s?= =?UTF-8?q?ub-agent=20=E6=8C=81=E4=B9=85=E5=8C=96=E6=97=B6=20SQLite=20data?= =?UTF-8?q?base=20is=20locked=20=E9=94=99=E8=AF=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将 7 个写事务从 BEGIN DEFERRED 改为 BEGIN IMMEDIATE,在事务开始即获取写锁,消除多 sub-agent 并发写入时的死锁路径。同时将 busy_timeout 从 5s 提升至 30s,为并发写者排队提供 100 倍余量。 根因:BEGIN DEFERRED 下多个事务可同时读 MAX(seq) 不持写锁,提交时互相阻塞,5s timeout 耗尽后返回 SQLITE_BUSY。BEGIN IMMEDIATE 强制写者串行排队,顺带消除 MAX(seq)+1 竞态导致的 UNIQUE 约束冲突。 --- src/storage/mod.rs | 35 ++++++++++++++++++----------------- 1 file changed, 18 insertions(+), 17 deletions(-) diff --git a/src/storage/mod.rs b/src/storage/mod.rs index a62cd34..ead4bbe 100644 --- a/src/storage/mod.rs +++ b/src/storage/mod.rs @@ -3,7 +3,7 @@ use std::path::{Path, PathBuf}; use r2d2::Pool; use r2d2_sqlite::SqliteConnectionManager; -use rusqlite::{Connection, OptionalExtension, params}; +use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params}; use crate::bus::ChatMessage; @@ -62,7 +62,7 @@ impl SessionStore { /// The connection is used for schema initialization only; the pool /// manages subsequent connections using the same file path. fn from_connection(conn: Connection, db_uri: &str) -> Result { - conn.busy_timeout(std::time::Duration::from_secs(5))?; + conn.busy_timeout(std::time::Duration::from_secs(30))?; conn.execute_batch( " PRAGMA journal_mode = WAL; @@ -230,7 +230,7 @@ impl SessionStore { let manager = SqliteConnectionManager::file(db_uri) .with_init(|c| { - c.busy_timeout(std::time::Duration::from_secs(5))?; + c.busy_timeout(std::time::Duration::from_secs(30))?; Ok(()) }); let pool = Pool::builder() @@ -575,8 +575,8 @@ impl SessionStore { topic_id: Option<&str>, message: &ChatMessage, ) -> Result<(), StorageError> { - let conn = self.pool.get()?; - let tx = conn.unchecked_transaction()?; + let mut conn = self.pool.get()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let seq: i64 = tx.query_row( "SELECT COALESCE(MAX(seq), 0) + 1 FROM messages WHERE session_id = ?1", @@ -651,8 +651,8 @@ impl SessionStore { return Ok(()); } - let conn = self.pool.get()?; - let tx = conn.unchecked_transaction()?; + let mut conn = self.pool.get()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let mut seq: i64 = tx.query_row( "SELECT COALESCE(MAX(seq), 0) + 1 FROM messages WHERE session_id = ?1", @@ -736,8 +736,8 @@ impl SessionStore { summary_message: &ChatMessage, preserved_messages: &[ChatMessage], ) -> Result { - let conn = self.pool.get()?; - let tx = conn.unchecked_transaction()?; + let mut conn = self.pool.get()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let current_max_seq: i64 = tx.query_row( "SELECT COALESCE(MAX(seq), 0) FROM messages WHERE session_id = ?1", @@ -833,8 +833,8 @@ impl SessionStore { session_id: &str, messages: &[ChatMessage], ) -> Result<(), StorageError> { - let conn = self.pool.get()?; - let tx = conn.unchecked_transaction()?; + let mut conn = self.pool.get()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = current_timestamp(); // Delete all existing messages for this session @@ -892,8 +892,8 @@ impl SessionStore { topic_id: &str, messages: &[ChatMessage], ) -> Result<(), StorageError> { - let conn = self.pool.get()?; - let tx = conn.unchecked_transaction()?; + let mut conn = self.pool.get()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = current_timestamp(); // Delete only messages belonging to this topic — other topics' @@ -1023,8 +1023,8 @@ impl SessionStore { pub fn put_memory(&self, input: &MemoryUpsert) -> Result { let now = current_timestamp(); - let conn = self.pool.get()?; - let tx = conn.unchecked_transaction()?; + let mut conn = self.pool.get()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let existing: Option<(String, i64)> = tx .query_row( @@ -1611,10 +1611,11 @@ impl SessionStore { items: &[TodoRecord], ) -> Result, StorageError> { let mut conn = self.pool.get()?; - // 用 transaction()(非 unchecked_transaction)保证严格事务语义: + // 用 BEGIN IMMEDIATE 事务保证严格语义:写锁在事务开始时获取, + // 避免并发写事务在提交时死锁导致 "database is locked"。 // 用户数据替换需保证原子性——中途失败必须回滚,避免 DELETE 后 INSERT // 异常导致 todos 列表丢失且无法恢复。 - let tx = conn.transaction()?; + let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = current_timestamp(); // Delete existing todos for this scope_key