// 评论变更的 SSE 端点:GET /api/comments/stream?post_id=N。 // 事件里只有「这篇文章的评论变了」——不带任何内容,前端收到后 // 自己 refetch。这样未审核的评论内容不会经 SSE 泄给围观者。 package api import ( "fmt" "net/http" "strconv" "time" "oneblog/internal/httpx" ) // commentsStream 长连推送。EventSource 断线会自动重连,所以这里 // 只需要:连上即发注释行、25s 心跳防代理掐线、事件触发刷新。 func (a *API) commentsStream(w http.ResponseWriter, r *http.Request) { postID, err := strconv.ParseInt(r.URL.Query().Get("post_id"), 10, 64) if err != nil || postID <= 0 { httpx.BadRequest(w, "post_id required") return } fl, ok := w.(http.Flusher) if !ok { httpx.Error(w, http.StatusInternalServerError, "streaming unsupported") return } w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache") // 反向代理(nginx 等)默认缓冲会吞掉流式响应 w.Header().Set("X-Accel-Buffering", "no") fmt.Fprint(w, ": connected\n\n") fl.Flush() ch, off := a.Hub.Subscribe(postID) defer off() heartbeat := time.NewTicker(25 * time.Second) defer heartbeat.Stop() for { select { case <-r.Context().Done(): return case <-heartbeat.C: fmt.Fprint(w, ": ping\n\n") fl.Flush() case <-ch: fmt.Fprintf(w, "event: comments\ndata: {\"post_id\":%d}\n\n", postID) fl.Flush() } } }