Agent2Agent over Apache Kafka

Agent messages that survive the restart.

ka2a is a Go SDK and command-line tool for durable agent-to-agent communication. It signs each A2A message, carries it through Kafka and records every stage in a local SQLite store. After a crash, each side knows what it sent, what it received and what is still unknown.

Developer preview. Not yet production ready. Public installation is not yet verified.

Two agents exchange a sealed message through a shared log. Each side keeps its own record.
  • Pure GoNo cgo. The SQLite driver is pure Go.
  • Signed recordsEd25519 over topic, key, headers and value.
  • SQLite recoveryOutbox, receive journal and task state survive a crash.
  • A2A 1.0 typesSendMessage, GetTask, ListTasks and CancelTask.

How it works

Four stages. None implies the others.

A queued request is not a delivered request. A delivered request is not a finished task. ka2a records each stage with its own evidence, so your code never has to guess.

A request from planner to reviewer passes four recorded stages: locally queued, broker acknowledged, remote recorded and protocol response. planner demo.example Kafka mailbox topic of reviewer reviewer demo.example 1 2 3 4 locally_queued outbox row + frozen bytes broker_acknowledged acks from all in-sync replicas remote_recorded verified, then journaled protocol_response
Figure 1 A request from planner to reviewer. Each numbered stage is a separate record in a store or on the broker.
  1. Locally queued

    One SQLite transaction stores the operation, its ID and the exact record bytes. Only then does publication start.

  2. Broker acknowledged

    Kafka acknowledges the record with acks=all. A timeout gives unknown, never a failure.

  3. Remote recorded

    The receiver checks the signature, the catalog grant and the schema. It journals the request before it accounts the offset.

  4. Protocol response

    The handler result travels back as a signed record. A response ends the exchange; it does not mean the task completed.

Examples

From a Go handler to a terminal session

The Go listings compile against the public API of the reviewed revision. The terminal sessions are real runs against a local single-broker Kafka with two synthetic endpoints. Their keys are demo keys that must never protect real data.

Go · pkg/ka2a

Embed a node that serves requests

Open takes exclusive ownership of the state directory and opens the store. Server registers one handler. Run runs restart recovery, checks the mailbox topic and starts the work. WaitReady returns when the node admits work, also before a broker answers.

  • The handler receives the verified principal, never a claimed one.
  • Status tells whether a broker answered the mailbox check, and the live broker state.
  • HoldOnAmbiguity holds a request for an operator when a crash leaves its effect unknown.
  • Unsupported operations answer with a typed A2A error.
reviewer/main.goGo
node, err := ka2a.Open(ctx, ka2a.Config{
	Endpoint:       ka2a.Endpoint{Domain: "demo.example", Name: "reviewer"},
	StateDir:       "/var/lib/ka2a/reviewer",
	CatalogFile:    "/etc/ka2a/catalog.json",
	SigningKeyFile: "/etc/ka2a/reviewer.key.json",
	Kafka: ka2a.Kafka{
		Brokers: []string{"kafka-1.demo.example:9093"},
		TLS:     &tls.Config{MinVersion: tls.VersionTLS13},
	},
})
if err != nil {
	log.Fatal(err)
}
defer func() { _ = node.Close(ctx) }()

// The handler gets the verified principal and a stable operation ID.
handler := ka2a.HandlerFunc(func(ctx context.Context, inv *ka2a.Invocation) (ka2a.Result, error) {
	if _, ok := inv.Request.SendMessage(); !ok {
		return ka2a.ErrorResult(a2a.ErrUnsupportedOperation), nil
	}
	log.Printf("request %s from %s", inv.OperationID, inv.Principal)
	return ka2a.TaskResult(&a2a.Task{
		ID:        a2a.TaskID("review-" + inv.OperationID),
		ContextID: "release-notes",
		Status:    a2a.TaskStatus{State: a2a.TaskStateCompleted},
	}), nil
})
err = node.Server(handler, ka2a.HandlerOptions{RecoveryClass: ka2a.HoldOnAmbiguity})
if err != nil {
	log.Fatal(err)
}
go func() { _ = node.Run(ctx) }()
if err := node.WaitReady(ctx); err != nil {
	log.Fatal(err)
}
Listing 1 A node of the endpoint reviewer with one handler.

Go · pkg/ka2a

Send a message and wait for the answer

SendMessage returns after the request is durable. Wait returns the correlated response, or ErrUnknownOutcome when no answer arrived in time.

  • An unknown outcome is not a failure. The request stays admitted and can still complete.
  • A retry with the same operation ID returns the original operation.
  • An A2A error from the peer is a response, not a transport failure.
planner/send.goGo
// Keep the operation ID with the intent. After a crash, the same ID
// returns the original operation instead of a second request.
operationID, err := ka2a.NewOperationID()
if err != nil {
	log.Fatal(err)
}
operation, err := node.Client().SendMessage(ctx,
	ka2a.Endpoint{Domain: "demo.example", Name: "reviewer"},
	&a2a.SendMessageRequest{Message: &a2a.Message{
		ID:    "msg-1",
		Role:  a2a.MessageRoleUser,
		Parts: a2a.ContentParts{a2a.NewTextPart("Draft release notes for v0.4")},
	}},
	ka2a.SubmitOptions{OperationID: operationID})
