Files
gitea/web_src/js/modules/worker.ts
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

131 lines
4.7 KiB
TypeScript

import type {
SharedWorkerControlMessage,
UserEventMessage,
UserEventType,
WorkerEventMessage,
WorkerInboundMessage,
} from '../types.ts';
const {appSubUrl, sharedWorkerUri} = window.config;
type EventOf<T extends UserEventType> = Extract<UserEventMessage, {eventType: T}>;
type Subscriber<T extends UserEventType = UserEventType> = (msg: EventOf<T>) => void;
const subscribers = new Map<UserEventType, Set<Subscriber>>();
let fallbackSignalled = false;
let sharedWorker: SharedWorker | null = null;
function dispatch(msg: UserEventMessage) {
if (msg.eventType === 'worker-connected') {
document.documentElement.setAttribute('data-user-events-connected', 'true');
}
const set = subscribers.get(msg.eventType);
if (!set) return;
for (const cb of set) cb(msg);
}
function signalFallback() {
if (fallbackSignalled) return;
fallbackSignalled = true;
dispatch({eventType: 'worker-unavailable'});
}
function init() {
try {
sharedWorker = new SharedWorker(sharedWorkerUri, {type: 'module', name: 'user-events'});
} catch (err) {
console.warn('SharedWorker unavailable, falling back to periodic polling', err);
queueMicrotask(signalFallback);
return;
}
// Browsers without module-SharedWorker support fail at parse time, before the WebSocket opens.
sharedWorker.addEventListener('error', (event) => {
console.error('worker error', event);
signalFallback();
});
sharedWorker.port.addEventListener('messageerror', () => {
console.error('unable to deserialize message');
});
sharedWorker.port.addEventListener('error', (e) => {
console.error('worker port error', e);
});
const postSharedWorkerControlMessage = (sw: SharedWorker, msg: SharedWorkerControlMessage) => {
sw.port.postMessage(msg);
};
const handleWorkerEvent = (msg: WorkerEventMessage) => {
if (msg.workerEvent === 'error') {
console.error('worker port event error', msg);
} else if (msg.workerEvent === 'close') {
postSharedWorkerControlMessage(sharedWorker!, {type: 'close'});
sharedWorker!.port.close();
} else {
console.error('unknown worker port event', msg);
}
};
const handleLogout = () => {
postSharedWorkerControlMessage(sharedWorker!, {type: 'close'});
sharedWorker!.port.close();
// slightly delay our "logout" for a short while, in case there are other logout requests in-flight.
// * if the logout is triggered by a page redirection (e.g.: user clicks "/user/logout")
// * "beforeunload" event is triggered, this code path won't execute
// * if the logout is triggered by a fetch call
// * "beforeunload" event is not triggered until JS does the redirection.
// * in this case, the logout fetch call already completes and has sent the "logout" message to the worker
// * there can be a data-race between the fetch call's redirection and the "logout" message from the worker
// * the fetch call's logout redirection should always win over the worker message, because it might have a custom location
setTimeout(() => { window.location.assign(`${appSubUrl}/`) }, 1000);
};
sharedWorker.port.addEventListener('message', (event: MessageEvent<WorkerInboundMessage>) => {
const msg = event.data;
if (msg?.msgType === 'worker-event') {
handleWorkerEvent(msg.msgData);
} else if (msg?.msgType === 'user-event') {
if (msg.msgData?.eventType === 'logout') {
handleLogout();
return;
}
dispatch(msg.msgData);
} else {
console.error('unknown inbound message', msg);
}
});
sharedWorker.port.start();
const wsProtocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:';
postSharedWorkerControlMessage(sharedWorker, {
type: 'start',
url: `${wsProtocol}//${window.location.host}${appSubUrl}/-/ws`,
showDebugLog: !window.config.runModeIsProd,
});
window.addEventListener('beforeunload', () => {
// FIXME: this logic is not quite right.
// "beforeunload" can be canceled by some actions like "are-you-sure" and the navigation can be cancelled.
// In this case: the worker port is incorrectly closed while the page is still there.
postSharedWorkerControlMessage(sharedWorker!, {type: 'close'});
sharedWorker!.port.close();
});
}
let initialized = false;
export function onUserEvent<T extends UserEventType>(type: T, cb: Subscriber<T>): () => void {
if (!initialized) {
initialized = true;
if (window.WebSocket && window.SharedWorker) {
init();
} else {
queueMicrotask(signalFallback);
}
}
let set = subscribers.get(type);
if (!set) {
set = new Set();
subscribers.set(type, set);
}
const wrapped: Subscriber = (msg) => cb(msg as EventOf<T>);
set.add(wrapped);
return () => { set.delete(wrapped) };
}