From a234afe4486721621755c490724d7e14f906154d Mon Sep 17 00:00:00 2001 From: Sakurasan <26715255+Sakurasan@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:42:11 +0800 Subject: [PATCH] =?UTF-8?q?=E8=AF=84=E8=AE=BA=E5=AE=9E=E6=97=B6=E6=9B=B4?= =?UTF-8?q?=E6=96=B0=EF=BC=9ASSE=20=E8=AE=A2=E9=98=85=E6=96=87=E7=AB=A0?= =?UTF-8?q?=E9=A2=91=E9=81=93=EF=BC=8C=E5=8F=91=E8=AF=84/=E7=BC=96?= =?UTF-8?q?=E8=BE=91/=E5=88=A0=E9=99=A4/=E5=AE=A1=E6=A0=B8=E9=80=9A?= =?UTF-8?q?=E8=BF=87=E5=8D=B3=E6=97=B6=E5=88=B7=E6=96=B0=E6=89=93=E5=BC=80?= =?UTF-8?q?=E4=B8=AD=E7=9A=84=E9=A1=B5=E9=9D=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - internal/hub:按文章 ID 的进程内发布/订阅,信号不携带内容(未审核评论不外泄) - GET /api/comments/stream?post_id=N:SSE 长连 + 25s 心跳 + X-Accel-Buffering off - 广播点:发评论、作者编辑/删除、后台审核通过、后台删除(与公开端共享同一 hub) - 前端 EventSource 静默刷新(不亮加载态、不干扰正在输入/编辑的状态) - requestLog 的 statusRecorder 补 Flush 透传——包装器没实现 Flusher 时流式端点整体不可用 --- backend/internal/admin/api.go | 14 +++++ backend/internal/api/api.go | 4 ++ backend/internal/api/comments.go | 3 ++ backend/internal/api/sse.go | 52 +++++++++++++++++++ backend/internal/hub/hub.go | 49 +++++++++++++++++ backend/main.go | 13 +++++ .../components/comments/CommentSection.vue | 16 +++++- frontend/src/composables/useComments.js | 19 +++++++ 8 files changed, 168 insertions(+), 2 deletions(-) create mode 100644 backend/internal/api/sse.go create mode 100644 backend/internal/hub/hub.go diff --git a/backend/internal/admin/api.go b/backend/internal/admin/api.go index dd8a8e9..595258a 100644 --- a/backend/internal/admin/api.go +++ b/backend/internal/admin/api.go @@ -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") diff --git a/backend/internal/api/api.go b/backend/internal/api/api.go index 3be969b..d19440b 100644 --- a/backend/internal/api/api.go +++ b/backend/internal/api/api.go @@ -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) diff --git a/backend/internal/api/comments.go b/backend/internal/api/comments.go index 9e467cc..e3d320d 100644 --- a/backend/internal/api/comments.go +++ b/backend/internal/api/comments.go @@ -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}) } diff --git a/backend/internal/api/sse.go b/backend/internal/api/sse.go new file mode 100644 index 0000000..c9bda5f --- /dev/null +++ b/backend/internal/api/sse.go @@ -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() + } + } +} diff --git a/backend/internal/hub/hub.go b/backend/internal/hub/hub.go new file mode 100644 index 0000000..0b1b3a2 --- /dev/null +++ b/backend/internal/hub/hub.go @@ -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: + } + } +} diff --git a/backend/main.go b/backend/main.go index 00ed714..5de1ef4 100644 --- a/backend/main.go +++ b/backend/main.go @@ -21,6 +21,7 @@ import ( "oneblog/internal/auth" "oneblog/internal/config" "oneblog/internal/db" + "oneblog/internal/hub" "oneblog/internal/storage" "oneblog/internal/store" "oneblog/internal/thumbs" @@ -66,6 +67,8 @@ func main() { readerSessions := auth.NewReaderSessions(cfg.SessionSec, 30*24*time.Hour) // 后台管理员会话(cookie one_session),前台评论也用它识别站主身份 adminSessions := admin.NewSessions(cfg.SessionSec, 7*24*time.Hour) + // 评论变更广播(SSE):公开端订阅、双端发布,共用同一实例 + commentHub := hub.New() // 缩略图懒生成 + 磁盘缓存(首次请求生成一次,之后直接供缓存) thumbCache := thumbs.NewStore(filepath.Join(cfg.DataDir, ".thumbnail_cache")) @@ -79,8 +82,10 @@ func main() { GG: auth.Google{ClientID: cfg.GoogleClientID, ClientSecret: cfg.GoogleClientSecret}, TG: auth.Telegram{Bot: cfg.TelegramBot, Token: cfg.TelegramToken}, AdminSessions: adminSessions, + Hub: commentHub, } adminAPI := admin.NewAPI(st, cfg, adminSessions) + adminAPI.Hub = commentHub adminAPI.Blobs = blobs adminAPI.Thumbs = thumbCache // 删文件时连带清掉它的缩略图 @@ -187,6 +192,14 @@ type statusRecorder struct { func (s *statusRecorder) WriteHeader(code int) { s.status = code; s.ResponseWriter.WriteHeader(code) } +// Flush 透传给底层 writer——SSE 这类流式响应靠它判定可刷新, +// 包装器不实现 Flusher 的话流式端点会直接不可用。 +func (s *statusRecorder) Flush() { + if f, ok := s.ResponseWriter.(http.Flusher); ok { + f.Flush() + } +} + func placeholderPage() string { return `