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