if err != nil {
	log.Fatal(err) // Nothing was queued.
}
response, err := operation.Wait(ctx)
switch {
case errors.Is(err, ka2a.ErrUnknownOutcome):
	// Not a failure: the request stays admitted. Look it up later.
	status, _ := operation.Status(ctx)
	log.Printf("still %s after %d attempts", status.Exchange, status.Send.Attempts)
case err != nil:
	log.Fatal(err)
case response.IsError():
	log.Printf("the peer answered with an error: %v", response.ErrorObject().A2AError())
default:
	result, _ := response.SendMessageResult()
	log.Printf("result: %T", result)
}
Listing 2 The client side: send, wait and handle each kind of result.

Configuration · ka2a.catalog.v1

Declare who may talk to whom

A static catalog lists the endpoints, their public keys and the operations that each source may call on each destination. The runtime never discovers peers by itself.

  • Strict JSON: unknown members, duplicate keys and trailing data fail.
  • The catalog digest becomes a configuration version in the store.
  • The mailbox topic name comes from a hash of the endpoint.
catalog.jsonSynthetic
{
  "format": "ka2a.catalog.v1",
  "mailbox": { "prefix": "ka2a", "rule": "sha256-v1" },
  "endpoints": [
    { "domain": "demo.example", "endpoint": "planner" },
    { "domain": "demo.example", "endpoint": "reviewer" }
  ],
  "keys": [
    {
      "key_id": "planner-2026-09",
      "algorithm": "ed25519",
      "public_key": "fqOnxXXHZaAEqupZK8jtSjpGZNv536-yffmbJQIvmPw",
      "domain": "demo.example",
      "endpoint": "planner",
      "not_before": "2026-09-01T00:00:00Z"
    },
    {
      "key_id": "reviewer-2026-09",
      "algorithm": "ed25519",
      "public_key": "BP6231K7kz18wX9sPQf2SkgOrgo6h2PaoiHK3z9qHAE",
      "domain": "demo.example",
      "endpoint": "reviewer",
      "not_before": "2026-09-01T00:00:00Z"
    }
  ],
  "grants": [
    {
      "source": { "domain": "demo.example", "endpoint": "planner" },
      "destination": { "domain": "demo.example", "endpoint": "reviewer" },
      "operations": ["SendMessage", "GetTask", "ListTasks", "CancelTask"]
    },
    {
      "source": { "domain": "demo.example", "endpoint": "reviewer" },
      "destination": { "domain": "demo.example", "endpoint": "planner" },
      "operations": ["SendMessage"]
    }
  ]
}
Listing 3 The catalog of the two demo endpoints. The keys are non-production demo keys.

CLI · ka2a init, topics, doctor

Set up an endpoint and check it

init creates the store and records the catalog. topics provision creates the mailbox topic and never changes an existing one. doctor reports each check as ok, warning, problem or unverified.

  • The runtime never creates or changes a topic. Provisioning is an explicit step.
  • Doctor never reports a setting as verified without a value from the broker.
  • The development profile says clearly what it does not prove.
terminalReal run · synthetic endpoints
$ ka2a init --state-dir planner --domain demo.example --endpoint planner \
    --catalog catalog.json --signing-key planner.key.json
Created the store of demo.example/planner in /srv/ka2a-demo/planner.
Node ID       5a6a9a52-8ab7-4076-a85e-f891a8299ac8
Generation    642c53fc-2033-407d-b89e-34fed810658a
Epoch         1
Catalog       sha256:536ccf39…87cbef45, configuration version 1
Mailbox topic ka2a.v1.mailbox.c5b7df51…6c16230e
Signing key   planner-2026-09, valid at 2026-10-02 18:29:31Z
Next: provision the mailbox topic with 'ka2a topics provision', then check the setup with 'ka2a doctor'.

$ ka2a topics provision --brokers 127.0.0.1:39192 --catalog catalog.json \
    --state-dir planner --profile development
Topic ka2a.v1.mailbox.c5b7df51…6c16230e was created. Profile development.
STATUS  CHECK                                  EXPECTED            ACTUAL     DETAIL
ok      broker.topic_exists                    true                true       Verified from the broker metadata.
ok      broker.partitions                      1                   1          Verified from the broker metadata.
ok      broker.replication_factor              1                   1          Verified from the broker metadata.
ok      broker.in_sync_replicas                >= 1                1          Verified from the broker metadata.
ok      broker.cleanup.policy                  delete              delete     Verified from the broker metadata.
ok      broker.retention.ms                    >= 691200000 or -1  691200000  Verified from the broker metadata.
ok      broker.max.message.bytes               >= 1065445          1065445    Verified from the broker metadata.
ok      broker.min.insync.replicas             >= 1                1          Verified from the broker metadata.
ok      broker.unclean.leader.election.enable  false               false      Verified from the broker metadata.
Result: verified.
Listing 4 Create the store and the mailbox topic of planner. An ellipsis (…) marks shortened digests.
terminalReal run · synthetic endpoints
$ ka2a doctor --state-dir planner --catalog catalog.json \
    --signing-key planner.key.json --brokers 127.0.0.1:39192 \
    --profile development
