Documentation
¶
Overview ¶
Package a2aserver provides a standalone A2A-protocol HTTP server that can be backed by any conversation implementation satisfying the interfaces defined here. It imports only runtime/ and has no dependency on sdk/.
Index ¶
- Variables
- type AgentCardProvider
- type Authenticator
- type Conversation
- type ConversationOpener
- type EventKind
- type HealthChecker
- type HealthCheckerFunc
- type InMemoryTaskStore
- func (s *InMemoryTaskStore) AddArtifacts(taskID string, artifacts []a2a.Artifact) error
- func (s *InMemoryTaskStore) Cancel(taskID string) error
- func (s *InMemoryTaskStore) Create(taskID, contextID string) (*a2a.Task, error)
- func (s *InMemoryTaskStore) CreateOwned(taskID, contextID, owner string) (*a2a.Task, error)
- func (s *InMemoryTaskStore) EvictTerminal(cutoff time.Time) []string
- func (s *InMemoryTaskStore) Get(taskID string) (*a2a.Task, error)
- func (s *InMemoryTaskStore) List(contextID string, limit, offset int) ([]*a2a.Task, error)
- func (s *InMemoryTaskStore) Owner(taskID string) (string, error)
- func (s *InMemoryTaskStore) Query(q TaskQuery) (TaskPage, error)
- func (s *InMemoryTaskStore) SetState(taskID string, state a2a.TaskState, msg *a2a.Message) error
- type MessageHandler
- type MessageRequest
- type Option
- func WithAuthenticator(auth Authenticator) Option
- func WithCard(card *a2a.AgentCard) Option
- func WithCardProvider(p AgentCardProvider) Option
- func WithConversationTTL(d time.Duration) Option
- func WithHealthCheck(name string, checker HealthChecker) Option
- func WithIdleTimeout(d time.Duration) Option
- func WithMaxBlockingWait(d time.Duration) Option
- func WithMaxBodySize(n int64) Option
- func WithPort(port int) Option
- func WithReadTimeout(d time.Duration) Option
- func WithTaskCanceler(c TaskCanceler) Option
- func WithTaskEventBus(bus TaskEventBus) Option
- func WithTaskOwner(owner OwnerFunc) Option
- func WithTaskStore(store TaskStore) Option
- func WithTaskTTL(d time.Duration) Option
- func WithWriteTimeout(d time.Duration) Option
- type OwnedTaskStore
- type OwnerFunc
- type PendingClientToolInfo
- type ResumableConversation
- type SendResult
- type Server
- type StaticCard
- type StreamEvent
- type StreamingConversation
- type TaskCanceler
- type TaskEvent
- type TaskEventBus
- type TaskPage
- type TaskQuerier
- type TaskQuery
- type TaskStore
- type ToolResult
- type ToolResultHandler
- type ToolResultRequest
Examples ¶
Constants ¶
This section is empty.
Variables ¶
var ( ErrTaskNotFound = errors.New("a2a: task not found") ErrTaskAlreadyExists = errors.New("a2a: task already exists") ErrInvalidTransition = errors.New("a2a: invalid state transition") ErrTaskTerminal = errors.New("a2a: task is in a terminal state") )
Task store errors.
var ErrTooManySubscribers = errors.New("a2a: too many subscribers")
ErrTooManySubscribers is returned when a task has reached its subscriber limit.
Functions ¶
This section is empty.
Types ¶
type AgentCardProvider ¶
AgentCardProvider returns the agent card to serve at /.well-known/agent-card.json (and the legacy /.well-known/agent.json).
type Authenticator ¶
Authenticator validates incoming requests. Return a non-nil error to reject.
type Conversation ¶
type Conversation interface {
Send(ctx context.Context, message any) (SendResult, error)
Close() error
}
Conversation is the non-streaming conversation interface the server uses.
type ConversationOpener ¶
type ConversationOpener func(contextID string) (Conversation, error)
ConversationOpener creates or retrieves a conversation for a given context ID.
type EventKind ¶
type EventKind int
EventKind discriminates the payload of a StreamEvent.
const ( // EventText indicates the event contains text content. EventText EventKind = iota // EventToolCall indicates a tool call (suppressed by the server). EventToolCall // EventMedia indicates the event contains media content. EventMedia // EventDone indicates the stream is complete. EventDone // EventClientTool indicates a client tool request awaiting fulfillment. EventClientTool // EventPending indicates the turn is paused on something the server cannot // see — a human approving a tool call, most often. Text carries the reason. // // It exists because "is this turn waiting?" cannot be inferred: a pause on // a server-side approval looks exactly like a finished turn from outside, // and reporting it as completed would be a control that did not run // reporting clean. A handler that never pauses never emits it. EventPending )
type HealthChecker ¶
HealthChecker performs a named health check. Implementations should return nil when healthy and a non-nil error describing the problem otherwise.
type HealthCheckerFunc ¶
HealthCheckerFunc adapts an ordinary function to the HealthChecker interface.
type InMemoryTaskStore ¶
type InMemoryTaskStore struct {
// contains filtered or unexported fields
}
InMemoryTaskStore is a concurrency-safe, in-memory implementation of TaskStore, TaskQuerier and OwnedTaskStore.
Example ¶
ExampleInMemoryTaskStore demonstrates the task lifecycle: create a task, transition it through states, and read it back. InMemoryTaskStore needs no external services, making it a good default for local development and tests; production deployments typically supply a persistent TaskStore.
package main
import (
"fmt"
a2aserver "github.com/AltairaLabs/PromptKit/server/a2a/v2"
)
func main() {
store := a2aserver.NewInMemoryTaskStore()
task, err := store.Create("task-1", "ctx-1")
if err != nil {
fmt.Println("create error:", err)
return
}
fmt.Println("created:", task.Status.State)
if err := store.SetState("task-1", "working", nil); err != nil {
fmt.Println("transition error:", err)
return
}
got, err := store.Get("task-1")
if err != nil {
fmt.Println("get error:", err)
return
}
fmt.Println("state:", got.Status.State)
}
Output: created: submitted state: working
func NewInMemoryTaskStore ¶
func NewInMemoryTaskStore() *InMemoryTaskStore
NewInMemoryTaskStore creates a new InMemoryTaskStore.
func (*InMemoryTaskStore) AddArtifacts ¶
func (s *InMemoryTaskStore) AddArtifacts(taskID string, artifacts []a2a.Artifact) error
AddArtifacts appends artifacts to a task.
func (*InMemoryTaskStore) Cancel ¶
func (s *InMemoryTaskStore) Cancel(taskID string) error
Cancel transitions the task to the canceled state from any non-terminal state.
func (*InMemoryTaskStore) Create ¶
func (s *InMemoryTaskStore) Create(taskID, contextID string) (*a2a.Task, error)
Create initializes a new task in the submitted state.
func (*InMemoryTaskStore) CreateOwned ¶ added in v2.8.0
func (s *InMemoryTaskStore) CreateOwned(taskID, contextID, owner string) (*a2a.Task, error)
CreateOwned implements OwnedTaskStore.
func (*InMemoryTaskStore) EvictTerminal ¶
func (s *InMemoryTaskStore) EvictTerminal(cutoff time.Time) []string
EvictTerminal removes tasks in a terminal state whose last status timestamp is older than cutoff. It returns the IDs of evicted tasks.
func (*InMemoryTaskStore) Get ¶
func (s *InMemoryTaskStore) Get(taskID string) (*a2a.Task, error)
Get retrieves a deep copy of a task by ID. The returned task is safe to read/modify without holding the store lock.
func (*InMemoryTaskStore) List ¶
List returns deep copies of tasks matching the given contextID with pagination. If contextID is empty, all tasks are returned. Results are ordered most recently updated first, ID breaking ties, for deterministic pagination. Offset and limit control pagination.
func (*InMemoryTaskStore) Owner ¶ added in v2.8.0
func (s *InMemoryTaskStore) Owner(taskID string) (string, error)
Owner implements OwnedTaskStore. A task created without an owner has "".
type MessageHandler ¶ added in v2.5.0
type MessageHandler interface {
Handle(ctx context.Context, req MessageRequest) <-chan StreamEvent
}
MessageHandler turns one inbound A2A message into a stream of events.
Handle is called once per message and the server keeps no state between calls. ctx is the HTTP request's context, so whatever the embedder's middleware put on it (caller identity, tenant, trace) is readable here, which is the thing ConversationOpener cannot offer.
The returned channel must be closed when the turn is over. A closing channel that sent no EventDone is treated as a completed turn.
type MessageRequest ¶ added in v2.5.0
type MessageRequest struct {
// ContextID is client-supplied. The server attaches no meaning to it: the
// embedder decides whether it identifies a session, and whether two callers
// presenting the same one share anything.
ContextID string
// TaskID is the server-assigned id of the task this message created.
TaskID string
// Message is the inbound A2A message.
Message a2a.Message
}
MessageRequest is one inbound message, as the handler sees it.
type Option ¶
type Option func(*Server)
Option configures a Server.
func WithAuthenticator ¶
func WithAuthenticator(auth Authenticator) Option
WithAuthenticator sets an authenticator for incoming requests.
func WithCard ¶
WithCard sets the agent card served at /.well-known/agent-card.json.
The server completes the card's JSON-RPC interface declarations for the protocol versions it speaks (see servedCard); declare SecuritySchemes and SecurityRequirements on it when WithAuthenticator is in use, so callers can discover how to authenticate.
func WithCardProvider ¶
func WithCardProvider(p AgentCardProvider) Option
WithCardProvider sets a dynamic agent card provider.
func WithConversationTTL ¶
WithConversationTTL sets how long idle conversations are retained before automatic eviction. A conversation is considered idle when its last-use timestamp exceeds this duration. Default: 1 hour. Set to 0 to disable.
func WithHealthCheck ¶
func WithHealthCheck(name string, checker HealthChecker) Option
WithHealthCheck registers a named health checker that is evaluated by the /readyz endpoint. Multiple checkers can be registered; each is reported individually in the response body.
func WithIdleTimeout ¶
WithIdleTimeout sets the maximum amount of time to wait for the next request when keep-alives are enabled. Default: 120s.
func WithMaxBlockingWait ¶ added in v2.8.0
WithMaxBlockingWait caps how long a blocking SendMessage holds its request open. When the turn has not finished or been interrupted by then, the server answers with the task in its current (working) state, and the turn runs on: the caller polls GetTask or subscribes for the rest. Default: 0, no cap — a 1.0 SendMessage waits for the turn, as the spec requires, until the turn ends or the caller disconnects.
Set it below the timeout of whatever sits in front of the server (a load balancer's idle timeout, a proxy's read timeout), so the caller gets a task to follow rather than a gateway error.
func WithMaxBodySize ¶
WithMaxBodySize sets the maximum allowed request body size in bytes. Default: 10 MB.
func WithReadTimeout ¶
WithReadTimeout sets the maximum duration for reading the entire request. Default: 30s.
func WithTaskCanceler ¶ added in v2.8.0
func WithTaskCanceler(c TaskCanceler) Option
WithTaskCanceler sets how CancelTask reaches the instance running a task. Default: in-process only.
func WithTaskEventBus ¶ added in v2.8.0
func WithTaskEventBus(bus TaskEventBus) Option
WithTaskEventBus sets how task updates reach SubscribeToTask callers. Default: in-process only.
func WithTaskOwner ¶ added in v2.8.0
WithTaskOwner scopes every task to the caller that created it. owner identifies the caller of each request; GetTask, CancelTask, ListTasks and SubscribeToTask then see only the caller's own tasks, and ListTasks without a contextId lists them all. With NewServer, a message into a conversation another caller opened is refused; with NewStatelessServer the handler owns contexts and decides.
The task store must implement OwnedTaskStore (the default in-memory store does); NewServer panics otherwise, since serving with scoping silently off would be worse than not starting. A request whose owner is empty is refused.
func WithTaskStore ¶
WithTaskStore sets a custom task store. Defaults to an in-memory store.
func WithTaskTTL ¶
WithTaskTTL sets how long completed/failed/canceled tasks are retained before automatic eviction. Default: 1 hour. Set to 0 to disable eviction.
func WithWriteTimeout ¶
WithWriteTimeout sets the maximum duration before timing out writes of the response. Default: 60s.
type OwnedTaskStore ¶ added in v2.8.0
type OwnedTaskStore interface {
TaskStore
// CreateOwned creates a task, as Create does, recording owner as its
// creator.
CreateOwned(taskID, contextID, owner string) (*a2a.Task, error)
// Owner returns the owner recorded for a task, or ErrTaskNotFound.
Owner(taskID string) (string, error)
}
OwnedTaskStore is a TaskStore that records which caller created each task. WithTaskOwner requires one; InMemoryTaskStore is one.
A store that also implements TaskQuerier must honor TaskQuery.Owner.
type OwnerFunc ¶ added in v2.8.0
OwnerFunc identifies the caller of a request, typically from what the host's authentication middleware put on the request context. It must return a non-empty identity for every caller allowed to use the server.
type PendingClientToolInfo ¶
type PendingClientToolInfo struct {
CallID string `json:"call_id"`
ToolName string `json:"tool_name"`
Args map[string]any `json:"args,omitempty"`
ConsentMsg string `json:"consent_message,omitempty"`
}
PendingClientToolInfo describes a client-side tool call awaiting fulfillment.
type ResumableConversation ¶
type ResumableConversation interface {
Conversation
SendToolResult(callID string, result any) error
RejectClientTool(callID string, reason string)
Resume(ctx context.Context) (SendResult, error)
ResumeStream(ctx context.Context) <-chan StreamEvent
}
ResumableConversation extends Conversation with the ability to submit client-side tool results and resume pipeline execution.
type SendResult ¶
type SendResult interface {
// HasPendingTools reports whether there are tools awaiting approval
// (HITL or client-side). When true the task transitions to input_required.
HasPendingTools() bool
// HasPendingClientTools reports whether there are client-side tools
// awaiting fulfillment by the caller.
HasPendingClientTools() bool
// PendingClientTools returns metadata for each pending client tool.
PendingClientTools() []PendingClientToolInfo
// Parts returns the content parts of the response.
Parts() []types.ContentPart
// Text returns the text content of the response as a fallback
// when Parts() is empty.
Text() string
}
SendResult is what the server needs from a completed conversation turn.
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server is an HTTP server that exposes a Conversation as an A2A-compliant JSON-RPC endpoint.
func NewServer ¶
func NewServer(opener ConversationOpener, opts ...Option) *Server
NewServer creates a new A2A server that OWNS its conversations: it opens one per context id through the supplied opener, caches it, and reuses it when that id returns. Suits an embedder running A2A and the runtime in one process.
For an embedder that owns conversations itself (because the runtime lives elsewhere, or it already tracks sessions), see NewStatelessServer.
func NewStatelessServer ¶ added in v2.5.0
func NewStatelessServer(h MessageHandler, opts ...Option) *Server
NewStatelessServer creates a server that holds no conversations.
It is NewServer's sibling: same protocol, same options, but each message goes to the handler with its request context and nothing is kept between calls. Use it when the embedder owns conversations — because it already tracks sessions, because the runtime is in another process, or because the server needs to scale horizontally with only the task store shared.
func (*Server) ListenAndServe ¶
ListenAndServe starts the HTTP server on the configured port.
WriteTimeout is set to 0 (disabled) because SSE streaming endpoints (message/stream, tasks/subscribe) hold the connection open indefinitely. A non-zero WriteTimeout would kill long-lived SSE connections. Non-streaming endpoints rely on the request context deadline for timeout enforcement.
type StaticCard ¶
StaticCard is an AgentCardProvider that always returns the same card.
type StreamEvent ¶
type StreamEvent struct {
Kind EventKind
Text string
Media *types.MediaContent
ClientTool *PendingClientToolInfo
Error error
}
StreamEvent is a single event on a streaming channel.
type StreamingConversation ¶
type StreamingConversation interface {
Conversation
Stream(ctx context.Context, message any) <-chan StreamEvent
}
StreamingConversation extends Conversation with streaming support.
type TaskCanceler ¶ added in v2.8.0
type TaskCanceler interface {
// Cancel asks whichever instance is running taskID to stop it. The task
// is already marked canceled in the store when this is called.
Cancel(ctx context.Context, taskID string) error
// Listen registers the function that stops a turn running on this
// instance; the implementation calls it for every cancel request, from
// any instance, and it is a no-op for tasks this instance is not
// running. The server calls Listen once, when it is created, and stop
// when it shuts down.
Listen(cancelLocal func(taskID string)) (stop func())
}
TaskCanceler stops a task's in-flight turn wherever it runs.
The default is in-process: CancelTask stops a turn only on the instance running it. A host running several replicas supplies a shared implementation with WithTaskCanceler, or pins each task's callers to one replica.
type TaskEvent ¶ added in v2.8.0
type TaskEvent struct {
StatusUpdate *a2a.TaskStatusUpdateEvent `json:"statusUpdate,omitempty"`
ArtifactUpdate *a2a.TaskArtifactUpdateEvent `json:"artifactUpdate,omitempty"`
}
TaskEvent is one update to a task, as SubscribeToTask callers receive it. Exactly one field is set.
It is version-neutral because a subscriber may speak a different protocol version, and carries a different JSON-RPC id, than the caller whose turn produced the event, so each subscriber encodes it for itself.
type TaskEventBus ¶ added in v2.8.0
type TaskEventBus interface {
// Publish delivers evt to the task's subscribers, wherever they are.
Publish(ctx context.Context, taskID string, evt TaskEvent) error
// Subscribe returns the task's events from now on. The server stops
// reading after a final event (TaskEvent.IsFinal); the implementation
// must stop delivering, and release the subscription, when ctx ends.
Subscribe(ctx context.Context, taskID string) (<-chan TaskEvent, error)
}
TaskEventBus carries task updates to SubscribeToTask callers.
The default is in-process: a subscriber sees updates only from turns this server instance runs. A host running several replicas behind one task store supplies a shared implementation (Redis pub/sub, NATS, ...) with WithTaskEventBus, or pins each task's callers to one replica.
type TaskPage ¶ added in v2.8.0
type TaskPage struct {
// Tasks are ordered most recently updated first (A2A 1.0 §3.1.4).
Tasks []*a2a.Task
// Total is the number of tasks matching the query across all pages.
Total int
}
TaskPage is one page of a TaskQuery's result.
type TaskQuerier ¶ added in v2.8.0
TaskQuerier is optionally implemented by a TaskStore that can filter, order and page tasks itself. Without it the server pages through List and does the filtering and ordering in memory, which is correct but reads every task in the context on each call.
type TaskQuery ¶ added in v2.8.0
type TaskQuery struct {
// Owner, when set, restricts the query to the tasks that caller created
// (see OwnedTaskStore).
Owner string
// ContextID, when set, restricts the query to one context.
ContextID string
// Status, when set, keeps only tasks in that state.
Status *a2a.TaskState
// StatusAfter, when set, keeps only tasks whose status changed at or after it.
StatusAfter *time.Time
// Limit and Offset select the page within the ordered result.
Limit int
Offset int
}
TaskQuery selects a page of tasks for ListTasks.
type TaskStore ¶
type TaskStore interface {
Create(taskID, contextID string) (*a2a.Task, error)
Get(taskID string) (*a2a.Task, error)
SetState(taskID string, state a2a.TaskState, msg *a2a.Message) error
AddArtifacts(taskID string, artifacts []a2a.Artifact) error
Cancel(taskID string) error
List(contextID string, limit, offset int) ([]*a2a.Task, error)
// EvictTerminal removes tasks in a terminal state whose last status
// timestamp is older than the given cutoff time. It returns the IDs
// of evicted tasks so callers can clean up associated resources.
EvictTerminal(olderThan time.Time) []string
}
TaskStore defines the interface for task persistence and lifecycle management.
type ToolResult ¶ added in v2.5.0
type ToolResult struct {
CallID string
Result any
// Rejected is true when the caller declined the tool. Reason carries their
// explanation, if any.
Rejected bool
Reason string
}
ToolResult is one fulfilled (or refused) client-side tool call.
type ToolResultHandler ¶ added in v2.5.0
type ToolResultHandler interface {
HandleToolResult(ctx context.Context, req ToolResultRequest) <-chan StreamEvent
}
ToolResultHandler is the optional client-tool half of MessageHandler. A handler that implements it can receive the results of client-side tool calls it asked for earlier and continue the turn.
Without it, a stateless server rejects tool-result messages rather than pretending to resume something it is not holding.
type ToolResultRequest ¶ added in v2.5.0
type ToolResultRequest struct {
ContextID string
TaskID string
// Results are the fulfilled tool calls, keyed by the call id the handler
// supplied in its [EventClientTool] events.
Results []ToolResult
}
ToolResultRequest carries client-tool results back to the handler.