后端新增 load_messages_for_topic_page(seq < cursor + limit 分页,走 (session_id, seq) 索引替代 OFFSET 深翻页),ChatMessage 增加 seq 游标;历史批次消息带 topic_id 下发,前端以 seq+topic_id 双重判定批次归属,规避切话题瞬间在途旧批次污染。 前端触顶增量加载:批次缓存在 pendingHistoryRef,收到 topic_history_end 一次性去重 prepend,scrollTop 按新增高度补偿锚定原头部消息;不足一屏自动补页,loading 超时 10s 自愈;流式输出中加载历史不清空流式累加器,避免已流出文本丢失。
282 lines
10 KiB
Rust
282 lines
10 KiB
Rust
//! Row <-> record mapping helpers and single-record lookups by connection.
|
||
//!
|
||
//! These free functions operate on a borrowed [`rusqlite::Connection`] (or
|
||
//! [`rusqlite::Row`]) and have no access to the [`super::SessionStore`] pool.
|
||
//! They are extracted from `mod.rs` to keep the main module focused on
|
||
//! `SessionStore` methods and repository implementations.
|
||
|
||
use rusqlite::{Connection, OptionalExtension, params};
|
||
|
||
use crate::bus::ChatMessage;
|
||
use crate::bus::message::MessageUsage;
|
||
|
||
use super::{
|
||
MemoryRecord, SchedulerJobRecord, SchedulerJobState, SchedulerJobStatus, SessionRecord,
|
||
SessionTokenStats, SkillEventRecord, StorageError,
|
||
};
|
||
|
||
/// 消息加载查询的共享列清单(列序与 map_chat_message_row / map_usage_row 的下标一一对应)。
|
||
/// 新增 usage 列时只需改这里 + map_usage_row,无需逐条 SELECT 手工对齐。
|
||
pub(super) const MESSAGE_LOAD_COLUMNS: &str = "id, role, content, system_context, reasoning_content, media_refs_json, created_at, tool_call_id, tool_name, tool_calls_json, tool_duration_ms, prompt_tokens, completion_tokens, total_tokens, context_window_tokens, cached_tokens, seq";
|
||
|
||
/// 从指定列索引读取 token usage 五元组(含 context_window_tokens、cached_tokens)。
|
||
pub(super) fn map_usage_row(
|
||
row: &rusqlite::Row<'_>,
|
||
prompt_idx: usize,
|
||
completion_idx: usize,
|
||
total_idx: usize,
|
||
context_window_idx: usize,
|
||
cached_idx: usize,
|
||
) -> rusqlite::Result<Option<MessageUsage>> {
|
||
let prompt: Option<i64> = row.get(prompt_idx)?;
|
||
let completion: Option<i64> = row.get(completion_idx)?;
|
||
let total: Option<i64> = row.get(total_idx)?;
|
||
let context_window: Option<i64> = row.get(context_window_idx)?;
|
||
let cached: Option<i64> = row.get(cached_idx)?;
|
||
if prompt.is_none()
|
||
&& completion.is_none()
|
||
&& total.is_none()
|
||
&& context_window.is_none()
|
||
&& cached.is_none()
|
||
{
|
||
Ok(None)
|
||
} else {
|
||
Ok(Some(MessageUsage {
|
||
prompt_tokens: prompt.unwrap_or(0) as u32,
|
||
completion_tokens: completion.unwrap_or(0) as u32,
|
||
total_tokens: total.unwrap_or(0) as u32,
|
||
cached_tokens: cached.unwrap_or(0) as u32,
|
||
context_window_tokens: context_window.map(|v| v as u32),
|
||
}))
|
||
}
|
||
}
|
||
|
||
/// token 用量聚合的共享 SUM 列清单:batch_topic_token_stats 与
|
||
/// get_session_token_stats 共用。新增累计指标只需改这里 + read_usage_sum_row,
|
||
/// 避免两份聚合 SQL 手工对齐(shotgun surgery)。
|
||
pub(super) const USAGE_SUM_COLUMNS: &str = "COALESCE(SUM(prompt_tokens), 0), \
|
||
COALESCE(SUM(completion_tokens), 0), \
|
||
COALESCE(SUM(total_tokens), 0), \
|
||
COALESCE(SUM(cached_tokens), 0)";
|
||
|
||
/// 从聚合行读取累计 usage 字段(SUM 列从 offset 开始)。
|
||
/// last_* 瞬时字段不在 SUM 中,此处置 None,由调用方在 last 查询后回填。
|
||
pub(super) fn read_usage_sum_row(
|
||
row: &rusqlite::Row<'_>,
|
||
offset: usize,
|
||
) -> rusqlite::Result<SessionTokenStats> {
|
||
Ok(SessionTokenStats {
|
||
prompt_tokens: row.get::<_, i64>(offset)? as u64,
|
||
completion_tokens: row.get::<_, i64>(offset + 1)? as u64,
|
||
total_tokens: row.get::<_, i64>(offset + 2)? as u64,
|
||
cached_tokens: row.get::<_, i64>(offset + 3)? as u64,
|
||
last_prompt_tokens: None,
|
||
context_window_tokens: None,
|
||
})
|
||
}
|
||
|
||
pub(super) fn get_session_with_conn(
|
||
conn: &Connection,
|
||
session_id: &str,
|
||
) -> Result<Option<SessionRecord>, StorageError> {
|
||
let mut stmt = conn.prepare(
|
||
"
|
||
SELECT id, title, channel_name, chat_id, summary,
|
||
created_at, updated_at, last_active_at,
|
||
archived_at, deleted_at, message_count,
|
||
user_turn_count, agent_prompt_reinjection_count
|
||
FROM sessions
|
||
WHERE id = ?1 AND deleted_at IS NULL
|
||
",
|
||
)?;
|
||
|
||
stmt.query_row(params![session_id], map_session_record)
|
||
.optional()
|
||
.map_err(StorageError::from)
|
||
}
|
||
|
||
pub(super) fn get_memory_with_conn(
|
||
conn: &Connection,
|
||
scope_kind: &str,
|
||
scope_key: &str,
|
||
namespace: &str,
|
||
memory_key: &str,
|
||
) -> Result<Option<MemoryRecord>, StorageError> {
|
||
let mut stmt = conn.prepare(
|
||
"
|
||
SELECT id, scope_kind, scope_key, namespace, memory_key, content,
|
||
source_type, source_session_id, source_message_id, source_message_seq,
|
||
source_channel_name, source_chat_id, created_at, updated_at
|
||
FROM memories
|
||
WHERE scope_kind = ?1 AND scope_key = ?2 AND namespace = ?3 AND memory_key = ?4
|
||
",
|
||
)?;
|
||
|
||
stmt.query_row(
|
||
params![scope_kind, scope_key, namespace, memory_key],
|
||
map_memory_record,
|
||
)
|
||
.optional()
|
||
.map_err(StorageError::from)
|
||
}
|
||
|
||
pub(super) fn get_scheduler_job_with_conn(
|
||
conn: &Connection,
|
||
job_id: &str,
|
||
) -> Result<Option<SchedulerJobRecord>, StorageError> {
|
||
let mut stmt = conn.prepare(
|
||
"
|
||
SELECT id, kind, schedule_json, interval_secs, startup_delay_secs,
|
||
target_json, payload_json, enabled, state, last_status, last_error,
|
||
run_count, max_runs, last_fired_at, next_fire_at, paused_at, completed_at,
|
||
created_at, updated_at
|
||
FROM scheduler_jobs
|
||
WHERE id = ?1
|
||
",
|
||
)?;
|
||
|
||
stmt.query_row(params![job_id], map_scheduler_job_record)
|
||
.optional()
|
||
.map_err(StorageError::from)
|
||
}
|
||
|
||
pub(super) fn map_session_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<SessionRecord> {
|
||
Ok(SessionRecord {
|
||
id: row.get(0)?,
|
||
title: row.get(1)?,
|
||
channel_name: row.get(2)?,
|
||
chat_id: row.get(3)?,
|
||
summary: row.get(4)?,
|
||
created_at: row.get(5)?,
|
||
updated_at: row.get(6)?,
|
||
last_active_at: row.get(7)?,
|
||
archived_at: row.get(8)?,
|
||
deleted_at: row.get(9)?,
|
||
message_count: row.get(10)?,
|
||
user_turn_count: row.get(11)?,
|
||
agent_prompt_reinjection_count: row.get(12)?,
|
||
})
|
||
}
|
||
|
||
pub(super) fn map_skill_event_record(
|
||
row: &rusqlite::Row<'_>,
|
||
) -> rusqlite::Result<SkillEventRecord> {
|
||
let payload_json: String = row.get(4)?;
|
||
let payload = serde_json::from_str(&payload_json).map_err(|err| {
|
||
rusqlite::Error::FromSqlConversionFailure(4, rusqlite::types::Type::Text, Box::new(err))
|
||
})?;
|
||
|
||
Ok(SkillEventRecord {
|
||
id: row.get(0)?,
|
||
session_id: row.get(1)?,
|
||
event_type: row.get(2)?,
|
||
skill_name: row.get(3)?,
|
||
payload,
|
||
created_at: row.get(5)?,
|
||
})
|
||
}
|
||
|
||
pub(super) fn map_chat_message_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ChatMessage> {
|
||
let media_refs_json: String = row.get(5)?;
|
||
let media_refs: Vec<String> = serde_json::from_str(&media_refs_json).map_err(|err| {
|
||
rusqlite::Error::FromSqlConversionFailure(
|
||
media_refs_json.len(),
|
||
rusqlite::types::Type::Text,
|
||
Box::new(err),
|
||
)
|
||
})?;
|
||
|
||
let tool_calls_json: Option<String> = row.get(9)?;
|
||
let tool_calls = tool_calls_json
|
||
.as_deref()
|
||
.map(serde_json::from_str)
|
||
.transpose()
|
||
.map_err(|err| {
|
||
rusqlite::Error::FromSqlConversionFailure(9, rusqlite::types::Type::Text, Box::new(err))
|
||
})?;
|
||
|
||
Ok(ChatMessage {
|
||
id: row.get(0)?,
|
||
role: row.get(1)?,
|
||
content: row.get(2)?,
|
||
system_context: row.get(3)?,
|
||
reasoning_content: row.get(4)?,
|
||
media_refs,
|
||
timestamp: row.get(6)?,
|
||
tool_call_id: row.get(7)?,
|
||
tool_name: row.get(8)?,
|
||
tool_state: None,
|
||
tool_duration_ms: row.get::<_, Option<i64>>(10)?.map(|v| v as u64),
|
||
tool_calls,
|
||
usage: map_usage_row(row, 11, 12, 13, 14, 15)?,
|
||
seq: row.get(16)?,
|
||
})
|
||
}
|
||
|
||
pub(super) fn map_memory_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<MemoryRecord> {
|
||
Ok(MemoryRecord {
|
||
id: row.get(0)?,
|
||
scope_kind: row.get(1)?,
|
||
scope_key: row.get(2)?,
|
||
namespace: row.get(3)?,
|
||
memory_key: row.get(4)?,
|
||
content: row.get(5)?,
|
||
source_type: row.get(6)?,
|
||
source_session_id: row.get(7)?,
|
||
source_message_id: row.get(8)?,
|
||
source_message_seq: row.get(9)?,
|
||
source_channel_name: row.get(10)?,
|
||
source_chat_id: row.get(11)?,
|
||
created_at: row.get(12)?,
|
||
updated_at: row.get(13)?,
|
||
})
|
||
}
|
||
|
||
pub(super) fn map_scheduler_job_record(
|
||
row: &rusqlite::Row<'_>,
|
||
) -> rusqlite::Result<SchedulerJobRecord> {
|
||
let schedule_json: String = row.get(2)?;
|
||
let target_json: String = row.get(5)?;
|
||
let payload_json: String = row.get(6)?;
|
||
let state: String = row.get(8)?;
|
||
let last_status: Option<String> = row.get(9)?;
|
||
|
||
let schedule = serde_json::from_str(&schedule_json).map_err(|err| {
|
||
rusqlite::Error::FromSqlConversionFailure(2, rusqlite::types::Type::Text, Box::new(err))
|
||
})?;
|
||
let target = serde_json::from_str(&target_json).map_err(|err| {
|
||
rusqlite::Error::FromSqlConversionFailure(5, rusqlite::types::Type::Text, Box::new(err))
|
||
})?;
|
||
let payload = serde_json::from_str(&payload_json).map_err(|err| {
|
||
rusqlite::Error::FromSqlConversionFailure(6, rusqlite::types::Type::Text, Box::new(err))
|
||
})?;
|
||
|
||
Ok(SchedulerJobRecord {
|
||
id: row.get(0)?,
|
||
kind: row.get(1)?,
|
||
schedule,
|
||
interval_secs: row.get(3)?,
|
||
startup_delay_secs: row.get(4)?,
|
||
target,
|
||
payload,
|
||
enabled: row.get::<_, i64>(7)? != 0,
|
||
state: SchedulerJobState::from_str(&state).ok_or_else(|| {
|
||
rusqlite::Error::FromSqlConversionFailure(
|
||
8,
|
||
rusqlite::types::Type::Text,
|
||
format!("invalid scheduler job state: {}", state).into(),
|
||
)
|
||
})?,
|
||
last_status: last_status.and_then(|value| SchedulerJobStatus::from_str(&value)),
|
||
last_error: row.get(10)?,
|
||
run_count: row.get(11)?,
|
||
max_runs: row.get(12)?,
|
||
last_fired_at: row.get(13)?,
|
||
next_fire_at: row.get(14)?,
|
||
paused_at: row.get(15)?,
|
||
completed_at: row.get(16)?,
|
||
created_at: row.get(17)?,
|
||
updated_at: row.get(18)?,
|
||
})
|
||
}
|