ka2a doctor: /srv/ka2a-demo/planner (checked at 2026-10-02 18:29:36Z)
STATUS   CHECK                                  DETAIL
…
ok       store.identity                         Endpoint demo.example/planner, node 5a6a9a52-…, generation 642c53fc-…, epoch 1.
info     store.owner                            No owner process holds the store lock.
ok       store.quarantine                       No quarantine.
ok       store.restore                          No restore evidence.
ok       store.storage                          Retained 0 B of 1.0 GiB (the default budget; the owner can use another). File system: ext.
…
ok       store.held_work                        No held or unknown work.
…
ok       catalog.load                           Strict ka2a.catalog.v1 with 2 endpoints; digest sha256:536ccf39…87cbef45.
ok       catalog.endpoint                       The catalog lists demo.example/planner with mailbox topic ka2a.v1.mailbox.c5b7df51…6c16230e.
ok       catalog.configuration                  The catalog is configuration version 1.
ok       signing_key                            Key planner-2026-09 signs for demo.example/planner, is private (mode 0600) and is valid at 2026-10-02 18:29:36Z.
warning  broker.profile                         Development profile: it does not demonstrate broker-loss durability.
warning  broker.security                        The broker connection has no TLS. Use plaintext only for a local development broker.
info     assumption.producer.acks               The producer requests acknowledgements from all in-sync replicas (acks=all). ka2a fixes this client setting.
…
ok       broker.topic_exists                    Verified from the broker metadata. (expected true, actual true)
ok       broker.partitions                      Verified from the broker metadata. (expected 1, actual 1)
ok       broker.replication_factor              Verified from the broker metadata. (expected 1, actual 1)
…
Result: ok. 0 problems, 2 warnings, 0 unverified.
Listing 5 Check the setup. An ellipsis (…) marks shortened identifiers and digests, and cut rows.

CLI · ka2a send, operation show

Send from the command line and follow the stages

The command line has no application handler, but it can send. Here the reviewer runs the handler of Listing 1. operation show then reads the store and prints each stage with its evidence.

  • Identifiers are pseudonyms by default. --raw-ids shows the raw values.
  • Bodies appear only with --show-result.
  • The task state is the state that the response carried. The peer owns the task.
terminalReal run · synthetic endpoints
$ ka2a send --state-dir planner --catalog catalog.json \
    --signing-key planner.key.json --brokers 127.0.0.1:39192 \
    --allow-plaintext --profile development --peer demo.example/reviewer \
    --operation-id 6f1c2a9e-4b7d-4c3e-9a51-2d8e7f0b3c14 \
    --text "Draft release notes for v0.4" --show-result
SendMessage to demo.example/reviewer: result.
Operation  op~r6avsxfdyqi5siuy (created)
STAGE                STATE
locally_queued       yes
broker_acknowledged  yes
remote_recorded      yes
protocol_response    responded
application_task     TASK_STATE_COMPLETED
Task       task~qc5gmwmtzthj2hwp, state TASK_STATE_COMPLETED, context ctx~45vgzjnfd4zhv6q2, 0 history messages, 0 artifacts
Result document:
{"id":"review-6f1c2a9e","contextId":"release-notes","status":{"message":{"messageId":"01a0fde1-1cfc-7943-b9db-131a96fe86a6","parts":[{"text":"Reviewed \"Draft release notes for v0.4\": 2 notes, no blocking findings."}],"role":"ROLE_AGENT"},"state":"TASK_STATE_COMPLETED"}}
Identifiers: pseudonyms keyed for this output only; --raw-ids shows raw identifiers.
Listing 6 One request with its result document.
terminalReal run · synthetic endpoints
$ ka2a operation show \
    --state-dir planner 6f1c2a9e-4b7d-4c3e-9a51-2d8e7f0b3c14
Operation op~leehqdk5tzfdrlob (SendMessage) to demo.example/reviewer, local sequence 1.
Read from /srv/ka2a-demo/planner at 2026-10-02 18:29:43Z. Historical: no owner process held the store lock at capture time. The data shows the last state that a process committed.
Identifiers: pseudonyms keyed for this snapshot only. Bodies are never shown.

STAGE                STATE                 EVIDENCE
locally_queued       yes                   Admitted to the local outbox at 2026-10-02 18:29:43Z.
broker_acknowledged  yes                   Kafka acknowledged partition 0 offset 0 at 2026-10-02 18:29:43Z.
remote_recorded      yes                   A response from demo.example/reviewer arrived at 2026-10-02 18:29:43Z.
protocol_response    responded             Outcome result. A response ends the exchange; it does not show that the task completed.
application_task     TASK_STATE_COMPLETED  The state that the response carried. The peer owns the task; this is not a live state.

Send         broker_ack, 1 attempts, since 2026-10-02 18:29:43Z
Exchange     responded since 2026-10-02 18:29:43Z
Response     result from demo.example/reviewer (key reviewer-2026-09), event ev~sdu26cvnlnb2n4x4, received 2026-10-02 18:29:43Z
Task         TASK_STATE_COMPLETED, task task~m27sq2twjnj5xjif, context ctx~bwxcmuoxy2x2liim (observed in the response; the peer owns the task)
Times        created 2026-10-02 18:29:43Z, admitted 2026-10-02 18:29:43Z, expires 2026-10-03 18:29:43Z, deadline none
Recovery     hold_on_ambiguity; request event ev~3qbpd3s3x2s45ien
Listing 7 The recorded stages of the same operation.

CLI · unknown outcomes

When the peer is away

An unanswered request is not lost and not failed. ka2a send exits with 6: the outcome is not known yet. The same request with the same operation ID later returns the original operation and its response.

  • A retry reuses the frozen record. It never signs a new request under the old ID.
  • A changed request under a used ID is an identity conflict.
  • Exit codes are stable, so scripts can react to each case.
terminalReal run · synthetic endpoints
# The reviewer is offline. Wait five seconds, then stop waiting.
$ ka2a send … --peer demo.example/reviewer \
    --operation-id 2b7e9c41-8f3a-4d62-b0e5-7a1c9d4f6e28 \
    --text "Summarize incident 311" --wait 5s
