Files
gitea/routers/web/websocket/websocket.go
mohammad rahimi 13d0f24423 feat: Replace SSE with WebSocket for UI notifications (#36965)
* Closes #36942
* Fixes #19265 

Replaces the SSE-based push channel (`/user/events`) with a WebSocket
endpoint (`/-/ws`).

### What changes

- **New `/-/ws` endpoint** (authenticated). One WebSocket per origin,
shared across tabs via a single `SharedWorker`.
- **Pubsub broker** (`services/pubsub`) for fan-out by topic, behind a
`Broker` interface. `MemoryBroker` is the default (single process); a
Redis backend is available for multi-process setups, configured via
`[websocket].PUBSUB_TYPE` / `PUBSUB_CONN_STR`. The internal Gitea queue
was not usable here because it has FIFO/single-consumer semantics.
- **Push-only event production.** Events are emitted by write-triggered
notifiers — `NotificationCountChange`, `PublishStopwatchesForUser`, and
the logout publisher — wired into the existing `notify.Notifier`
interface. No server-side pollers.
- **Typed pub/sub on the client.** `web_src/js/modules/worker.ts` is a
singleton transport; features subscribe per event type via
`onUserEvent('notification-count', cb)` instead of branching on
`event.data.type`.
- **Wire contract** (`UserEventType` union) is shared between the worker
and consumers via `web_src/js/types.ts`, kept in sync with
`services/websocket/events.go`.
- **Client-side periodic polling fallback** kicks in only when the
WebSocket cannot be established (e.g. proxy blocks WS, browser lacks
module-SharedWorker support).

### What's removed

- `modules/eventsource` (SSE manager, run loop, messenger).
- `/user/events` route and `tests/integration/eventsource_test.go`.
- All server-side polling for stopwatches and notification counts.

### Stopwatch multi-tab fix

The navbar stopwatch icon was previously rendered conditionally on `{{if
$activeStopwatch}}`, so tabs loaded before the timer started had no DOM
element to update. The icon and popup are now always rendered (toggled
with `tw-hidden`), and the start/stop/cancel handlers POST silently so
all open tabs reflect the change in real time.

### Deployment note

WebSocket needs the upgrade headers to pass through a reverse proxy,
e.g. for nginx:

```nginx
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
```

Without them the WebSocket cannot be established, and after 3
consecutive failed opens the shared worker signals `push-unavailable`:
the notification count and stopwatch fall back to periodic polling on
the existing `[ui.notification]` timeouts. Real-time push is lost, the
features keep working. The reverse-proxy docs need the same note (see
the `docs-update-needed` label).

---------

Co-authored-by: silverwind <me@silverwind.io>
Co-authored-by: wxiaoguang <wxiaoguang@gmail.com>
Co-authored-by: Epid <rexmrj@gmail.com>
2026-07-27 07:46:00 +00:00

118 lines
3.5 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// Copyright 2026 The Gitea Authors. All rights reserved.
// SPDX-License-Identifier: MIT
package websocket
import (
gocontext "context"
"net/http"
"time"
"gitea.dev/modules/graceful"
"gitea.dev/modules/json"
"gitea.dev/modules/log"
"gitea.dev/services/context"
"gitea.dev/services/pubsub"
websocket_service "gitea.dev/services/websocket"
gitea_ws "github.com/coder/websocket"
)
const (
pingInterval = 30 * time.Second
pingTimeout = 10 * time.Second
writeTimeout = 10 * time.Second
// First code in the IANA library/framework reserved range (30003999).
// Sentinel for an unauthenticated session so the SharedWorker can tell
// "your cookie is gone" apart from a transient network failure and stop
// reconnecting in a tight loop.
closeCodeUnauthenticated gitea_ws.StatusCode = 3000
)
// filterLogout forwards a session-free logout only to the targeted connection
// (its own session, or every session when SessionID is empty) and drops it for
// the rest. Non-logout messages pass through untouched.
func filterLogout(eventType string, eventDataBytes []byte, connSessionID string) []byte {
if eventType != websocket_service.EventLogout {
return eventDataBytes
}
var lm websocket_service.UserEventMessage[websocket_service.LogoutEventData]
if err := json.Unmarshal(eventDataBytes, &lm); err != nil {
return eventDataBytes
}
if lm.EventData.SessionID == "" || lm.EventData.SessionID == connSessionID {
return []byte(`{"eventType":"logout"}`)
}
return nil
}
func Serve(ctx *context.Context) {
// Answer plain GETs (health checks, crawlers) here; letting Accept reject them
// would log an error per request. Same reply it would have sent.
if ctx.Req.Header.Get("Upgrade") == "" {
ctx.Resp.Header().Set("Connection", "Upgrade")
ctx.Resp.Header().Set("Upgrade", "websocket")
ctx.Resp.WriteHeader(http.StatusUpgradeRequired)
return
}
conn, err := gitea_ws.Accept(ctx.Resp, ctx.Req, nil)
if err != nil {
log.Error("websocket: accept failed: %v", err)
return
}
defer conn.CloseNow() //nolint:errcheck // best-effort close
if !ctx.IsSigned {
_ = conn.Close(closeCodeUnauthenticated, "unauthenticated")
return
}
sessionID := ctx.Session.ID()
ch, cancel := pubsub.DefaultBroker.Subscribe(pubsub.UserTopic(ctx.Doer.ID))
defer cancel()
// Ping requires a concurrent reader to observe the pong frame; CloseRead
// spawns one and cancels its context when the peer goes away.
wsCtx := conn.CloseRead(ctx.Req.Context())
shutdownCtx := graceful.GetManager().ShutdownContext()
pingTicker := time.NewTicker(pingInterval)
defer pingTicker.Stop()
for {
select {
case <-wsCtx.Done():
return
case <-shutdownCtx.Done():
return
case <-pingTicker.C:
pingCtx, cancelPing := gocontext.WithTimeout(wsCtx, pingTimeout)
err := conn.Ping(pingCtx)
cancelPing()
if err != nil {
log.Trace("websocket: ping failed: %v", err)
return
}
case brokerPayload, ok := <-ch:
if !ok {
return
}
eventType, eventDataBytes := websocket_service.ExtractUserEventMessage(brokerPayload)
eventDataBytes = filterLogout(eventType, eventDataBytes, sessionID)
if eventDataBytes == nil {
continue
}
// Bound the write so a stalled/slow peer can't block this goroutine
// indefinitely and starve the ping ticker.
writeCtx, cancelWrite := gocontext.WithTimeout(wsCtx, writeTimeout)
err := conn.Write(writeCtx, gitea_ws.MessageText, eventDataBytes)
cancelWrite()
if err != nil {
log.Trace("websocket: write failed: %v", err)
return
}
}
}
}