Wrapping EventSource in an RxJS Observable Permalink to this section
Part of Angular & Svelte SSE Integration, under Frontend Consumption & Client Patterns.
An EventSource already behaves like an Observable: it pushes values over time, it can fail, and it should stop when nobody is listening. Wrapping it in RxJS makes that explicit and lets the stream compose with the rest of an Angular (or any RxJS-based) application — filtering by event type, folding into state, combining with HTTP results. The wrapper is short, but each line encodes a decision: when the connection opens, when it closes, how many connections exist, and what happens on errors. This guide writes it carefully.
Symptom & Developer Intent Permalink to this section
- Every
asyncpipe orsubscribe()opens a new stream; DevTools shows five identical requests. - The connection stays open after leaving the page.
- A single malformed event kills the whole Observable chain, and nothing reconnects.
retry()added to the pipe causes a reconnect storm, fightingEventSource’s own reconnection.- Components re-render far more often than the event rate would suggest.
The intent is one connection per stream, opened lazily and closed when unused, typed events, resilience to bad data, a single owner of reconnection logic, and minimal change detection.
Root Cause Analysis Permalink to this section
A new Observable(...) is cold: its subscribe function runs once per subscriber. Wrapping new EventSource(...) in it without sharing means one connection per subscriber. Sharing operators turn it hot, with a reference count that runs teardown when the last subscriber unsubscribes.
The retry storm comes from two reconnection mechanisms competing. EventSource reconnects by itself after network errors and does not signal an error to the Observable in that case — it fires error events while readyState is CONNECTING. Only when the connection fails permanently (readyState === CLOSED) should the Observable error, and only then does an RxJS retry make sense.
Step-by-Step Resolution Permalink to this section
Step 1 — Write the cold Observable with teardown Permalink to this section
import { Observable } from 'rxjs';
export interface SseEvent<T> { type: string; id: string; data: T; }
export function fromEventSource<T>(url: string, types: string[], init?: EventSourceInit): Observable<SseEvent<T>> {
return new Observable<SseEvent<T>>((subscriber) => {
const es = new EventSource(url, init);
const onMessage = (e: MessageEvent) => {
try {
subscriber.next({ type: e.type, id: e.lastEventId, data: JSON.parse(e.data) as T });
} catch {
// A bad frame is skipped, not fatal: report it without ending the stream.
console.warn('unparseable SSE event', e.type);
}
};
types.forEach((t) => es.addEventListener(t, onMessage as EventListener));
es.onerror = () => {
if (es.readyState === EventSource.CLOSED) {
subscriber.error(new Error('SSE connection closed')); // permanent: let RxJS decide
}
// readyState CONNECTING: the browser is already retrying; do nothing.
};
return () => es.close(); // teardown on unsubscribe
});
}
Step 2 — Share it with a reference count Permalink to this section
import { share, timer } from 'rxjs';
const orders$ = fromEventSource<OrderEvent>('/api/stream', ['order.updated', 'order.deleted'], { withCredentials: true })
.pipe(share({ resetOnRefCountZero: () => timer(5000) })); // keep open 5 s after last unsubscribe
resetOnRefCountZero with a short timer avoids tearing down and reconnecting when a component is destroyed and recreated during navigation. For state-shaped streams where late subscribers need the latest value, use shareReplay({ bufferSize: 1, refCount: true }).
Step 3 — Retry only permanent failures, with backoff Permalink to this section
import { retry, timer } from 'rxjs';
const resilient$ = orders$.pipe(
retry({
delay: (_err, attempt) => timer(Math.min(30_000, 1000 * 2 ** attempt) * (1 + Math.random() * 0.3)),
resetOnSuccess: true,
}),
);
Because the Observable only errors when EventSource has given up (a non-200 response, for instance), this retry never competes with the browser’s own reconnection. For 401 responses, which EventSource cannot distinguish, check authentication state in the retry delay function and refresh the session before resubscribing.
Step 4 — Derive typed streams per event type Permalink to this section
const updates$ = resilient$.pipe(filter((e): e is SseEvent<OrderUpdated> => e.type === 'order.updated'));
const openCount$ = resilient$.pipe(
scan((acc, e) => applyOrderEvent(acc, e), initialOrders),
map((s) => s.openCount),
distinctUntilChanged(), // only emit when the displayed value changes
);
distinctUntilChanged on derived values is a cheap and effective way to avoid rendering when an event does not change what the template shows.
Step 5 — Provide the wrapper through dependency injection Permalink to this section
In Angular, put the wrapper behind an injectable service and an injection token for the EventSource constructor. Production code injects the real browser class; tests inject a fake; server-side rendering injects nothing and receives EMPTY. Keeping a map of shared Observables inside the service, keyed by URL and event types, guarantees that two components asking for the same stream get the same instance, even if they were written by different teams and never coordinate:
export const EVENT_SOURCE = new InjectionToken<typeof EventSource>('EventSource', {
providedIn: 'root',
factory: () => (typeof EventSource === 'function' ? EventSource : (null as any)),
});
The service then calls new (inject(EVENT_SOURCE))(url, init) inside the Observable’s subscribe function.
Step 6 — Keep change detection quiet in Angular Permalink to this section
Create the EventSource outside Angular’s zone and re-enter only for emissions, or rely on signals in zoneless applications:
const zoned$ = new Observable<SseEvent<T>>((sub) => {
const inner = this.zone.runOutsideAngular(() => fromEventSource<T>(url, types).subscribe({
next: (v) => this.zone.run(() => sub.next(v)),
error: (e) => this.zone.run(() => sub.error(e)),
}));
return () => inner.unsubscribe();
});
Combine with bufferTime(0, animationFrameScheduler) for high-rate streams so that a burst becomes one emission per frame.
Validation & Monitoring Permalink to this section
it('opens one connection for many subscribers and closes after the last', fakeAsync(() => {
const subs = [orders$.subscribe(), orders$.subscribe(), orders$.subscribe()];
expect(FakeEventSource.instances.length).toBe(1);
subs.forEach((s) => s.unsubscribe());
tick(5000);
expect(FakeEventSource.instances[0].closed).toBe(true);
}));
In the browser, the Network panel should show one stream request per stream URL. Record permanent failures and retries as telemetry, with the reason if known.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Is there an RxJS operator for EventSource built in?
No. fromEvent can listen to EventSource events, but it does not create or close the connection. A small custom Observable like the one here owns the connection's lifecycle.
Why not error the Observable on every EventSource error event?
EventSource fires error while it is reconnecting after transient drops. Erroring then would end the stream and trigger RxJS retries on top of the browser's own reconnection.
share or shareReplay?
share for event-shaped data where late subscribers only need new events; shareReplay with bufferSize 1 and refCount true for state-shaped data where a late subscriber needs the current value.
How do I send a bearer token?
Replace EventSource inside the wrapper with a fetch-based client that sets the Authorization header. The Observable's interface and the rest of the app are unchanged.