SendMessage to demo.example/reviewer: pending.
Operation  op~aasno2j4u3mp7no2 (created)
STAGE                STATE
locally_queued       yes
broker_acknowledged  yes
remote_recorded      unknown
protocol_response    awaiting_response
application_task     none
No response yet. The request stays admitted and can still complete; a later start of the node records the response.
Identifiers: pseudonyms keyed for this output only; --raw-ids shows raw identifiers.
# exit status 6: the outcome is not known yet. It is not a failure.

# The reviewer is back. Send the same request with the same operation ID.
$ ka2a send … --peer demo.example/reviewer \
    --operation-id 2b7e9c41-8f3a-4d62-b0e5-7a1c9d4f6e28 \
    --text "Summarize incident 311" --wait 20s
SendMessage to demo.example/reviewer: result.
Operation  op~6ca7irrorzjmkgfs (the existing operation of this key)
STAGE                STATE
locally_queued       yes
broker_acknowledged  yes
remote_recorded      yes
protocol_response    responded
application_task     TASK_STATE_COMPLETED
Task       task~wo3ozqewwijocygx, state TASK_STATE_COMPLETED, context ctx~7dvzjhzt7dzrzzy2, 0 history messages, 0 artifacts
Identifiers: pseudonyms keyed for this output only; --raw-ids shows raw identifiers.
# exit status 0. The reviewer ran the request once.

# The same operation ID with a different request is refused.
$ ka2a send … --peer demo.example/reviewer \
    --operation-id 2b7e9c41-8f3a-4d62-b0e5-7a1c9d4f6e28 \
    --text "Summarize incident 312" --wait 5s
ka2a send: identity_conflict: the operation ID is already used by another request (the request fingerprint differs)
# exit status 1
Listing 8 The reviewer is offline, then back. The ellipsis (…) stands for the connection flags of Listing 6.

Wire · KA2A-WIRE-1

What travels on the wire

Each request is one Kafka record in the CloudEvents binary mode. The value is canonical JSON (RFC 8785). The signature covers the topic, the key, the sorted headers and the value.

  • The record key is the operation ID, so retries land on the same key.
  • ce_ka2aexpires bounds how long a receiver accepts the record.
  • Frozen test vectors keep the byte format stable across releases.
testdata/wire/request.jsonFixture
# Frozen test vector "request" of KA2A-WIRE-1. NON-PRODUCTION fixture key.
topic    ka2a.v1.mailbox.ef118f190fbbb51fbf1c872b534476430cc1caf73ed676d1065bb16f697d50a8
key      0f010000-0000-4000-8000-00000000a001
headers  ce_dataschema   urn:sha256:673d5641e3a3d163271068fd6abdd570c8439659f8a91811c4b4a0ef4056bedc
         ce_id           0e010000-0000-4000-9000-00000000b001
         ce_ka2aa2a      1.0
         ce_ka2adomain   fixture.test
         ce_ka2aexpires  2026-10-01T12:00:00.000000Z
         ce_ka2akey      fixture-client-a-1
         ce_ka2aop       0f010000-0000-4000-8000-00000000a001
         ce_ka2asig      xG-guTuDp-Bs2bm1VgVXLfWu4Zm_qipqG6s-8OnkCv6eYZS_ablaU9vJ8im2n7f0oB3-zTmySN9D6lo0RkFBDg
         ce_ka2atarget   server-b
         ce_ka2awire     1
         ce_source       urn:ka2a:endpoint:fixture.test:client-a
         ce_specversion  1.0
         ce_time         2026-09-30T12:00:00.000000Z
         ce_type         io.ka2a.request.v1
         content-type    application/json
value    {"operation":"SendMessage","params":{"message":{"messageId":"msg-0001","parts":[{"text":"hello, server"}],"role":"ROLE_USER"}}}
Listing 9 The frozen test vector request. It is signed with a public, non-production fixture key.

CLI · ka2a observe

Inspect a store without touching it

observe opens the store read-only and prints one consistent snapshot. It never takes the owner lock, writes a row or moves a broker offset. --output json prints the ka2a.snapshot/1 document.

  • The snapshot says if an owner process held the lock. Without one, the data is historical.
  • A store snapshot does not observe the broker, and it says so.
terminalReal run · synthetic endpoints
$ ka2a observe --state-dir planner
Snapshot of /srv/ka2a-demo/planner, captured at 2026-10-02 18:30:03Z.
Historical: no owner process held the store lock at capture time. The data shows the last state that a process committed.
Identifiers: pseudonyms keyed for this snapshot only. Bodies are never shown.

Endpoint     demo.example/planner
Node ID      5a6a9a52-8ab7-4076-a85e-f891a8299ac8
Generation   642c53fc-2033-407d-b89e-34fed810658a
Epoch        6
Quarantine   none
Storage      8.6 KiB of 1.0 GiB (default budget), normal, 0.0 %
Broker       unknown: A store snapshot does not observe the broker. The reader does not connect to Kafka.
Config       version 1, sha256:536ccf393ba0b9b1206b8e0a81386bce01a5df66fd9de31bc8b9addd87cbef45, catalog catalog.json: 2 endpoints
Outbox       0 records wait for publication

Counts
  outbox     queued 0, publishing 0, broker_ack 3, held 0, error 0
  exchanges  awaiting_response 0, responded 3, held 0, abandoned 0
  dispatch   queued 0, started 0, completed 0, held 0, abandoned 0
  other      tasks 0, rejections 0, tombstones 0, partitions 1 (committed ahead 0)

