
Cosmonic Control v0.12 ships native Kafka through cosmonic:kafka@0.5.0. Workloads can
publish records, process events, use a stateful pull consumer, or commit Kafka transactions
without carrying a broker address or credential in component code.
For serverless event processing, start with the handler interface. The host owns the Kafka
consumer and keeps its group membership stable. It creates component instances when partitions
have work, runs different partitions concurrently, and can reclaim every guest instance when the
stream is quiet. Kafka stays connected while your compute scales independently, all the way to
zero.
Start with the handler
A handler exports one function. The host polls Kafka and calls that function with a partition-ordered batch. This deployment permits up to 32 concurrent component calls and lets the warm pool return to zero after 30 seconds:
apiVersion: runtime.wasmcloud.dev/v1alpha1
kind: WorkloadDeployment
metadata:
name: order-transformer
spec:
replicas: 1
template:
spec:
hostInterfaces:
- namespace: cosmonic
package: kafka
version: '0.5.0'
interfaces: [handler, producer]
config:
bootstrap.servers: kafka.kafka.svc:9092
topics: orders.input,orders.output,orders.dlq
handler.topics: orders.input
handler.group.id: order-transformer
handler.batch.size: '100'
dead-letter.topic: orders.dlq
acks: all
secretFrom:
- name: kafka-orders-credentials
components:
- name: processor
image: ghcr.io/example/order-transformer:1.0.0
poolSize: 32
maxConcurrency: 1
reclaimWindowSeconds: 30
reclaimMinInstances: 0The component exports cosmonic:kafka/handler@0.5.0. This example also imports producer, so a
handler can transform an input record and publish its result without opening another Kafka client.
Useful concurrency for one deployment replica is bounded by:
min(assigned partitions, 64, poolSize x maxConcurrency)poolSize controls elastic guest compute, not Kafka group membership. Increasing it does not add
consumers or trigger a rebalance. Increase workload replicas when you need another host or group
member for availability or more assigned partitions.
That separation makes handler the best default for serverless and elastic workloads. A burst can
use more component instances without waiting for pod scheduling, and those instances can disappear
while the host preserves the consumer session.
Delivery is explicit
Handler delivery is at least once. The host advances offsets only after the workload reports progress, so a crash between an external side effect and an offset commit can deliver a record again. Make handlers idempotent, or use a Kafka transaction when output records and consumed offsets must commit atomically.
A handler can report exactly how much of a batch succeeded:
ok(none)accepts the full batch.ok(some(offset))accepts records through that offset and retries the remaining suffix.- A
transienterror retries with exponential backoff. - A
permanenterror isolates the head record and sends it to the configured dead-letter topic. - A trap or timeout retries the head record and dead-letters it after five consecutive failures.
Calls contain records from one partition in offset order. The default maximum batch size is 100 records, and different partitions can run concurrently. Treat malformed input as permanent rather than trapping, and keep required state outside the component instance: the next call may run on a different instance.
Least privilege is part of the binding
In 0.5.0, workloads do not pass connection or client settings to open. Producer functions are
binding-scoped, and consumer.open() takes no arguments. Broker addresses, credentials, group IDs,
transaction IDs, limits, and client policy come from hostInterfaces, outside the guest.
The topics setting is an authorization boundary, not only a convenience. It must grant every
subscription, dead-letter topic, and topic named by a producer call. handler.topics selects the
subscription from that grant. A component cannot redirect itself to another broker, substitute a
credential, or publish to an ungranted topic.
Keep plain policy in config and credentials in a Kubernetes Secret referenced by secretFrom.
The host resolves the Secret and supplies it to the native client; the component never receives the
password or TLS material.
You can also use named bindings to give one component deliberately separate Kafka authorities. For example, two producer imports can route to different clusters or principals, each with its own topic grant.
Choose the interface that matches the workload
The 0.5.0 package has four distinct application patterns:
handleris the default for serverless event processing. The host owns polling, group membership, retries, offsets, and dead-letter delivery while guest compute scales independently.produceris for publishing.send,send-batch, andsend-streamcall a producer attached to the binding; there is no guest-owned producer resource,open,close, or globalflush.transactionis for atomic Kafka work. Its resource can publish records and commit consumed offsets together. Each live replica needs a distinct, stabletransactional.id.consumeris for long-running services that need direct session control. Choose it for manual assignment, rebalance events, pause, seek, explicit commits, or guest-controlled pull pace. Those operations make the workload stateful, so this is not the default serverless shape.
The native backend
Cosmonic Control's production backend runs on librdkafka through the Rust rdkafka crate. The host
reuses a native producer for each binding, so concurrent send calls can benefit from librdkafka's
own batching. Blocking metadata and transaction operations stay off the async execution path.
The backend supports TLS, SASL, transactions, pull consumers, and host-driven handler delivery. Its broker-backed test suite covers authentication success and rejection, unavailable brokers, disconnects, queue backpressure, transactions, group rebalancing, retry behavior, and dead-letter delivery.
Native Kafka allocations live outside Wasm memory, so the plugin bounds clients, streams, queues, and handler buffers per component. It also emits bounded-cardinality OpenTelemetry metrics for active clients, guest operations, operation duration, and pull-consumer records. Dedicated handler throughput and outcome metrics are not part of 0.5.0 yet.
The second backend: a Kafka client inside the sandbox
Because Kafka is exposed as a component interface, we can explore a very different implementation.
The experimental backend is a host component plugin: a Wasm component that speaks the Kafka wire
protocol over wasi:sockets. There is no native Kafka client in the host process. The long-lived
poller is itself sandboxed, versioned, and shipped like any other component.
That placement matters for serverless Kafka. Broker connections, metadata, heartbeats, and group membership must outlive an ephemeral handler invocation. Most platforms put that state in trusted native infrastructure. This experiment preserves the same split between stable polling and elastic work while moving the poller into a sandboxed component.
The plugin's push path dispatches partition-ordered batches to a workload exporting
cosmonic:kafka/handler, with retries, partial progress, offset commits, and dead-letter delivery.
Its Kafka protocol implementation is written directly in Rust and tested end to end from Wasm
against Redpanda and Apache Kafka.
This backend remains experimental. It has no TLS, SASL, or transactions, and its current
cosmonic:kafka@0.2.0 contract is not drop-in compatible with the native 0.5.0 interface. The
source and runnable example
show what becomes possible when even stateful infrastructure can run as a sandboxed component.
Try it out
Native Kafka is available in
Cosmonic Control v0.12 and later. The Kafka
interface is published at
ghcr.io/cosmonic-labs/cosmonic/kafka:0.5.0. The
awesome-cosmonic Kafka templates
include complete Rust projects and manifests for all four patterns. Start with
kafka-handler-consumer
for serverless or elastic event processing, then reach for the pull service or transaction template
only when the workload needs the state they own.
The larger point is not merely that WebAssembly can call Kafka. It is that Kafka access can be a least-privilege capability: platform policy owns the connection and authority, the application sees a typed interface, and the handler keeps stable broker state separate from compute that can scale up in milliseconds and return to zero.
