gRPC
Das Framework stellt zwei Actor-Klassen für gRPC bereit:
| Klasse | Rolle |
|---|---|
GrpcServerActor | Hostet einen gRPC-Server; stellt Service-Methoden bereit. |
GrpcClientActor | Verbindet sich mit einem Remote-gRPC-Service; ruft Methoden auf. |
Beide umschließen @grpc/grpc-js. Nützlich, wenn du
Protobuf-typisierte Verträge hast und den Broker-Actor-
Lifecycle (Reconnect, Buffer, Subscriber-Fan-out) für clientseitige
Aufrufe willst.
Beide Actors werden über einen fluenten Options-Builder
konfiguriert — GrpcServerOptions.create()… /
GrpcClientOptions.create()… — und wie jeder andere Actor
gespawnt.
Der Server ist ein konkreter Actor: spawne ihn direkt. Er lädt das
Proto, registriert die Methoden-Handler und bindet in
preStart — der spawn-Aufruf kehrt zurück, bevor der Bind
abgeschlossen ist, also gib ihm vor dem ersten Client-Call einen
Moment.
Jede Methode wird auf einen Handler-Actor abgebildet, der den
eingehenden Call als Nachricht empfängt. Handler sind in der
handlers-Map nach Methodennamen indiziert; jede Methode, die
nicht in der Map steht, wird als UNIMPLEMENTED gemeldet.
import { Actor,} from 'actor-ts';import { GrpcServerActor, GrpcServerOptions, type GrpcUnaryCall, type GrpcServerStreamCall,} from 'actor-ts/io';
// A unary handler: reply once via `respond`.class GetSensorHandler extends Actor<GrpcUnaryCall> { override onReceive(call: GrpcUnaryCall): void { const id = (call.request as { id: string }).id; call.respond({ id, label: `sensor-${id}` }); }}
// A server-stream handler: emit chunks via `send`, finish via `complete`.class WatchSensorHandler extends Actor<GrpcServerStreamCall> { override onReceive(call: GrpcServerStreamCall): void { const limit = (call.request as { limit?: number }).limit ?? 5; for (let i = 0; i < limit; i++) { call.send({ value: 20 + i, ts: Date.now() }); } call.complete(); }}
const getHandler = system.spawn(GetSensorHandler, 'get');const watchHandler = system.spawn(WatchSensorHandler, 'watch');
const grpcServerOptions = GrpcServerOptions.create() .withProtoPath(protoPath) .withPackageName('sensor.v1') .withServiceName('SensorService') .withBind('127.0.0.1:50051') .withHandlers({ GetSensor: { kind: 'unary', target: getHandler }, WatchSensor: { kind: 'serverStream', target: watchHandler }, });const server = system.spawn( () => new GrpcServerActor( grpcServerOptions, ), 'grpc-server',);Die Bind-Adresse ist ein einzelner host:port-String über
withBind; packageName und serviceName werden separat gesetzt.
Jeder Handler-Deskriptor ist { kind, target }, wobei kind
entweder 'unary', 'serverStream', 'clientStream' oder 'bidi'
ist.
Call-Formen der Handler
Abschnitt betitelt „Call-Formen der Handler“Das Framework deserialisiert das Protobuf-Payload und übergibt dem Handler ein typisiertes Call-Objekt:
Handler-kind | Call-Typ | Reply-API |
|---|---|---|
unary | GrpcUnaryCall | respond(res) / respondError(message, code?) |
serverStream | GrpcServerStreamCall | send(chunk) / complete() / fail(message, code?) |
clientStream | GrpcClientStreamCall | onData(target) / respond(res) / respondError(message, code?) |
bidi | GrpcBidiCall | onData(target) / send(chunk) / complete() |
Jeder Call trägt method, den deserialisierten request (außer
den Client-Stream- und Bidi-Calls, die ihre Requests streamen) und
ein metadata-Record. respondError / fail nehmen einen
optionalen numerischen Status-Code, der auf 13 (INTERNAL) fällt.
metadata enthält die Request-Header des Clients, ein Handler kann
also pro Call entscheiden — einen authorization-Header prüfen, eine
Tenant-Id lesen. Zwei Einschränkungen folgen daraus, dass der Record
Strings hält: ein mehrfach gesendeter Header fällt auf seinen
ersten Wert zusammen, und binäre (-bin) Header entfallen
vollständig. Ein nicht gesendeter Header liefert undefined, auch bei
Namen wie constructor — der Record hat keinen Prototyp, ein Lookup
gelingt also nie versehentlich.
Ein onData-Subscriber erhält GrpcRequestStreamInbound — entweder
{ kind: 'chunk', chunk } oder { kind: 'end' }.
import { match, P } from 'ts-pattern';import { Actor,} from 'actor-ts';import { GrpcClientActor, GrpcClientOptions, type GrpcInbound, type ReplyMessage, type RpcErrorMessage, type StreamDataMessage, type StreamErrorMessage, type StreamStartedMessage,} from 'actor-ts/io';
const grpcClientOptions = GrpcClientOptions.create() .withProtoPath(protoPath) .withPackageName('sensor.v1') .withServiceName('SensorService') .withEndpoint('127.0.0.1:50051');const client = system.spawn( () => new GrpcClientActor( grpcClientOptions, ), 'sensor-client',);
// A collector actor receives every reply / stream frame as `GrpcInbound`.class ReplyCollector extends Actor<GrpcInbound> { override onReceive(message: GrpcInbound): void { match(message) .with({ kind: 'reply' }, (m) => this.onReply(m)) .with({ kind: 'stream-started' }, (m) => this.onStreamStarted(m)) .with({ kind: 'stream-data' }, (m) => this.onStreamData(m)) .with({ kind: 'stream-end' }, () => this.onStreamEnd()) // Unary and streaming failures are reported the same way. .with(P.union({ kind: 'rpc-error' }, { kind: 'stream-error' }), (m) => this.onError(m)) .exhaustive(); }
private onReply(message: ReplyMessage): void { console.log('unary reply:', message.response); }
// A client stream is open — `message.handle` addresses its writes. private onStreamStarted(message: StreamStartedMessage): void { console.log('client stream open:', message.handle.streamId); }
private onStreamData(message: StreamDataMessage): void { console.log('stream chunk:', message.chunk); }
private onStreamEnd(): void { console.log('stream complete'); }
private onError(message: RpcErrorMessage | StreamErrorMessage): void { console.error('error:', message.error.message); }}
const collector = system.spawn(ReplyCollector, 'collector');
// Make a unary call — the reply is delivered to `target`.client.tell({ kind: 'unary', method: 'GetSensor', request: { id: 'rt-7' }, target: collector });Der Endpoint ist ein einzelner host:port-String über
withEndpoint. Jeder Call benennt einen target-Actor; der
Actor routet die Antwort und alle Stream-Frames dorthin.
.withDeadlineMs(30000) (der Standardwert) begrenzt ausschließlich
Unary-Calls. Ein gRPC-Deadline umfasst den gesamten RPC, und ein
einzelner Wert kann nicht zugleich einen Request/Response-Call zügig
scheitern lassen und einen langlebigen Stream laufen lassen — die drei
Streaming-Modi bleiben deshalb unbegrenzt.
Eingehende Frames
Abschnitt betitelt „Eingehende Frames“Antworten und Stream-Frames treffen beim target-Actor als
GrpcInbound-Discriminated-Union ein:
kind | Felder | Bedeutung |
|---|---|---|
reply | response | Unärer **oder Client-Stream-**Abschluss. |
stream-started | handle | Ein Client-Stream ist offen — siehe unten. |
stream-data | streamId, chunk | Ein Stream-Chunk. |
stream-end | streamId | Stream sauber geschlossen. |
stream-error | streamId, error | Stream fehlgeschlagen. |
rpc-error | error | Unärer Call / Call-Setup fehlgeschlagen. |
Streaming-Modi
Abschnitt betitelt „Streaming-Modi“gRPC hat vier Call-Typen, und das Framework deckt alle vier ab:
| Typ | Client sendet | Server liefert | Clientseitige kinds |
|---|---|---|---|
| Unär | Einen Request | Eine Response | unary |
| Server-Streaming | Einen Request | Stream von Responses | serverStream |
| Client-Streaming | Stream von Requests | Eine Response | clientStreamStart / clientStreamSend / clientStreamClose |
| Bidirektionales Streaming | Stream von Requests | Stream von Responses | bidiStart / bidiSend / bidiClose |
Jeder clientseitige Call ist ein Kommando, das du dem Client-Actor
per tell schickst; das kind wählt den Call-Typ.
// Server-streaming: one request, N chunks routed to `target`.client.tell({ kind: 'serverStream', method: 'WatchSensor', request: { id: 'rt-7', limit: 5 }, target: collector,});
// Client-streaming: open a stream, then push requests into it.client.tell({ kind: 'clientStreamStart', method: 'ReportReadings', target: collector });Client-Streaming und das Stream-Handle
Abschnitt betitelt „Client-Streaming und das Stream-Handle“clientStreamStart liefert nichts direkt zurück. Der Actor stellt
dem target einen stream-started-Frame zu, der ein
GrpcStreamHandle trägt — und dieses Handle adressiert jeden
weiteren Schreibvorgang:
// Inside the collector, on the `stream-started` frame:const handle = message.handle;
client.tell({ kind: 'clientStreamSend', handle, chunk: { value: 21.5 } });client.tell({ kind: 'clientStreamClose', handle });clientStreamClose schließt den Request-Stream halbseitig; die
einzelne Server-Response trifft danach als ganz normaler reply
ein, ein Fehlschlag als rpc-error.
Das Handle hat zwei Felder mit zwei verschiedenen Aufgaben.
streamId ist die Korrelations-ID — dieselbe Zahl, die die Frames
dieses Streams tragen, sodass ein Collector mehrere gleichzeitige
Streams multiplexen kann. token ist die Capability: ein tell
trägt keinen verifizierten Absender, eine fortlaufende ID würde
also jedem, der den Client-Actor erreicht, das Schreiben in einen
nie selbst geöffneten Stream erlauben. Der Token besteht aus 64 Bit
kryptografischer Zufälligkeit — damit ist der Map-Lookup selbst die
Ownership-Prüfung. Gib das Handle genauso sorgsam weiter, wie du es
mit einem Schreib-Handle tätest.
Serverseitig wird eine Client-Streaming-Methode mit
kind: 'clientStream' registriert; ihr Handler konsumiert den
Request-Stream über onData und antwortet einmal:
class ReportReadingsHandler extends Actor<GrpcClientStreamCall> { override onReceive(call: GrpcClientStreamCall): void { let count = 0; const sink = { tell: (m: GrpcRequestStreamInbound): void => { match(m) .with({ kind: 'chunk' }, () => { count++; }) .with({ kind: 'end' }, () => call.respond({ count })) .exhaustive(); } } as unknown as ActorRef<GrpcRequestStreamInbound>; call.onData(sink); }}Chunks, die eintreffen, bevor der Handler subscribed hat, werden
gepuffert und beim ersten onData nachgeliefert — in dem Turn
zwischen “Call landet in der Mailbox” und “Handler läuft” geht also
nichts verloren.
Bidirektionales Streaming
Abschnitt betitelt „Bidirektionales Streaming“Bidi öffnet genauso, hält aber beide Richtungen offen:
client.tell({ kind: 'bidiStart', method: 'Chat', target: collector });Es benutzt weiterhin den älteren In-Band-Handshake: der Actor
stellt einen ersten stream-data-Frame zu, dessen Chunk
{ __streamId } ist, und bidiSend / bidiClose adressieren diese
nackte Zahl.
// Inside the collector, after receiving the streamId hint:const streamId = (message.chunk as { __streamId: number }).__streamId;
client.tell({ kind: 'bidiSend', streamId, chunk: { text: 'hello' } });client.tell({ kind: 'bidiClose', streamId });Dieser Handshake ist eine bekannte Schwachstelle — der Hinweis ist
von echten Stream-Daten nicht unterscheidbar, und die ID ist
erratbar — nachverfolgt als
#788. Der
stream-started-Frame und das Capability-Handle von oben sind das,
was er übernehmen wird; schreib neuen Code gegen die
Client-Stream-Form, wo du die Wahl hast.
Server-Stream-Calls tragen ebenfalls eine streamId auf ihren
stream-data- / stream-end-Frames, sodass ein einzelner
Collector mehrere gleichzeitige Streams multiplexen kann. Streams
bleiben offen, bis eine Seite sie schließt oder der Actor stoppt.
TLS ist über withCredentials opt-in. Wird es weggelassen,
verwenden beide Actors standardmäßig insecure
({ kind: 'insecure' }). Zertifikate werden als
Uint8Array-Werte übergeben, nicht als Dateipfade:
import { readFileSync } from 'node:fs';
// Server: supply cert + key. Add `rootCerts` to require client// certs (mTLS).const grpcServerOptions = GrpcServerOptions.create() .withProtoPath(protoPath) .withPackageName('sensor.v1') .withServiceName('SensorService') .withBind('0.0.0.0:50051') .withHandlers({ /* … */ }) .withCredentials({ kind: 'tls', cert: readFileSync('./server.crt'), key: readFileSync('./server.key'), });new GrpcServerActor( grpcServerOptions,);
// Client: pass `rootCerts` to verify the server; add cert + key for mTLS.const grpcClientOptions = GrpcClientOptions.create() .withProtoPath(protoPath) .withPackageName('sensor.v1') .withServiceName('SensorService') .withEndpoint('sensor.svc:50051') .withCredentials({ kind: 'tls', rootCerts: readFileSync('./ca.crt'), });new GrpcClientActor( grpcClientOptions,);Für Mutual TLS (mTLS) füge dem Credentials-Objekt cert + key
der Gegenseite hinzu.
Health-Checking
Abschnitt betitelt „Health-Checking“Der Server kann den Standard-Dienst grpc.health.v1.Health
neben deinem eigenen bereitstellen, sodass grpc_health_probe, die
gRPC-Probe von Kubernetes und gRPC-Load-Balancer den Knoten fragen
können, ob er bereit ist.
Der Status ist kein zweiter Begriff von „gesund“: er stammt aus
derselben HealthCheckRegistry,
die auch den /ready-Endpunkt des Management-Servers speist — der
einen pro ActorSystem, erreichbar über healthChecksOf(system)
und bereits mit den Cluster-Readiness-Checks des Frameworks
bestückt. Genau diese Registry an withHealth zu reichen ist das
Opt-in:
import { healthChecksOf } from 'actor-ts/management';
const health = healthChecksOf(system);health.addReadiness(() => ({ name: 'journal', status: journal.isConnected() }));
const grpcServerOptions = GrpcServerOptions.create() .withProtoPath(protoPath) .withPackageName('sensor.v1') .withServiceName('SensorService') .withBind('0.0.0.0:50051') .withHandlers({ /* … */ }) .withHealth(health);Check antwortet nur dann SERVING, wenn jeder
Readiness-Check besteht, und NOT_SERVING, sobald einer
fehlschlägt — buchstäblich dieselbe Regel, die /ready anwendet, das
exportierte
isHealthy.
Ausgewertet wird pro Aufruf, es wird also nichts zwischengespeichert.
Weil beide healthChecksOf(system) lesen, bleiben der gRPC-Dienst und
/ready schon von der Konstruktion her im Gleichschritt; ein
blankes new HealthCheckRegistry() an dieser Stelle würde sie
auseinanderlaufen lassen, und nur eine der beiden wäre die Antwort,
auf die ein Load Balancer reagiert.
Das gilt auch für den Weg aus dem Betrieb: Nach cluster.leave()
bleiben die Readiness-Checks des Clusters registriert und schlagen
fehl. Check antwortet für einen drainierten Knoten also
NOT_SERVING statt SERVING auf einer leeren Registry.
Das Feld HealthCheckRequest.service wählt aus, wonach gefragt
wird:
service | Antwort |
|---|---|
'' (leer) | Der ganze Server — die übliche Probe. |
sensor.v1.SensorService | Der bereitgestellte Dienst, voll qualifiziert. |
SensorService | Dasselbe; der bloße Name wird der Bequemlichkeit halber akzeptiert. |
grpc.health.v1.Health | Der Health-Dienst selbst. |
| alles andere | NOT_FOUND (Statuscode 5). |
Implementiert ist nur Check; Watch antwortet mit
UNIMPLEMENTED — das dokumentierte Signal für einen Client, auf
das Pollen von Check zurückzufallen.
Konfiguration
Abschnitt betitelt „Konfiguration“Instanzbezogene Builder-Optionen haben Vorrang vor HOCON, das
wiederum Vorrang vor den eingebauten Defaults hat. Die
HOCON-Schlüssel liegen unter actor-ts.io.broker.grpc.server und
actor-ts.io.broker.grpc.client:
actor-ts.io.broker.grpc { server { protoPath = "./proto/sensor.proto" packageName = "sensor.v1" serviceName = "SensorService" bind = "0.0.0.0:50051" } client { protoPath = "./proto/sensor.proto" packageName = "sensor.v1" serviceName = "SensorService" endpoint = "sensor.svc:50051" deadlineMs = 30s }}Handler und TLS-Credentials werden über den Builder bereitgestellt, nicht über HOCON.
Peer-Dependency
Abschnitt betitelt „Peer-Dependency“npm install @grpc/grpc-js @grpc/proto-loader# or: bun add @grpc/grpc-js @grpc/proto-loaderBeide Pakete sind Peer-Dependencies.
Wann gRPC
Abschnitt betitelt „Wann gRPC“Zwei primäre Einsatzfälle:
- Service-zu-Service innerhalb eines Clusters, wenn Protobuf-typisierte Verträge für Evolution wichtig sind.
- Externe Clients, die bereits gRPC sprechen (Mobile-Apps, Services in anderen Sprachen).
Für interne Actor-zu-Actor-Kommunikation innerhalb eines Clusters ist der Cluster-Transport besser — direkt über TypeScript typisiert, kein Protobuf erforderlich. gRPC ist für sprachübergreifende oder externe Vertrags-Fälle.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- I/O-Übersicht — das große Bild.
- BrokerActor-Basis — der gemeinsame Lifecycle.
- Refs zwischen Nodes — für clusterinterne RPCs (TypeScript-typisierte Alternative).
- HTTP-Übersicht — für HTTP-basierten RPC stattdessen.