Operations (newest first; all states; 3 shown)
SEQ  OPERATION            NAME         TARGET                 SEND        EXCHANGE   TASK
3    op~rimzfcmwt2bbietx  SendMessage  demo.example/reviewer  broker_ack  responded  TASK_STATE_COMPLETED
2    op~pr4rjw3j3uh3vukg  ListTasks    demo.example/reviewer  broker_ack  responded  none
1    op~3nqavsg5mgo5sfau  SendMessage  demo.example/reviewer  broker_ack  responded  TASK_STATE_COMPLETED

Dispatch (oldest first; states started, held; 0 shown)
  none

Outbox records (oldest first; states queued, publishing, held, error; 0 shown)
  none

Rejections (newest first; 0 shown)
  none

Partition progress (topic partition; 1 shown)
TOPIC                                                                             PARTITION  ACCOUNTED_NEXT  COMMITTED_NEXT  UPDATED
ka2a.v1.mailbox.c5b7df518114d8dd3f3c8c2efd196681e22d370262550787a019dccc6c16230e  0          3               3               2026-10-02 18:30:00Z

Incidents (newest first; 0 shown)
  none
Listing 10 A snapshot of planner. The partition progress shows the accounted and the committed offsets.

CLI · ka2a run --metrics-listen

Scrape metrics that hold no secrets

ka2a run can serve the node metrics in the Prometheus text format on a loopback address. Embedders mount the same exporter with Node.MetricsHandler.

  • Labels carry fixed states only, plus the configured domain and endpoint. No operation ID, context, body, key or broker address appears.
  • The broker state is live: reachable, unreachable or unknown.
  • A non-loopback address is a usage error before the node opens.
  • The last credential reload and the state of the catalog on disk are metrics too.
terminalReal run · synthetic endpoints
$ ka2a run --state-dir planner --catalog catalog.json \
    --signing-key planner.key.json --brokers 127.0.0.1:39192 \
    --allow-plaintext --profile development --metrics-listen 127.0.0.1:39501
ka2a run: metrics at http://127.0.0.1:39501/metrics
ka2a run: node demo.example/planner is ready on mailbox ka2a.v1.mailbox.c5b7df518114d8dd3f3c8c2efd196681e22d370262550787a019dccc6c16230e; press Ctrl-C to stop.

# In a second terminal: read the metrics.
$ curl -s http://127.0.0.1:39501/metrics
ka2a_node_info{domain="demo.example",endpoint="planner"} 1
ka2a_node_running 1
ka2a_node_ready 1
ka2a_node_closed 0
ka2a_quarantine_active 0
ka2a_status_observed_timestamp_seconds 1.790965810311e+09
ka2a_store_retained_bytes 8798
ka2a_store_budget_bytes 1.073741824e+09
ka2a_outbox_records{state="queued"} 0
ka2a_outbox_records{state="publishing"} 0
ka2a_outbox_records{state="held"} 0
ka2a_outbox_records{state="error"} 0
ka2a_outbox_oldest_queued_age_seconds 0
ka2a_exchanges{state="awaiting"} 0
ka2a_exchanges{state="held"} 0
ka2a_dispatch{state="queued"} 0
ka2a_dispatch{state="started"} 0
ka2a_dispatch{state="held"} 0
ka2a_tasks 0
ka2a_rejections 0
ka2a_mailbox_check{result="not_checked"} 0
ka2a_mailbox_check{result="verified"} 1
ka2a_mailbox_check{result="unverified"} 0
ka2a_mailbox_check{result="broker_unreachable"} 0
ka2a_mailbox_check{result="mismatch"} 0
ka2a_broker_state{state="unknown"} 0
ka2a_broker_state{state="reachable"} 1
ka2a_broker_state{state="unreachable"} 0
ka2a_broker_last_contact_timestamp_seconds 1.790965809873e+09
ka2a_broker_last_failure_timestamp_seconds 0
ka2a_credentials_last_reload{result="none"} 1
ka2a_credentials_last_reload{result="reloaded"} 0
ka2a_credentials_last_reload{result="refused"} 0
…
ka2a_catalog_disk{state="same"} 1
ka2a_catalog_disk{state="differs"} 0
ka2a_catalog_disk{state="refused"} 0
ka2a_submit_accepted_total 0
…
ka2a_publish_acknowledged_total 0
…
ka2a_publish_unknown_total 0
…
ka2a_requests_recorded_total 0
…
ka2a_checkpoint_commits_total 0
…
ka2a_observer_dropped_total 0
…
Listing 11 The live metrics of planner. The HELP and TYPE comment lines are left out, and an ellipsis (…) marks cut metric lines.

Publication

Every publish attempt has one of four outcomes

ka2a never reports a remote failure from a timeout. It classifies what the broker told it and acts on that class.

OutcomeWhat it meansWhat ka2a does
acknowledgedAll in-sync replicas acknowledged the record. Partition and offset are known.Records the broker stage with its offset and waits for the response.
rejectedThe record is not written, and the same bytes fail again. Examples: too large, not authorized.Marks the operation as an error with a stable reason code.
not_sentThe record did not leave the process, for example after a cancellation.Keeps it queued. A retry of the same bytes is safe.
unknownThe broker possibly wrote the record, for example after a timeout or a disconnect.Publishes the same frozen bytes again within the admission horizon. The receiver deduplicates.

Recovery

Crash anywhere. Restart. Know where you are.

The store is the source of truth for local work. Kafka offsets follow the store, never the other way round.

Restart republishes the same bytes

A crash after admission leaves the record in the outbox. The next start publishes the identical bytes. The receiver sees a duplicate and answers from its journal.

Ambiguous work is held, not guessed

A handler crash can leave an effect unknown. With HoldOnAmbiguity, ka2a holds the request for an operator. With IdempotentWithKey, it runs the handler again with the same key.

