评论实时更新:SSE 订阅文章频道,发评/编辑/删除/审核通过即时刷新打开中的页面

- internal/hub:按文章 ID 的进程内发布/订阅,信号不携带内容(未审核评论不外泄)
- GET /api/comments/stream?post_id=N:SSE 长连 + 25s 心跳 + X-Accel-Buffering off
- 广播点:发评论、作者编辑/删除、后台审核通过、后台删除(与公开端共享同一 hub)
- 前端 EventSource 静默刷新(不亮加载态、不干扰正在输入/编辑的状态)
- requestLog 的 statusRecorder 补 Flush 透传——包装器没实现 Flusher 时流式端点整体不可用
This commit is contained in:
Sakurasan
2026-09-28 21:42:11 +08:00
parent 661d157b6c
commit a234afe448
8 changed files with 168 additions and 2 deletions
+14
View File
@@ -22,6 +22,7 @@ import (
"oneblog/internal/config"
"oneblog/internal/httpx"
"oneblog/internal/hub"
"oneblog/internal/model"
"oneblog/internal/storage"
"oneblog/internal/store"
@@ -32,6 +33,9 @@ type API struct {
Store *store.Store
Cfg *config.Config
Sessions *Sessions
// Hub 是评论变更广播(与公开 API 共享同一实例,main.go 装配);
// 审核通过 / 后台删除评论时让前台打开着的页面即时刷新。
Hub *hub.Hub
// 文件上传的存储后端与公开域名(main.go 装配,两个 API 共享同一实例)
Blobs storage.BlobStore
PublicBase string
@@ -849,6 +853,11 @@ func (a *API) adminCommentByID(w http.ResponseWriter, r *http.Request) {
httpx.ServerError(w, err)
return
}
if in.Status == "visible" && a.Hub != nil {
if c, err := a.Store.GetComment(id); err == nil {
a.Hub.Broadcast(c.PostID)
}
}
httpx.OK(w, map[string]any{"ok": true})
case http.MethodDelete:
// 后台删除同样走软删(墓碑保楼层)
@@ -860,6 +869,11 @@ func (a *API) adminCommentByID(w http.ResponseWriter, r *http.Request) {
httpx.ServerError(w, err)
return
}
if a.Hub != nil {
if c, err := a.Store.GetComment(id); err == nil {
a.Hub.Broadcast(c.PostID)
}
}
httpx.OK(w, map[string]any{"ok": true})
default:
httpx.Error(w, http.StatusMethodNotAllowed, "PUT/DELETE required")
+4
View File
@@ -15,6 +15,7 @@ import (
"oneblog/internal/auth"
"oneblog/internal/config"
"oneblog/internal/httpx"
"oneblog/internal/hub"
"oneblog/internal/model"
"oneblog/internal/ratelimit"
"oneblog/internal/storage"
@@ -40,6 +41,8 @@ type API struct {
// 其余登录方式(main.go 装配,未配置的自动不开放)
GG auth.Google
TG auth.Telegram
// Hub 是评论变更的进程内广播(SSE 用;与后台 admin 共享同一实例)
Hub *hub.Hub
// 限流(Routes 里惰性初始化):读者登录失败按 IP 计、评论写入按读者计。
// 登录入口此前裸奔——OAuth 跳转本身难刷,但 state 校验失败、
@@ -78,6 +81,7 @@ func (a *API) Routes() http.Handler {
mux.HandleFunc("/api/auth/callback/google", a.googleCallback)
mux.HandleFunc("/api/auth/telegram", a.telegramAuth)
mux.HandleFunc("/api/comments", a.comments)
mux.HandleFunc("/api/comments/stream", a.commentsStream)
mux.HandleFunc("/api/comments/", a.commentSub)
mux.HandleFunc("/api/site", a.site)
mux.HandleFunc("/api/posts", a.listPosts)
+3
View File
@@ -285,6 +285,7 @@ func (a *API) createComment(w http.ResponseWriter, r *http.Request) {
httpx.ServerError(w, err)
return
}
a.Hub.Broadcast(in.PostID)
if rateKey != "" {
a.commentNew.Add(rateKey) // 只计成功写入:空正文这类手滑不扣配额
}
@@ -390,6 +391,7 @@ func (a *API) editComment(w http.ResponseWriter, r *http.Request, id int64) {
httpx.ServerError(w, err)
return
}
a.Hub.Broadcast(updated.PostID)
httpx.OK(w, updated)
}
@@ -417,6 +419,7 @@ func (a *API) deleteComment(w http.ResponseWriter, r *http.Request, id int64) {
httpx.ServerError(w, err)
return
}
a.Hub.Broadcast(c.PostID)
httpx.OK(w, map[string]any{"ok": true})
}
+52
View File
@@ -0,0 +1,52 @@
// 评论变更的 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()
}
}
}
+49
View File
@@ -0,0 +1,49 @@
// Package hub 是一个极简的进程内发布/订阅:按主题(文章 ID)广播
// 「有变更」信号。订阅者只拿到一个空 struct 信号——不携带内容,
// 消费方(SSE 前端)收到后自己 refetch,未审核的评论内容不会经由
// 这条通道外泄。
package hub
import "sync"
type Hub struct {
mu sync.Mutex
subs map[int64]map[chan struct{}]struct{}
}
func New() *Hub {
return &Hub{subs: make(map[int64]map[chan struct{}]struct{})}
}
// Subscribe 订阅某主题;返回信号 channel 和退订函数。
// channel 容量 1:订阅者处理不过来时新信号直接丢弃(合并刷新)。
func (h *Hub) Subscribe(topic int64) (<-chan struct{}, func()) {
ch := make(chan struct{}, 1)
h.mu.Lock()
if h.subs[topic] == nil {
h.subs[topic] = make(map[chan struct{}]struct{})
}
h.subs[topic][ch] = struct{}{}
h.mu.Unlock()
off := func() {
h.mu.Lock()
delete(h.subs[topic], ch)
if len(h.subs[topic]) == 0 {
delete(h.subs, topic)
}
h.mu.Unlock()
}
return ch, off
}
// Broadcast 唤醒某主题的全部订阅者;积压的订阅者不阻塞。
func (h *Hub) Broadcast(topic int64) {
h.mu.Lock()
defer h.mu.Unlock()
for ch := range h.subs[topic] {
select {
case ch <- struct{}{}:
default:
}
}
}