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.
Statustells whether a broker answered the mailbox check, and the live broker state.HoldOnAmbiguityholds a request for an operator when a crash leaves its effect unknown.- Unsupported operations answer with a typed A2A error.
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)
}
reviewer with one handler.