Checkpoints follow the store

A consumer commits only the contiguous prefix of records that the store accounted. An epoch fence stops a stale consumer from committing after a rebalance.

Restores are detected

A witness file and the store epoch detect a rolled-back state directory. The store then enters quarantine until an operator records a decision with ka2a quarantine release.

One owner per store

An owner lock keeps a second process out of a state directory. Read-only tools work beside the owner and never take the lock.

Bounded retention

Held, unknown and unpublished work is never removed automatically. Completed rows age out under an explicit storage budget.

Operate

Read-only tools for the people on call

ka2a serve-ui serves a small operational interface on a loopback address. It renders on the server, needs no JavaScript and cannot change task, protocol or store state.

For deployment, recovery and every flag, read the operator documentation: production deployment, runbooks, the configuration reference, the metrics, error and output format references, the upgrade notes, the versioning policy, the release procedure and the threat model.

Tests keep the metrics, error and output format references complete. The documentation also has example alert rules, a systemd unit and a Containerfile. The systemd unit did not run under systemd, and the Containerfile was not built.

The ka2a operations overview for the endpoint demo.example/planner. Panels show the endpoint identity, storage pressure, broker connectivity marked unknown and the configuration digest, followed by the counts of each stage.
Figure 2 Synthetic demo data. The overview page of ka2a serve-ui for the demo endpoint planner. A banner says that the page shows a snapshot, not a live process.

Commands that read state

ka2a doctor
Checks the store, the catalog, the signing key, the broker connection and the mailbox topic. It warns when the state directory is on tmpfs or overlay, because that state is lost at a reboot or with the container.
ka2a observe
Prints one consistent snapshot of the store.
ka2a operation show
Shows the stages of one operation with evidence.
ka2a snapshot export
Writes a ka2a.snapshot/1 file for a later review.
ka2a topics verify
Compares the mailbox topic with the profile. It changes nothing.

Stable exit codes

CodeMeaning
0Success.
1The command failed. The message names the error code.
2Usage error: unknown command, flag or argument.
3The checks found problems.
4The command is not available in this build.
5A running node owns the state directory.
6A request is admitted, but its outcome is not known yet.

Broker security

Secure the broker connection. Rotate credentials without a restart.

A node and the administrative commands share one set of TLS and SASL flags. A refused connection or a missing permission is a classified problem with a fixed hint, never a vague failure.

TLS, SASL and mutual TLS

ka2a connects over TLS with SASL SCRAM-SHA-512 or PLAIN, or with a client certificate (--tls-cert, --tls-key). The key file must be private (mode 0600). The node commands refuse an expired certificate before they contact a broker.

Classified refusals and denials

A refused or revoked handshake is tls_failed, a refused login is authentication_failed, and a missing ACL is authorization_failed with a scope such as topic_write or group_read. doctor and topics verify then exit with status 3.

The minimal ACLs, as a plan

ka2a topics acl-plan derives the minimal ACLs of a node from the catalog and prints the kafka-acls.sh commands. It never contacts a broker and never applies an ACL. An operator reviews and applies them.

Credential reload without a restart

SIGHUP makes ka2a run read the TLS and SASL files again. Embedders call Node.ReloadCredentials. Each bootstrap broker in --brokers gets a check with the new material before the node uses it. A refused reload keeps the previous material, and the node continues.

Catalog and signing key need a restart

A running node never adopts a changed catalog or signing key, because the catalog is the root of trust of work in flight. The node compares the catalog on disk with its running digest and shows same, differs or refused.

Keys and catalogs from the CLI

ka2a keys generate writes a new Ed25519 key file and prints only its public catalog entry. catalog validate classifies every finding, and catalog digest prints the digest that a node records. There is no command that edits a catalog.

Revocation lists for broker certificates

--tls-crl reads local revocation lists of the CAs in --tls-ca. Embedders use the same check through ka2a.ParseBrokerRevocationLists. Every certificate of the chain is checked. ka2a refuses a stale list and a list that is not complete, and a reload cannot remove the check. doctor warns when the next update of the lists is less than 24 hours away, and a metric gives that time. ka2a never fetches a list and does not use OCSP.

Broker CA rotation in two reloads

First reload a bundle of the old and the new CA. Then the brokers change their certificates. Then reload only the new CA. A real-broker test did these steps, and no exchange failed.

A refused login does not fail the work

When the broker refuses the SASL authentication of a produce, the request stays queued. After the reload of the correct password, the same request completes. A real-broker test shows this. With broker re-authentication on, open connections use the new password at their next re-authentication.

CLI · keys generate, catalog validate

Rotate a signing key and check the catalog

A key rotation changes two keys: the new key, and the end and the rotation link of the old key. The validator finds a catalog that adds the new key without the link.

  • A node accepts the first catalog, but two keys can sign for one endpoint at the same time. catalog validate reports this as a problem.
  • The commands need no broker and no store. They only read files, except keys generate, which writes the new key file and never replaces a file.
terminalReal run · synthetic endpoints
# Make the next key of planner. The private key never appears in the output.
$ ka2a keys generate --endpoint demo.example/planner \
    --key-id planner-2026-10 --not-before 2026-10-01T00:00:00Z \
    --out keys/planner-2026-10.key.json
Wrote the private signing key planner-2026-10 to /srv/ka2a-demo/keys/planner-2026-10.key.json (mode 0600). Add this entry to the "keys" list of the catalog:
{
  "key_id": "planner-2026-10",
  "algorithm": "ed25519",
  "public_key": "XbSECvRVfJI5ofb9WbO81as8buZkdkf-GhcQrFsZ6gU",
  "domain": "demo.example",
  "endpoint": "planner",
  "not_before": "2026-10-01T00:00:00Z"
}

