Zum Inhalt springen
Deutsch

SSE (Server-Sent Events)

SseActor verbindet sich mit einem Server-Sent-Events-Endpoint und konsumiert ihn. Es ist ein Nur-Lese-Client, kein Server: Er öffnet ein langlebiges HTTP-GET, parst das text/event-stream-Wire-Format und leitet jedes geparste Event an einen Ziel-Actor weiter. Es gibt keine Kommandos und keinen ausgehenden Pfad — SSE ist einseitig, vom Server zum Client.

import { ActorSystem, Actor } from 'actor-ts';
import { SseActor, type SseEvent } from 'actor-ts/io';
import { SseOptions } from 'actor-ts/io';
// A collector receives each parsed event.
class Collector extends Actor<SseEvent> {
override onReceive(ev: SseEvent): void {
// ev.event — the `event:` field ('message' by default)
// ev.data — the `data:` payload (multiline joined with '\n')
// ev.id — the `id:` field, when the server sent one
console.log(ev.event, ev.data, ev.id);
}
}
const target = sys.spawnAnonymous(Collector);
const sseOptions = SseOptions.create()
.withUrl('https://example.com/events')
.withTarget(target)
.withHeaders({ authorization: 'Bearer …' });
sys.spawnAnonymous(() => new SseActor(
sseOptions,
));

Der Actor nutzt das globale fetch (Bun, Node ≥ 18 und Deno bringen alle eines mit), es gibt also keine zusätzliche Abhängigkeit zu installieren.

interface SseOptionsType extends BrokerCommonOptionsType {
url?: string; // required — the SSE endpoint
headers?: Readonly<Record<string, string>>; // optional custom request headers
target?: ActorRef<SseEvent>; // required — subscriber for inbound events
idleTimeoutMs?: number; // read-idle deadline; default 0 (off)
connectTimeoutMs?: number; // connect deadline; default 0 (off)
}
SettingErforderlichBeschreibung
urljaDer SSE-Endpoint, mit dem der Actor sich verbindet (GET, Accept: text/event-stream).
targetjaDer Actor, der jedes geparste SseEvent empfängt.
headersneinZusätzliche Request-Header — z. B. authorization für einen geschützten Feed.
idleTimeoutMsneinDen Stream als verloren melden, wenn so lange kein einziges Byte eintrifft.
connectTimeoutMsneinEinen Connect abbrechen, der nicht rechtzeitig Response-Header liefert.

requiredOptions ist ['url', 'target']; beide müssen vorhanden sein, sonst schlägt der Actor beim Start sofort fehl.

Jedes geparste Event wird als SseEvent zugestellt:

type SseEvent = {
event: string; // the `event:` field, or 'message' by default
data: string; // the `data:` payload
id?: string; // the `id:` field, when present
};

Der Client parst das SSE-Wire-Format — Events durch eine Leerzeile (\n\n) getrennt, Felder als field: value:

data: hello
event: tick
data: {"n":1}
id: 100
data: line-1
data: line-2

Parsing-Regeln, die der SSE-Spezifikation folgen:

  • Default-Event-Name. Ein Block ohne event:-Feld kommt mit event: 'message' an.
  • Mehrzeilige Daten. Mehrere data:-Zeilen in einem Block werden mit '\n' verbunden — der letzte Block oben ergibt data: 'line-1\nline-2'.
  • Kommentare ignoriert. Zeilen, die mit : beginnen (Keepalive-Kommentare), werden übersprungen.
  • id. Als ev.id zugestellt, wenn der Server eine sendet; sonst undefined.

Der obige Stream stellt somit drei Events zu: { event: 'message', data: 'hello' }, ein tick-Event mit id '100' und das mehrzeilige { event: 'message', data: 'line-1\nline-2' }.

SseActor erweitert BrokerActor, daher werden Reconnect-mit-Backoff und der Circuit Breaker geerbt — konfiguriere sie am Builder:

SseOptions.create()
.withUrl('https://example.com/events')
.withTarget(target)
.withReconnect({ /* base delay, max delay, jitter … */ })
.withCircuitBreaker(5, 30_000); // failureThreshold, resetMs

Wenn der Stream endet oder die Verbindung beobachtbar abbricht, greift die Reconnect-Maschinerie der Basisklasse. Übergib .withReconnect(false), um nach dem Stream-Ende zu stoppen statt erneut zu versuchen — praktisch für endliche Feeds oder Tests:

SseOptions.create()
.withUrl(url)
.withTarget(target)
.withReconnect(false); // stop when the stream closes

„Beobachtbar“ leistet in diesem Satz echte Arbeit. consume parkt auf reader.read(), und ein Server, der verschwindet, ohne die Response zu schließen — ein beendeter Container, ein verfallener NAT-Eintrag, eine verschluckende Route — liefert kein done, keinen Fehler und keine Bytes. SSE ist read-only, also kann nichts sonst im Actor das von einem Feed unterscheiden, der gerade nichts zu melden hat: der Stream bleibt connected, solange der Prozess läuft, und das Monitoring ist blind.

idleTimeoutMs ist die Uhr, die dem ein Ende setzt:

const sseOptions = SseOptions.create()
.withUrl('https://example.com/events')
.withTarget(target)
.withIdleTimeoutMs(60_000) // no bytes for 60 s → reconnect
.withConnectTimeoutMs(10_000); // no response headers in 10 s → fail the attempt

Jeder Chunk setzt sie zurück, auch ein Kommentar — : ping\n\n parst zu gar keinem Event, und genau so hält ein wohlerzogener Server einen leerlaufenden Feed offen. Setze idleTimeoutMs daher deutlich über das Keepalive-Intervall des Servers (15–30 s ist die übliche Konvention); darunter kappt die Deadline gesunde Streams in Endlosschleife. Standardmäßig aus, weil nur du dieses Intervall kennst.

connectTimeoutMs ist die andere Hälfte: fetch hat keine eigene Deadline, also hält ein Server, der die TCP-Verbindung annimmt und den Request nie beantwortet, den Actor unbegrenzt in connecting.

Beide melden den Verlust über den gewöhnlichen handleConnectionLost-Pfad, sodass Backoff, Circuit Breaker und BrokerDisconnected sich exakt so verhalten wie bei einem beendeten Stream — siehe BrokerActor-Basis.

Settings lösen mit der üblichen Präzedenz auf — explizite Options > HOCON > eingebaute Defaults. Der HOCON-Pfad ist actor-ts.io.broker.sse:

actor-ts.io.broker.sse {
url = "https://example.com/events"
headers { authorization = "Bearer …" }
idleTimeoutMs = 60000
connectTimeoutMs = 10000
}

target ist eine Actor-Referenz und kann daher nur im Code gesetzt werden, nicht in HOCON.

SseActor passt zum Konsumieren eines externen SSE-Feeds — überall dort, wo ein Server einen einseitigen Stream über einfaches HTTP pusht und du jedes Event ins Actor-System geroutet haben möchtest:

  1. LLM-Token-Streams — einen streamenden Completion-Endpoint konsumieren und Tokens an einen Collector weiterleiten.
  2. Marktdaten / Kurs-Ticks — ein als Events gepushter Live-Feed.
  3. CI-/Build-Logs — ein langlaufender Log-Stream, der beim Entstehen mitgelesen wird.

Weil es einfach HTTP ist, funktioniert es durch Proxies, Load-Balancer und CDNs, die kein WebSocket verstehen.