feat: Phase 5.2 - 实现 query 处理逻辑
- 创建可取消的 context 并存储 cancel func - 获取对话历史并传递给 orchestrator - 启动 orchestrator.ProcessQuery goroutine 处理查询 - 请求完成后自动清理 cancel func 和活跃请求标记
This commit is contained in:
@@ -170,19 +170,43 @@ func serveWS(c *gin.Context, sessionMgr session.Manager, orch orchestrator.Orche
|
|||||||
logger.Log.Infow("query received", "session", sessionID, "request", msg.RequestID)
|
logger.Log.Infow("query received", "session", sessionID, "request", msg.RequestID)
|
||||||
|
|
||||||
// 刷新会话 TTL
|
// 刷新会话 TTL
|
||||||
if err := sessionMgr.Touch(context.Background(), sessionID); err != nil {
|
if err := client.sessionMgr.Touch(context.Background(), sessionID); err != nil {
|
||||||
logger.Log.Warnw("touch session failed", "session", sessionID, "error", err)
|
logger.Log.Warnw("touch session failed", "session", sessionID, "error", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// 标记活跃请求
|
// 标记活跃请求
|
||||||
if err := sessionMgr.SetActiveRequest(context.Background(), sessionID, msg.RequestID); err != nil {
|
if err := client.sessionMgr.SetActiveRequest(context.Background(), sessionID, msg.RequestID); err != nil {
|
||||||
logger.Log.Warnw("set active request failed", "session", sessionID, "error", err)
|
logger.Log.Warnw("set active request failed", "session", sessionID, "error", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// 获取对话历史(供后续 Orchestrator 使用)
|
// 获取对话历史
|
||||||
_, _ = sessionMgr.GetHistory(context.Background(), sessionID, 20)
|
history, _ := client.sessionMgr.GetHistory(context.Background(), sessionID, 20)
|
||||||
|
|
||||||
// TODO: 解码 audio Base64 → 启动 orchestrator.ProcessQuery goroutine
|
// 创建可取消的 context
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
client.mu.Lock()
|
||||||
|
client.cancelFuncs[msg.RequestID] = cancel
|
||||||
|
client.mu.Unlock()
|
||||||
|
|
||||||
|
// 创建 sender
|
||||||
|
sender := &WSClient{client: client, requestID: msg.RequestID}
|
||||||
|
|
||||||
|
// 启动 orchestrator 处理 goroutine
|
||||||
|
go func() {
|
||||||
|
defer func() {
|
||||||
|
// 清理 cancel func
|
||||||
|
client.mu.Lock()
|
||||||
|
delete(client.cancelFuncs, msg.RequestID)
|
||||||
|
client.mu.Unlock()
|
||||||
|
cancel()
|
||||||
|
// 清除活跃请求
|
||||||
|
_ = client.sessionMgr.ClearActiveRequest(context.Background(), sessionID)
|
||||||
|
}()
|
||||||
|
|
||||||
|
if err := client.orchestrator.ProcessQuery(ctx, sessionID, msg, history, sender); err != nil {
|
||||||
|
logger.Log.Errorw("process query failed", "session", sessionID, "request", msg.RequestID, "error", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
case "config":
|
case "config":
|
||||||
var msg models.WsConfig
|
var msg models.WsConfig
|
||||||
|
|||||||
Reference in New Issue
Block a user