# Not a ka2a command: a copy of the catalog gets the printed entry, without a rotation link.
$ ka2a catalog validate catalog-unlinked.json
Catalog /srv/ka2a-demo/catalog-unlinked.json: 2 endpoints, 3 keys, 2 grants; digest sha256:313c9aeb…9048658a.
SEVERITY  CLASS                 PATH       FINDING
problem   overlapping_validity  $.keys[2]  keys planner-2026-09 and planner-2026-10 of demo.example/planner are valid at the same time, and neither replaces the other
1 problem, 0 warnings. A node accepts this catalog.
# exit status 3: the checks found a problem.

# A second copy also gives the old key an end (not_after) and a link to the new key (replaced_by).
$ ka2a catalog validate catalog-next.json
Catalog /srv/ka2a-demo/catalog-next.json: 2 endpoints, 3 keys, 2 grants; digest sha256:ef930caf…07d6fc4e.
0 problems, 0 warnings. A node accepts this catalog.

$ ka2a catalog digest catalog-next.json
sha256:ef930caf158270e5db083ba33eb9cd195c785742cdcc37eb0efe5a4e07d6fc4e
Listing 12 Make the next key of planner and check two copies of the catalog. The catalog copies are made by a script, not by ka2a. An ellipsis (…) marks shortened digests.

CLI · topics acl-plan

Print the minimal broker ACLs of an endpoint

The plan gives Read on the own mailbox topic and on the consumer group, DescribeConfigs for the mailbox check, and Write on the mailbox of each peer. Responses go to the mailbox of the requester, so both directions need Write.

  • With mutual TLS, the principal comes from the client certificate. With SASL, it is User: and the user name.
  • A denied record is not sent again. After you add the ACL, send the request again. See the runbook Missing broker ACL.
terminalReal run · synthetic endpoints
$ ka2a topics acl-plan --catalog catalog.json --state-dir planner \
    --principal User:planner
Minimal ACLs of demo.example/planner for principal User:planner. ka2a prints them; it never applies them.
RESOURCE  NAME                               OPERATION        PURPOSE
topic     ka2a.v1.mailbox.c5b7df51…6c16230e  Read             consume the own mailbox (Read implies Describe)
topic     ka2a.v1.mailbox.c5b7df51…6c16230e  DescribeConfigs  check the mailbox against the declared profile
group     ka2a.owner.demo.example.planner    Read             join the consumer group and commit offsets (Read implies Describe)
topic     ka2a.v1.mailbox.79e3dba7…9c370722  Write            send requests and responses to demo.example/reviewer (Write implies Describe)

Commands (set BOOTSTRAP and ADMIN_CONFIG for an administrative principal):
kafka-acls.sh --bootstrap-server "$BOOTSTRAP" --command-config "$ADMIN_CONFIG" --add --allow-principal 'User:planner' --allow-host '*' --operation Read --topic 'ka2a.v1.mailbox.c5b7df51…6c16230e' --resource-pattern-type literal
…
Listing 13 The ACL plan of planner. An ellipsis (…) marks shortened topic names, and the table columns are aligned again after the cut. The last ellipsis marks three cut commands and the cut notes.

What the tests show. Integration tests run ka2a against real brokers that require TLS with SASL, or mutual TLS with ACLs. They rotate a client certificate, a SASL password and the broker CA, and no exchange fails. Load profiles with a leader loss passed on three brokers with mutual TLS and only the planned ACLs, on a shared, contended host; they are not a capacity claim. A real broker with a revoked certificate was refused, and accepted again after it renewed its certificate. No test rotates the client CA of a mutual TLS listener, no real broker has a revoked intermediate CA, no real-broker test reads the new revocation status fields, and the reload check does not contact brokers that only the cluster metadata names. Read production deployment for the setup.

Built from specifications

Every behavior starts as a requirement

ka2a uses spec-driven development with OpenSpec. Each requirement has a stable ID and scenarios. Each scenario has a row in an acceptance ledger, and each automated row names the test that exercises it.

acceptance scenarios with automated tests
60 of 65
scenarios with partial evidence and a named gap
5
requirements with stable IDs
25
frozen wire test vectors
29
Go test functions in the module
576

One requirement, traced to its tests

specs/communication-contract/spec.mdRequirement
### Requirement: Frozen records and identity conflicts

Requirement-ID: KA2A-R1-004

After durable local acceptance the exact key, headers,
value, identity, signature and expiry SHALL be retained
across retries. Reusing an operation identity with a
changed logical request SHALL conflict rather than
mutate the original.

#### Scenario: Changed request
- GIVEN an operation key already names one
  target and request
- WHEN a caller submits different params
  with that key
- THEN it receives an identity-conflict error
  without replacing or sending the old operation
Listing 14 The requirement and one of its scenarios.
acceptance.mdLedger row
| ID | Requirement | Task | Scenario | Status |
| K-004-B | KA2A-R1-004 | K2 |
  Changed request | automated |

Evidence:
  Codec:   TestIdentityConflict
  Store:   TestSameOperationKeyWithChangedRequestConflicts
  Runtime: TestOperationKeyReuseAndConflict
           (same key and request returns the original
           without new bytes; changed body, class,
           horizon or operation conflicts; one record
           on the broker, one dispatch)
