fix(storage): 清理旧版本遗留的空 cli 会话并降级会话列表日志
- 早期版本每次 WebSocket 连接都会新建一个 cli 通道空会话(CLI Session xxxxxxxx),长期运行积累大量空壳记录,导致会话列表冗长、网关日志逐条刷屏 - 新增 user_version=3 一次性迁移:删除无消息/话题/待办/技能事件的 cli 空会话,有真实内容的会话不受影响(事务包裹,含回归测试) - ws.rs 会话列表逐条日志由 INFO 降级为 debug,仅保留条数汇总
This commit is contained in:
parent
c666c48f09
commit
b0c24d64f0
@ -238,9 +238,17 @@ async fn handle_socket(ws: WebSocket, state: Arc<GatewayState>) {
|
|||||||
websocket_sessions.sort_by_key(|s| -(s.last_active_at));
|
websocket_sessions.sort_by_key(|s| -(s.last_active_at));
|
||||||
}
|
}
|
||||||
|
|
||||||
tracing::info!("Sending {} sessions to client", websocket_sessions.len());
|
tracing::info!(
|
||||||
|
session_count = websocket_sessions.len(),
|
||||||
|
"Sending session list to client"
|
||||||
|
);
|
||||||
for s in &websocket_sessions {
|
for s in &websocket_sessions {
|
||||||
tracing::info!(" - {}: {} (channel: {})", s.id, s.title, s.channel_name);
|
tracing::debug!(
|
||||||
|
session_id = %s.id,
|
||||||
|
title = %s.title,
|
||||||
|
channel = %s.channel_name,
|
||||||
|
"Session list entry"
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
let session_summaries: Vec<crate::protocol::SessionSummary> = websocket_sessions
|
let session_summaries: Vec<crate::protocol::SessionSummary> = websocket_sessions
|
||||||
|
|||||||
@ -383,6 +383,46 @@ fn repair_session_id_prefix_pollution_inner(conn: &Connection) -> Result<(), Sto
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 清理历史遗留的空 cli 会话(user_version 3,一次性迁移)。
|
||||||
|
///
|
||||||
|
/// 早期版本每次 WebSocket 连接建立都会新建一个 cli 通道空会话
|
||||||
|
/// (自动标题 "CLI Session xxxxxxxx"),长期运行后 sessions 表积累大量
|
||||||
|
/// 空壳记录:它们会被并入发送给前端的会话列表,也在网关日志中逐条打印。
|
||||||
|
/// 本迁移删除没有任何消息、话题、待办、技能事件的 cli 会话(纯空壳),
|
||||||
|
/// 有真实内容的会话不受影响。事务包裹,通过 PRAGMA user_version 仅执行一次。
|
||||||
|
pub(super) fn cleanup_legacy_empty_cli_sessions(conn: &mut Connection) -> Result<(), StorageError> {
|
||||||
|
const EMPTY_CLI_CLEANUP_VERSION: i64 = 3;
|
||||||
|
|
||||||
|
let current_version: i64 = conn.query_row("PRAGMA user_version", [], |row| row.get(0))?;
|
||||||
|
if current_version >= EMPTY_CLI_CLEANUP_VERSION {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
let tx = conn.transaction()?;
|
||||||
|
let deleted = tx.execute(
|
||||||
|
"DELETE FROM sessions
|
||||||
|
WHERE channel_name = 'cli'
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM messages m WHERE m.session_id = sessions.id)
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM topics t WHERE t.session_id = sessions.id)
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM todos d WHERE d.session_id = sessions.id)
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM skill_events s WHERE s.session_id = sessions.id)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
if deleted > 0 {
|
||||||
|
tracing::info!(
|
||||||
|
deleted_count = deleted,
|
||||||
|
"Cleaned up legacy empty cli sessions"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
tx.execute(
|
||||||
|
&format!("PRAGMA user_version = {EMPTY_CLI_CLEANUP_VERSION}"),
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
tx.commit()?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
pub(super) fn ensure_todos_schema(conn: &Connection) -> Result<(), StorageError> {
|
pub(super) fn ensure_todos_schema(conn: &Connection) -> Result<(), StorageError> {
|
||||||
let table_exists: bool = conn
|
let table_exists: bool = conn
|
||||||
.query_row(
|
.query_row(
|
||||||
|
|||||||
@ -238,6 +238,7 @@ impl SessionStore {
|
|||||||
ensure_todos_schema(&conn)?;
|
ensure_todos_schema(&conn)?;
|
||||||
ensure_pending_subagents_schema(&conn)?;
|
ensure_pending_subagents_schema(&conn)?;
|
||||||
repair_session_id_prefix_pollution(&mut conn)?;
|
repair_session_id_prefix_pollution(&mut conn)?;
|
||||||
|
cleanup_legacy_empty_cli_sessions(&mut conn)?;
|
||||||
|
|
||||||
drop(conn);
|
drop(conn);
|
||||||
|
|
||||||
|
|||||||
@ -794,3 +794,33 @@ fn test_repair_session_id_prefix_pollution() {
|
|||||||
assert_eq!(store.get_session("abc").unwrap().unwrap().message_count, 2);
|
assert_eq!(store.get_session("abc").unwrap().unwrap().message_count, 2);
|
||||||
assert_eq!(store.list_sessions("websocket", false).unwrap().len(), 2);
|
assert_eq!(store.list_sessions("websocket", false).unwrap().len(), 2);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_cleanup_legacy_empty_cli_sessions() {
|
||||||
|
let store = SessionStore::in_memory().unwrap();
|
||||||
|
|
||||||
|
// 空 cli 会话(旧版本遗留的空壳):应被删除
|
||||||
|
let empty = store.create_cli_session(None).unwrap();
|
||||||
|
// 有消息的 cli 会话:必须保留
|
||||||
|
let with_data = store.create_cli_session(None).unwrap();
|
||||||
|
store
|
||||||
|
.append_message(&with_data.id, &ChatMessage::user("hello"))
|
||||||
|
.unwrap();
|
||||||
|
// 空 websocket 会话:不在 cli 清理范围内
|
||||||
|
let ws = store.ensure_channel_session("websocket", "chat-1").unwrap();
|
||||||
|
|
||||||
|
let mut conn = store.pool.get().unwrap();
|
||||||
|
// from_connection 已把 user_version 推到 3,回退到 2 模拟"尚未清理"
|
||||||
|
conn.execute("PRAGMA user_version = 2", []).unwrap();
|
||||||
|
|
||||||
|
super::migrations::cleanup_legacy_empty_cli_sessions(&mut conn).unwrap();
|
||||||
|
|
||||||
|
assert!(store.get_session(&empty.id).unwrap().is_none());
|
||||||
|
assert!(store.get_session(&with_data.id).unwrap().is_some());
|
||||||
|
assert!(store.get_session(&ws.id).unwrap().is_some());
|
||||||
|
|
||||||
|
// 幂等:版本守卫使第二次调用直接跳过,已有会话不受影响
|
||||||
|
super::migrations::cleanup_legacy_empty_cli_sessions(&mut conn).unwrap();
|
||||||
|
assert!(store.get_session(&with_data.id).unwrap().is_some());
|
||||||
|
assert!(store.get_session(&ws.id).unwrap().is_some());
|
||||||
|
}
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user