mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-08 00:25:05 +00:00
feat: SSE Last-Event-ID 断线重放 + GUI 表单防刷新
- webui handler: 新增 sseEventRing 环状缓冲区(200条),断线重连按 Last-Event-ID 重放遗漏事件 - GUI app.js: doRenderAll 检测连接表单打开时改走 refreshDataOnly,修复 15s 定时器擦掉用户输入的 bug - 附带 dashboard.html/index.html 前端调整 + handler_sse_test.go 单测
This commit is contained in:
@ -66,6 +66,51 @@ func init() {
|
||||
}
|
||||
}
|
||||
|
||||
// sseEventRecord 保存一条 SSE 事件元数据,供断线重连时按 Last-Event-ID 重放遗漏事件。
|
||||
type sseEventRecord struct {
|
||||
id string // SSE 事件 id 值(如 "1234567890-5")
|
||||
eventType string // 事件类型(agent_output, reasoning 等)
|
||||
data json.RawMessage // 序列化后的 payload JSON
|
||||
}
|
||||
|
||||
// sseEventRing 是一个固定大小的环状缓冲区,保持最近 cap 条 SSE 事件。
|
||||
type sseEventRing struct {
|
||||
mu sync.Mutex
|
||||
buf []sseEventRecord
|
||||
cap int
|
||||
}
|
||||
|
||||
func newSSEEventRing(cap int) *sseEventRing {
|
||||
return &sseEventRing{cap: cap}
|
||||
}
|
||||
|
||||
// Append 追加一条事件,超过容量时丢弃最旧条目。
|
||||
func (r *sseEventRing) Append(id, eventType string, data json.RawMessage) {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.buf = append(r.buf, sseEventRecord{id: id, eventType: eventType, data: data})
|
||||
if len(r.buf) > r.cap {
|
||||
r.buf = r.buf[len(r.buf)-r.cap:]
|
||||
}
|
||||
}
|
||||
|
||||
// After 返回所有在指定 id 之后的事件(按写入顺序),若 id 不在缓冲区中则返回全部。
|
||||
func (r *sseEventRing) After(id string) []sseEventRecord {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
for i := len(r.buf) - 1; i >= 0; i-- {
|
||||
if r.buf[i].id == id {
|
||||
result := make([]sseEventRecord, len(r.buf)-i-1)
|
||||
copy(result, r.buf[i+1:])
|
||||
return result
|
||||
}
|
||||
}
|
||||
// ID 不在缓冲区(可能是太旧或从未收到),返回全部
|
||||
result := make([]sseEventRecord, len(r.buf))
|
||||
copy(result, r.buf)
|
||||
return result
|
||||
}
|
||||
|
||||
type Handler struct {
|
||||
sdk *sdk.PluginSDK
|
||||
supervisor sdk.SupervisorAPI
|
||||
@ -85,6 +130,8 @@ type Handler struct {
|
||||
sessionMu sync.Mutex
|
||||
sessions map[string]time.Time
|
||||
|
||||
sseEvents *sseEventRing // SSE 事件环状缓冲区,Last-Event-ID 重放用
|
||||
|
||||
chatMu sync.Mutex
|
||||
chatHistory []ChatMsg
|
||||
pendingIdx int // chatHistory 中正在进行的 assistant 消息索引,-1 表示无
|
||||
@ -214,6 +261,7 @@ func NewHandler(s *sdk.PluginSDK) *Handler {
|
||||
termStates: make(map[string]*termState),
|
||||
pendingIdx: -1,
|
||||
chatMsgCache: make(map[string]*chatMsgEntry),
|
||||
sseEvents: newSSEEventRing(200),
|
||||
}
|
||||
h.loadChatHistory()
|
||||
if s != nil {
|
||||
@ -1344,10 +1392,24 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
log.Printf("[SSE] handler started, subscribing to events")
|
||||
|
||||
// 解析 Last-Event-ID(断线重连时客户端携带)
|
||||
lastEventID := r.Header.Get("Last-Event-ID")
|
||||
if lastEventID != "" {
|
||||
// 解析 Last-Event-ID(断线重连时客户端携带),重放期间内遗漏的事件。
|
||||
// 注意:本 handler 的 Last-Event-ID 重放仅为 GUI (cmd/gui/renderer/app.js) 服务。
|
||||
// 浏览器原生 EventSource (webui/dashboard.html 使用) 由浏览器自动处理 Last-Event-ID 重连。
|
||||
if lastEventID := r.Header.Get("Last-Event-ID"); lastEventID != "" {
|
||||
log.Printf("[SSE] client reported Last-Event-ID: %s", lastEventID)
|
||||
if h.sseEvents != nil {
|
||||
replayed := h.sseEvents.After(lastEventID)
|
||||
if len(replayed) == 0 {
|
||||
log.Printf("[SSE] replay: nothing after id %s (id not in ring or already at tip)", lastEventID)
|
||||
} else {
|
||||
log.Printf("[SSE] replay: sending %d events after id %s", len(replayed), lastEventID)
|
||||
for _, rec := range replayed {
|
||||
fmt.Fprintf(w, "id: %s\nevent: %s\ndata: %s\n", rec.id, rec.eventType, string(rec.data))
|
||||
flusher.Flush()
|
||||
}
|
||||
log.Printf("[SSE] replay complete, wrote %d events", len(replayed))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
subTypes := []string{"agent_output", "reasoning", "agent_error", "tool_call", "stage", "agent_llm_chain", "terminal_output"}
|
||||
@ -1362,8 +1424,13 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
data, _ := json.Marshal(evt)
|
||||
seq++
|
||||
id := fmt.Sprintf("%d-%d", evt.Timestamp, seq)
|
||||
// 写入环状缓冲区,供断线重连重放
|
||||
if h.sseEvents != nil {
|
||||
h.sseEvents.Append(id, string(evt.Type), data)
|
||||
}
|
||||
select {
|
||||
case writeCh <- fmt.Sprintf("id: %d-%d\nevent: %s\ndata: %s\n", evt.Timestamp, seq, evt.Type, string(data)):
|
||||
case writeCh <- fmt.Sprintf("id: %s\nevent: %s\ndata: %s\n", id, evt.Type, string(data)):
|
||||
if evt.Type == sdk.EventToolCall {
|
||||
toolName, _ := evt.Payload["tool"].(string)
|
||||
log.Printf("[SSE] wrote tool_call to writeCh: tool=%s", toolName)
|
||||
|
||||
Reference in New Issue
Block a user