Adding a language SDK¶
Any language can host a streamlet: the process speaks the streamlet protocol
to the sidecar in its pod, and the sidecar does everything else. An SDK is compatible when it passes two
things, the descriptor fixtures and the conformance suite. Nothing else about it is prescribed.
The Python SDK in
sdks/python is the worked
example; its reference shows one shape an SDK can take.
What an SDK implements¶
- Declaration. A way to declare a streamlet: its name, its inlets and outlets with JSON contracts
(a schema name, fingerprinted as the Base64 of the SHA-256 of that name), and typed parameters of the
six
ConfigTypes. Registration is explicit; the SDK never discovers streamlets by scanning. - The descriptor writer. The declaration as canonical JSON, exactly as the Descriptor page defines it: snake_case field names, sorted keys, ports sorted by name, defaults omitted, two-space indentation, one trailing newline. It should also refuse what the validation rules refuse.
- Two gRPC services, served on
127.0.0.1:$FLOW_PROCESS_PORTand nowhere else:Discovery(answerDiscoverwith the descriptor; log the problemsReportErrorcarries) andStreamlet.Run. - The conversation. Wait for
Startand apply itsconfig_json. For eachBatch, call the user's code, send eachEmitas it is produced, then exactly oneAck, or oneFailif the code raised. Run batches of different partitions concurrently; never reorder one batch's messages. A newRunvoids everything tied to the previous one. OnStop, finish in-flight batches and complete the stream. An emit without a key leavesRecord.keyunset.
Everything else — codecs, test harnesses, project templates, the build tool — is the SDK's own business.
Copying the protocol¶
Copy the whole protocol/ directory into the SDK verbatim: the .proto files, README.md,
DESCRIPTOR.md and fixtures/. Generate code from the copy. CI diffs every copy against protocol/, so
a change to the protocol is a change to every SDK in the same commit. The Python SDK's
scripts/proto.py does the copy and the generation.
The descriptor fixtures¶
protocol/fixtures/declarations/*.md describe six streamlets in prose. Declare each in the SDK's
language, write its descriptor with the sdk block pinned to {"name": "fixture", "version": "0.0.0"},
and assert that the bytes equal protocol/fixtures/descriptors/<name>.json.
The conformance suite¶
Implement the reference streamlet of protocol/fixtures/declarations/conformance.md. It behaves by each
record's key: echo emits the record to out; fan emits it to out and other; skip emits
nothing; fail fails the batch; late waits 300 ms, then echoes; rogue-outlet emits to an outlet
named nope; multiply emits the record factor times with a header n=<i>; unkeyed emits it with
no key; header-echo emits it with its headers reversed; any other key echoes.
Serve it on a port and run the suite from the ankka-flow repository:
sbt 'sidecar/testOnly *ConformanceSuite' -Dflow.conformance.target=127.0.0.1:9010
sbt 'sidecar/testOnly *ConformanceSuite' -Dflow.conformance.target=127.0.0.1:9010 -Dflow.conformance.only=run.fan-out
Every case is named for the conversation it checks, such as run.emits-precede-ack or
run.two-partitions-interleave, so a failure says what the SDK got wrong. The scripted inputs each case
sends are in protocol/fixtures/conversations/, readable without Scala. The Python SDK wraps this in
uv run conformance; a new SDK can offer the same.
The violation.* and version.* cases are skipped against a process. They prove that the sidecar
refuses misbehaviour a correct SDK cannot produce — a double acknowledgement, an emit after its ack, an
unknown batch, another major version — using a scriptable double inside the suite's own JVM.