Listing 15 The ledger row of that scenario, with its evidence.
internal/runtime/sender_test.goTest
// Acceptance: K-004-B
// KH-08 (sender): a caller-supplied operation key names one logical
// request. The same request returns the original operation without new
// bytes; a changed request conflicts and neither replaces nor sends the
// original again.
func TestOperationKeyReuseAndConflict(t *testing.T) {
	leak.Check(t)
	ctx := testContext(t, time.Minute)
	f := newFixture(t, []string{"client", "server"})
	…
}
Listing 16 The runtime test that the row names.

A check script fails the build when a ledger row names a test that does not exist, or when a test names a ledger row that does not exist.

Real Kafka, real SQLite

Integration tests start a digest-pinned Apache Kafka image in Docker. They kill and restart processes and brokers, then check the stores and the offsets.

Properties and fuzzing

Rapid property tests drive random histories through the store and the checkpoint coordinator. Fuzz targets attack the codec, the catalog and the snapshot reader.

Boundaries as tests

Architecture tests keep internal types out of the public API and keep notarizing out of the dependency graph.

Release artifacts and qualification runs

Reproducible release artifacts

make release-artifacts builds one pure-Go binary per supported platform, with one SPDX 2.3 SBOM per binary, the notices, an evidence summary and SHA256SUMS. make release-check builds everything twice from git archive and compares each byte. The tooling publishes nothing and signs nothing.

65 minutes on three secured brokers

A soak ran for 65 minutes at 20 operations per second on three brokers with TLS and SCRAM-SHA-512 and the qualified durability profile. All 73,689 operations completed, the handler ran once for each operation, and retention reached a steady state. The host was shared and contended, so the numbers are not a capacity claim.

Model traces replayed against Go

18 traces of the Quint model of the receive path (1,058 steps) replay against the real checkpoint coordinator and the real SQLite store in make verify. This shows agreement on these traces only. It is not a proof, and it says nothing about concurrency or liveness.

Limits of the fault campaign. A protected evaluator runs the fault histories from a pinned reference, and it refuses a changed or shadowed pinned test. The same agent can still change a target and its pin, so the evaluator has no real authority boundary yet. The failure records are local files of the user who runs the evaluator.

Status

What works today, and what does not

This is a developer preview. The list below describes the reviewed revision in source-state.json.

Works in this revision

  • Go API: Open, Run, WaitReady, Close, a typed client and a server with handlers.
  • A2A 1.0 SendMessage, GetTask, ListTasks and CancelTask, without streaming.
  • KA2A-WIRE-1 records with Ed25519 signatures and frozen test vectors.
  • SQLite outbox, receive journal, dispatch, held work, retention and restore detection.
  • CLI for setup, sending, inspection, backup and restore.
  • Broker TLS with SASL SCRAM-SHA-512 or PLAIN, and mutual TLS client certificates, tested against real brokers.
  • Classified broker refusals and ACL denials, and a minimal ACL plan from the catalog.
  • Reload of the broker TLS and SASL files without a restart (SIGHUP or Node.ReloadCredentials), checked against each bootstrap broker, and the catalog drift status.
  • Local revocation lists for broker certificates (--tls-crl, or ka2a.ParseBrokerRevocationLists in the Go API), tested against a real broker with a revoked certificate. doctor, the node status and a metric show the next update of the lists.
  • Key generation and catalog checks: keys generate, keys public, catalog validate, catalog digest.
  • A read-only loopback UI without JavaScript, and Prometheus metrics with fixed labels.
  • Reproducible release artifacts with SBOMs and checksums. Nothing is published.
  • A known-vulnerability scan (make vulncheck). With an offline snapshot of the Go vulnerability database of 2026-10-01, it found nothing in the Go code. The advisories of adm-zip in the development-only model tooling are fixed.
  • A versioning and stability policy and a changelog. No version is released yet.
  • A purge leaves zeros in the freed pages of the database file (SQLite secure_delete). Fuzz targets for the parsers of operator input ran without a finding.
  • A 65-minute soak on three secured brokers, load runs on three brokers with mutual TLS and minimal ACLs, and model traces replayed against the Go code.
  • Operator documentation checked against the binary. A test validates its JSON examples, including a complete catalog.
  • Fault histories on a real broker, a fresh-container run of the unit and property tests, and a separate consumer module that uses only the public API.

Not in this release

  • Streaming and push notifications.
  • Dynamic discovery and federation.
  • Kubernetes integration and transparent failover.
  • Other brokers and plugins.
  • A reload of the catalog or the signing key. Both need a restart.
  • Commands that edit a catalog.

Not yet verified or blocked

  • A release: no tag, release or upload exists. Publishing and signing are manual decisions of the owner.
  • Public installation: the source repositories are private during the preview.
  • An independent evaluator: the protected evaluator has no real authority boundary.
  • Tests on linux/arm64 and macOS. Of the four shipping platforms, three only cross-compile.
  • A vulnerability scan with a newer database. The scan covers Go code only, and it is not a security review.
  • An online backup, an HTTP health endpoint and the erasure of one message. A backup needs a stopped node.
  • OCSP, a revoked intermediate CA on a real broker, a rotation of the client CA of a mutual TLS listener, and a real-broker run of the revocation status fields and the doctor warning.
  • A fresh-environment run of the integration tests, the fault campaign, lint, the race detector, fuzzing and the benchmarks.
  • A leader loss during a soak, and a soak on a quiet host.
  • Formal proofs. The Quint model has bounded runs and replayed traces only.

Also in this family

Notarizing: review a change with its evidence in view

A local-first workspace for spec-driven review. It reads OpenSpec changes, shows their premises as a graph and links each requirement to the check reports you import. It can receive evidence over ka2a, and it works without it.