Files
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

70 lines
2.6 KiB
TypeScript

import type {UserEventMessage, WorkerInboundMessage} from '../types.ts';
// Minimal SharedWorker/MessagePort doubles: worker.ts wires user events onto
// `sharedWorker.port`, so we capture the port to feed it messages and assert dispatch behavior.
type PortListener = (ev: {data: WorkerInboundMessage}) => void;
class MockMessagePort {
listeners: Record<string, PortListener[]> = {};
posted: unknown[] = [];
addEventListener(type: string, cb: PortListener) {
(this.listeners[type] ||= []).push(cb);
}
removeEventListener() {}
postMessage(msg: unknown) { this.posted.push(msg) }
start() {}
close() {}
// Simulate the underlying worker delivering a message to the page.
deliver(msg: UserEventMessage) {
for (const cb of this.listeners['message'] || []) cb({data: {msgType: 'user-event', msgData: msg}});
}
}
let lastWorker: MockSharedWorker;
class MockSharedWorker {
port = new MockMessagePort();
// eslint-disable-next-line unicorn/no-this-assignment
constructor() { lastWorker = this }
addEventListener() {}
}
// worker.ts caches module-scope state (subscribers, initialized), so re-import
// a fresh module per test after stubbing the globals it reads on init.
async function freshWorker() {
vi.resetModules();
vi.stubGlobal('WebSocket', class {});
vi.stubGlobal('SharedWorker', MockSharedWorker);
return await import('./worker.ts');
}
afterEach(() => {
vi.unstubAllGlobals();
});
// sequential: freshWorker resets the module registry and stubs globals, which is unsafe
// to interleave with the other test under the repo's `sequence.concurrent` vitest config.
// Every push follows a real DB write, so a repeated value still means the
// underlying list changed (e.g. a new comment on an already-unread issue).
test('dispatches every push, including repeated values', {concurrent: false}, async () => {
const {onUserEvent} = await freshWorker();
const received: number[] = [];
onUserEvent('notification-count', (msg) => { received.push(msg.eventData.count) });
lastWorker.port.deliver({eventType: 'notification-count', eventData: {count: 1}});
lastWorker.port.deliver({eventType: 'notification-count', eventData: {count: 1}});
expect(received).toEqual([1, 1]);
});
test('worker-connected flags the page and reaches its subscribers', {concurrent: false}, async () => {
const {onUserEvent} = await freshWorker();
let connects = 0;
onUserEvent('worker-connected', () => { connects++ });
lastWorker.port.deliver({eventType: 'worker-connected'});
expect(connects).toBe(1);
expect(document.documentElement.getAttribute('data-user-events-connected')).toBe('true');
});