消息级溯源 P1(溯源字段升级):
- ContextManager 加 last_assistant/last_user_message_id 取消息 id
- audit 6 处 message_id 接真值(替换 P0 None 占位)
- 知识提炼 source_ref 升级 conv_msg:{id}(知识生命线 extracted/referenced 双事件)
- 4 处 build_knowledge_context 调用方传 user message id
- 老数据 id=None 降级 conv:{id} 兼容
流式失败重试修复(UX-260619-06):
- 症状1(回答一半消失):AiError 收尾前 flushCurrentText 保留部分回复到占位气泡
(原仅 AiAgentRound/AiCompleted 回填,AiError 清 currentText 致部分回复丢失)
- 症状2(重试报"没有可重新生成的回复"):ai_regenerate pop_last_assistant_round 返 false 时
末尾是 user(流式中途失败这轮 assistant 未持久化)直接从该 user 重跑(等价重发)
1214 连续 user 已由 merge_consecutive_users(7c98134)修,本批补失败重试链路
自验: workspace EXIT 0 + cargo test 全绿 + vue-tsc EXIT 0
147 lines
5.2 KiB
Rust
147 lines
5.2 KiB
Rust
//! 知识生命线记录器 — 统一入口,各业务流程只调一行
|
|
//!
|
|
//! 设计目标:
|
|
//! - 事件记录逻辑不散落在各业务流程内部(提取/注入/状态变更),集中于此模块
|
|
//! - fire-and-forget 友好:便捷方法内部吞掉错误(只 warn),调用方无需处理 Result
|
|
//! - 易扩展:新事件类型只加一个便捷方法 + 一行 context 构造,不改主流程
|
|
|
|
use std::sync::Arc;
|
|
|
|
use df_types::types::new_id;
|
|
use df_storage::crud::KnowledgeEventsRepo;
|
|
use df_storage::db::Database;
|
|
use df_storage::models::KnowledgeEventRecord;
|
|
|
|
use super::{err_str, now_millis};
|
|
|
|
/// 事件类型常量(生命线节点;归档复用 status_changed 的 to=archived,不单列)
|
|
pub const EVENT_CREATED: &str = "created";
|
|
pub const EVENT_EXTRACTED: &str = "extracted";
|
|
pub const EVENT_STATUS_CHANGED: &str = "status_changed";
|
|
pub const EVENT_REFERENCED: &str = "referenced";
|
|
|
|
/// 知识生命线记录器 — 持有 DB 句柄,提供便捷事件写入方法
|
|
///
|
|
/// 用法:`KnowledgeTimeline::new(&state.db).record_xxx(...).await`
|
|
///
|
|
/// 异步策略(按路径热度,非 bug,刻意不一致):
|
|
/// - 热路径(对话注入 build_knowledge_context / 提炼 extract):整块已在 spawn 任务内,
|
|
/// 此处直接 await = 对用户 fire-and-forget,不阻塞 AI 对话
|
|
/// - 低频命令(状态变更 / 创建):事件写入 ms 级,直接 await 比 detached spawn 更简单
|
|
/// 便捷方法内部 `fire()` 已吞错误(warn 不阻断),调用方无需处理 Result。
|
|
pub struct KnowledgeTimeline {
|
|
db: Arc<Database>,
|
|
}
|
|
|
|
impl KnowledgeTimeline {
|
|
pub fn new(db: &Arc<Database>) -> Self {
|
|
Self { db: db.clone() }
|
|
}
|
|
|
|
/// 通用写入(底层入口,便捷方法均经此)
|
|
pub async fn record(
|
|
&self,
|
|
knowledge_id: &str,
|
|
event_type: &str,
|
|
source_ref: Option<String>,
|
|
context_json: Option<String>,
|
|
) -> Result<(), String> {
|
|
let record = KnowledgeEventRecord {
|
|
id: new_id(),
|
|
knowledge_id: knowledge_id.to_string(),
|
|
event_type: event_type.to_string(),
|
|
source_ref,
|
|
context_json,
|
|
timestamp: now_millis(),
|
|
};
|
|
KnowledgeEventsRepo::new(&self.db)
|
|
.insert(record)
|
|
.await
|
|
.map_err(err_str)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// 手动录入产生(context: method)
|
|
pub async fn record_created(&self, knowledge_id: &str, method: &str) {
|
|
let ctx = serde_json::json!({ "method": method }).to_string();
|
|
self.fire(knowledge_id, EVENT_CREATED, Some("manual".into()), Some(ctx)).await;
|
|
}
|
|
|
|
/// AI 对话提炼产生(context: conv_id / conv_title / reasoning)
|
|
///
|
|
/// F-260619-04 P1 消息级溯源:`message_id` 为末条 assistant 消息 id(本轮 AI 产出知识的载体)。
|
|
/// - Some(id) → source_ref = `conv_msg:{id}`(消息级)
|
|
/// - None(老数据无 id) → source_ref = `conv:{conv_id}`(对话级,向前兼容)
|
|
pub async fn record_extracted(
|
|
&self,
|
|
knowledge_id: &str,
|
|
conv_id: &str,
|
|
conv_title: &str,
|
|
reasoning: &str,
|
|
message_id: Option<&str>,
|
|
) {
|
|
let ctx = serde_json::json!({
|
|
"conv_id": conv_id,
|
|
"conv_title": conv_title,
|
|
"reasoning": reasoning,
|
|
})
|
|
.to_string();
|
|
let source_ref = match message_id {
|
|
Some(mid) => format!("conv_msg:{}", mid),
|
|
None => format!("conv:{}", conv_id),
|
|
};
|
|
self.fire(
|
|
knowledge_id,
|
|
EVENT_EXTRACTED,
|
|
Some(source_ref),
|
|
Some(ctx),
|
|
)
|
|
.await;
|
|
}
|
|
|
|
/// 被检索命中/注入(context: conv_id / query)
|
|
///
|
|
/// F-260619-04 P1 消息级溯源:`message_id` 为触发检索的 user 消息 id(用户问题命中知识)。
|
|
/// - Some(id) → source_ref = `conv_msg:{id}`(消息级)
|
|
/// - None(老数据无 id) → source_ref = `conv:{conv_id}`(对话级,向前兼容)
|
|
pub async fn record_referenced(
|
|
&self,
|
|
knowledge_id: &str,
|
|
conv_id: &str,
|
|
query: &str,
|
|
message_id: Option<&str>,
|
|
) {
|
|
let ctx = serde_json::json!({ "conv_id": conv_id, "query": query }).to_string();
|
|
let source_ref = match message_id {
|
|
Some(mid) => format!("conv_msg:{}", mid),
|
|
None => format!("conv:{}", conv_id),
|
|
};
|
|
self.fire(
|
|
knowledge_id,
|
|
EVENT_REFERENCED,
|
|
Some(source_ref),
|
|
Some(ctx),
|
|
)
|
|
.await;
|
|
}
|
|
|
|
/// 状态变更(context: from / to;to=published=审核通过,to=archived=归档)
|
|
pub async fn record_status_change(&self, knowledge_id: &str, from: &str, to: &str) {
|
|
let ctx = serde_json::json!({ "from": from, "to": to }).to_string();
|
|
self.fire(knowledge_id, EVENT_STATUS_CHANGED, None, Some(ctx)).await;
|
|
}
|
|
|
|
/// 内部:写入并吞错误(fire-and-forget,失败只 warn 不阻断主流程)
|
|
async fn fire(
|
|
&self,
|
|
knowledge_id: &str,
|
|
event_type: &str,
|
|
source_ref: Option<String>,
|
|
context_json: Option<String>,
|
|
) {
|
|
if let Err(e) = self.record(knowledge_id, event_type, source_ref, context_json).await {
|
|
tracing::warn!("知识生命线事件写入失败(非阻断) [{}]: {}", event_type, e);
|
|
}
|
|
}
|
|
}
|