a2aserver

package module
v2.13.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 9, 2026 License: Apache-2.0 Imports: 20 Imported by: 0

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

Examples

Constants

This section is empty.

Variables

View Source
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.

View Source
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

type AgentCardProvider interface {
	AgentCard(r *http.Request) (*a2a.AgentCard, error)
}

AgentCardProvider returns the agent card to serve at /.well-known/agent-card.json (and the legacy /.well-known/agent.json).

type Authenticator

type Authenticator interface {
	Authenticate(r *http.Request) error
}

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

type HealthChecker interface {
	Check(ctx context.Context) error
}

HealthChecker performs a named health check. Implementations should return nil when healthy and a non-nil error describing the problem otherwise.

type HealthCheckerFunc

type HealthCheckerFunc func(ctx context.Context) error

HealthCheckerFunc adapts an ordinary function to the HealthChecker interface.

func (HealthCheckerFunc) Check

func (f HealthCheckerFunc) Check(ctx context.Context) error

Check calls f(ctx).

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

func (s *InMemoryTaskStore) List(contextID string, limit, offset int) ([]*a2a.Task, error)

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 "".

func (*InMemoryTaskStore) Query added in v2.8.0

func (s *InMemoryTaskStore) Query(q TaskQuery) (TaskPage, error)

Query implements TaskQuerier.

func (*InMemoryTaskStore) SetState

func (s *InMemoryTaskStore) SetState(taskID string, state a2a.TaskState, msg *a2a.Message) error

SetState transitions the task to a new state with an optional status message.

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

func WithCard(card *a2a.AgentCard) Option

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

func WithConversationTTL(d time.Duration) Option

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

func WithIdleTimeout(d time.Duration) Option

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

func WithMaxBlockingWait(d time.Duration) Option

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

func WithMaxBodySize(n int64) Option

WithMaxBodySize sets the maximum allowed request body size in bytes. Default: 10 MB.

func WithPort

func WithPort(port int) Option

WithPort sets the TCP port for ListenAndServe.

func WithReadTimeout

func WithReadTimeout(d time.Duration) Option

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

func WithTaskOwner(owner OwnerFunc) Option

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

func WithTaskStore(store TaskStore) Option

WithTaskStore sets a custom task store. Defaults to an in-memory store.

func WithTaskTTL

func WithTaskTTL(d time.Duration) Option

WithTaskTTL sets how long completed/failed/canceled tasks are retained before automatic eviction. Default: 1 hour. Set to 0 to disable eviction.

func WithWriteTimeout

func WithWriteTimeout(d time.Duration) Option

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

type OwnerFunc func(r *http.Request) string

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) Handler

func (s *Server) Handler() http.Handler

Handler returns an http.Handler implementing the A2A protocol.

func (*Server) ListenAndServe

func (s *Server) ListenAndServe() error

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.

func (*Server) Serve

func (s *Server) Serve(ln net.Listener) error

Serve starts the HTTP server on the given listener. See ListenAndServe for the rationale behind WriteTimeout: 0.

func (*Server) Shutdown

func (s *Server) Shutdown(ctx context.Context) error

Shutdown gracefully shuts down the server: stops the eviction goroutine, drains HTTP requests, cancels in-flight tasks, and closes all conversations.

type StaticCard

type StaticCard struct {
	Card a2a.AgentCard
}

StaticCard is an AgentCardProvider that always returns the same card.

func (*StaticCard) AgentCard

func (s *StaticCard) AgentCard(*http.Request) (*a2a.AgentCard, error)

AgentCard returns the static 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.

func (TaskEvent) IsFinal added in v2.8.0

func (e TaskEvent) IsFinal() bool

IsFinal reports whether the event ends a task's stream: a status update to a terminal or interrupted state.

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

type TaskQuerier interface {
	Query(q TaskQuery) (TaskPage, error)
}

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL