agents

package
v1.55.0 Latest Latest
Warning

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

Go to latest
Published: Mar 20, 2026 License: MIT Imports: 20 Imported by: 0

Documentation

Index

Constants

View Source
const (
	HeaderRequestID      = "X-Agent-Request-ID"      // Unique request identifier
	HeaderCancelCallback = "X-Agent-Cancel-Callback" // URL to call for cancellation
)

Headers for request cancellation in agent-to-agent communication

View Source
const (
	HeaderCallChain   = "X-Agent-Call-Chain"  // JSON array of agent IDs in call path
	HeaderCallDepth   = "X-Agent-Call-Depth"  // Current depth (integer)
	HeaderOriginAgent = "X-Agent-Origin"      // Original requesting agent
	HeaderCallID      = "X-Agent-Call-ID"     // Unique ID for this call tree
	HeaderTTL         = "X-Agent-TTL"         // Time-to-live in seconds
	HeaderStartTime   = "X-Agent-Start-Time"  // Unix timestamp (nanoseconds) when call chain started
	HeaderCallerID    = "X-Agent-Caller-ID"   // Direct caller agent ID (immediate parent)
	HeaderCallerName  = "X-Agent-Caller-Name" // Direct caller agent name (human-readable)
)

Headers for call chain tracking between agents

Variables

This section is empty.

Functions

func AddCallContextHeaders

func AddCallContextHeaders(req *http.Request, ctx *CallContext, currentAgentID string)

AddCallContextHeaders adds call context headers to an outgoing request

func AddCancellationHeaders

func AddCancellationHeaders(req *http.Request, requestID, cancelCallbackURL string)

AddCancellationHeaders adds cancellation-related headers to an outgoing request

func ContextWithCancellation

func ContextWithCancellation(parent context.Context, tracker *RequestTracker, requestID string) context.Context

ContextWithCancellation creates a child context that respects both parent cancellation and the request tracker's cancellation

func ExtractCancellationContext

func ExtractCancellationContext(r *http.Request) (requestID, cancelCallback string)

ExtractCancellationContext extracts cancellation info from incoming request headers

func GenerateRequestID

func GenerateRequestID() string

GenerateRequestID creates a unique request ID Format: req_{timestamp_base32}_{random_6chars}

func IsCancelled

func IsCancelled(tracker *RequestTracker, requestID string) bool

IsCancelled checks if a request has been cancelled

Types

type ActiveRequest

type ActiveRequest struct {
	RequestID       string        `json:"request_id"`
	ThreadID        string        `json:"thread_id,omitempty"`
	Status          RequestStatus `json:"status"`
	StartTime       time.Time     `json:"start_time"`
	CancelledAt     *time.Time    `json:"cancelled_at,omitempty"`
	CompletedAt     *time.Time    `json:"completed_at,omitempty"`
	CancelCallbacks []string      `json:"-"` // URLs to call when cancelled (for agent propagation)
	// contains filtered or unexported fields
}

ActiveRequest represents an ongoing chat request

func (*ActiveRequest) Context

func (r *ActiveRequest) Context() context.Context

Context returns the request's context

func (*ActiveRequest) IsCancelled

func (r *ActiveRequest) IsCancelled() bool

IsCancelled returns true if the request has been cancelled

type AgentClient

type AgentClient struct {
	// contains filtered or unexported fields
}

func NewAgentClient

func NewAgentClient(cfg *config.AgentsConfig, eventBus *events.EventBus, db *sql.DB, selfAgentID, selfAgentName string) *AgentClient

func (*AgentClient) CallAgent

func (c *AgentClient) CallAgent(agentID, message, contextType, threadID string) (*CallResult, error)

func (*AgentClient) CallAgentStreaming

func (c *AgentClient) CallAgentStreaming(agentID, message, contextType, threadID string, callback tools.StreamCallback) (*CallResult, error)

CallAgentStreaming calls another agent with streaming support The callback is called for each chunk of the response

func (*AgentClient) CallAgentStreamingWithContext

func (c *AgentClient) CallAgentStreamingWithContext(ctx context.Context, agentID, message, contextType, threadID string, callback tools.StreamCallback) (*CallResult, error)

CallAgentStreamingWithContext calls another agent with streaming and context support

func (*AgentClient) CallAgentWithContext

func (c *AgentClient) CallAgentWithContext(ctx context.Context, agentID, message, contextType, threadID string) (*CallResult, error)

CallAgentWithContext calls another agent with context support for cancellation

func (*AgentClient) CheckDelegatedTask

func (c *AgentClient) CheckDelegatedTask(agentID, taskID string) (map[string]interface{}, error)

CheckDelegatedTask checks the status of a task on a remote agent

func (*AgentClient) ClearCallContext

func (c *AgentClient) ClearCallContext()

ClearCallContext clears the current call context

func (*AgentClient) ClearRequestContext

func (c *AgentClient) ClearRequestContext()

ClearRequestContext clears the current request context

func (*AgentClient) DelegateTask

func (c *AgentClient) DelegateTask(agentID, title, description string, priority int, executeAt string, testMode bool) (map[string]interface{}, error)

DelegateTask creates a task on a remote agent for async execution

func (*AgentClient) GetAgentActivity added in v1.38.0

func (c *AgentClient) GetAgentActivity(agentID, since string, limit int) (map[string]interface{}, error)

GetAgentActivity fetches the activity summary from a remote agent's /activity endpoint

func (*AgentClient) GetAgentName added in v1.37.0

func (c *AgentClient) GetAgentName(agentID string) string

GetAgentName returns the human-readable name for an agent ID, or the ID itself if not found

func (*AgentClient) GetAvailableAgents

func (c *AgentClient) GetAvailableAgents(filterTags, filterCapabilities []string) []map[string]interface{}

func (*AgentClient) GetGuardRails

func (c *AgentClient) GetGuardRails() *GuardRails

GetGuardRails returns the guard rails instance for external use

func (*AgentClient) SetCallContext

func (c *AgentClient) SetCallContext(ctx *CallContext)

SetCallContext sets the call context for the current request This should be called when handling an incoming agent-to-agent call

func (*AgentClient) SetDiscoveryService

func (c *AgentClient) SetDiscoveryService(ds DiscoveryService)

SetDiscoveryService injects the discovery service

func (*AgentClient) SetRequestContext

func (c *AgentClient) SetRequestContext(requestID, selfURL string)

SetRequestContext sets the current request ID and self URL for cancellation propagation

func (*AgentClient) SetSelfAgentURL

func (c *AgentClient) SetSelfAgentURL(url string)

SetSelfAgentURL sets this agent's URL for cancel callbacks

type CallAgentTool

type CallAgentTool struct {
	Client *AgentClient
}

CallAgentTool implements the call_agent tool with streaming support

func (*CallAgentTool) Description

func (t *CallAgentTool) Description() string

func (*CallAgentTool) DisplayName

func (t *CallAgentTool) DisplayName() string

func (*CallAgentTool) DynamicDisplayName added in v1.37.0

func (t *CallAgentTool) DynamicDisplayName(params map[string]interface{}) string

DynamicDisplayName returns "Calling {AgentName}" based on the input params

func (*CallAgentTool) Execute

func (t *CallAgentTool) Execute(params map[string]interface{}) (interface{}, error)

Execute falls back to non-streaming for backwards compatibility

func (*CallAgentTool) ExecuteStreaming

func (t *CallAgentTool) ExecuteStreaming(params map[string]interface{}, callback tools.StreamCallback) (interface{}, error)

ExecuteStreaming calls another agent and streams the response

func (*CallAgentTool) InputSchema

func (t *CallAgentTool) InputSchema() map[string]interface{}

func (*CallAgentTool) Name

func (t *CallAgentTool) Name() string

func (*CallAgentTool) SupportsStreaming

func (t *CallAgentTool) SupportsStreaming() bool

SupportsStreaming returns true - this tool streams agent responses

type CallContext

type CallContext struct {
	CallChain   []string  `json:"call_chain"`   // List of agent IDs in the call path
	CallDepth   int       `json:"call_depth"`   // Current depth in the call tree
	OriginAgent string    `json:"origin_agent"` // The agent that started the chain
	CallID      string    `json:"call_id"`      // Unique identifier for this call tree
	StartTime   time.Time `json:"start_time"`   // When the call chain started
}

CallContext holds the context of a call chain

func ExtractCallContext

func ExtractCallContext(r *http.Request) *CallContext

ExtractCallContext extracts call context from incoming request headers

type CallResult

type CallResult struct {
	Success    bool
	AgentID    string
	AgentName  string
	Response   string
	ThreadID   string
	DurationMS int64
	TokensUsed int
	Error      string
}

type CancelRequest

type CancelRequest struct {
	Reason string `json:"reason,omitempty"`
}

CancelRequest is the request body for cancel endpoint

type CancelResponse

type CancelResponse struct {
	Success   bool   `json:"success"`
	RequestID string `json:"request_id"`
	Cancelled bool   `json:"cancelled"`
	Error     string `json:"error,omitempty"`
}

CancelResponse is the response from cancel endpoint

type CancellationConfig

type CancellationConfig struct {
	Enabled           bool   `json:"enabled"`
	MaxActiveRequests int    `json:"max_active_requests"`
	CleanupInterval   string `json:"cleanup_interval"`
	PropagateToAgents bool   `json:"propagate_to_agents"`
}

CancellationConfig holds configuration for request cancellation

type ChatRequest

type ChatRequest struct {
	Message  string `json:"message"`
	Stream   bool   `json:"stream"` // Always false for agent-to-agent
	ThreadID string `json:"thread_id,omitempty"`
}

ChatRequest matches the /chat endpoint structure

type ChatResponse

type ChatResponse struct {
	Success  bool   `json:"success"`
	Response string `json:"response"`
	ThreadID string `json:"thread_id"`
	Model    string `json:"model,omitempty"`
	Error    string `json:"error,omitempty"`
}

ChatResponse matches the /chat endpoint non-streaming response

type CheckDelegatedTaskTool

type CheckDelegatedTaskTool struct {
	Client *AgentClient
}

CheckDelegatedTaskTool implements the check_delegated_task tool

func (*CheckDelegatedTaskTool) Description

func (t *CheckDelegatedTaskTool) Description() string

func (*CheckDelegatedTaskTool) DisplayName

func (t *CheckDelegatedTaskTool) DisplayName() string

func (*CheckDelegatedTaskTool) DynamicDisplayName added in v1.37.0

func (t *CheckDelegatedTaskTool) DynamicDisplayName(params map[string]interface{}) string

DynamicDisplayName returns "Checking task on {AgentName}" based on the input params

func (*CheckDelegatedTaskTool) Execute

func (t *CheckDelegatedTaskTool) Execute(params map[string]interface{}) (interface{}, error)

func (*CheckDelegatedTaskTool) InputSchema

func (t *CheckDelegatedTaskTool) InputSchema() map[string]interface{}

func (*CheckDelegatedTaskTool) Name

func (t *CheckDelegatedTaskTool) Name() string

type CircuitBreaker

type CircuitBreaker struct {
	// contains filtered or unexported fields
}

CircuitBreaker prevents repeated calls to failing agents

func NewCircuitBreaker

func NewCircuitBreaker() *CircuitBreaker

NewCircuitBreaker creates a new circuit breaker with default settings

func (*CircuitBreaker) CanCall

func (cb *CircuitBreaker) CanCall(agentID string) bool

CanCall checks if a call to the agent is allowed

func (*CircuitBreaker) Configure

func (cb *CircuitBreaker) Configure(threshold int, resetAfter, openFor time.Duration)

Configure updates circuit breaker settings

func (*CircuitBreaker) GetState

func (cb *CircuitBreaker) GetState(agentID string) string

GetState returns the current state for an agent

func (*CircuitBreaker) RecordFailure

func (cb *CircuitBreaker) RecordFailure(agentID string)

RecordFailure records a failed call and potentially opens the circuit

func (*CircuitBreaker) RecordSuccess

func (cb *CircuitBreaker) RecordSuccess(agentID string)

RecordSuccess records a successful call and resets the failure count

func (*CircuitBreaker) TimeUntilReset

func (cb *CircuitBreaker) TimeUntilReset(agentID string) time.Duration

TimeUntilReset returns how long until the circuit might close for an agent

type DelegateTaskTool

type DelegateTaskTool struct {
	Client *AgentClient
}

DelegateTaskTool implements the delegate_task tool for async task delegation

func (*DelegateTaskTool) Description

func (t *DelegateTaskTool) Description() string

func (*DelegateTaskTool) DisplayName

func (t *DelegateTaskTool) DisplayName() string

func (*DelegateTaskTool) DynamicDisplayName added in v1.37.0

func (t *DelegateTaskTool) DynamicDisplayName(params map[string]interface{}) string

DynamicDisplayName returns "Delegating to {AgentName}" based on the input params

func (*DelegateTaskTool) Execute

func (t *DelegateTaskTool) Execute(params map[string]interface{}) (interface{}, error)

func (*DelegateTaskTool) InputSchema

func (t *DelegateTaskTool) InputSchema() map[string]interface{}

func (*DelegateTaskTool) Name

func (t *DelegateTaskTool) Name() string

type DiscoveryService

type DiscoveryService interface {
	GetAgents() []config.AgentInfo
	IsRunning() bool
}

DiscoveryService interface (avoid circular import)

type GetAgentActivityTool added in v1.38.0

type GetAgentActivityTool struct {
	Client *AgentClient
}

GetAgentActivityTool implements the get_agent_activity tool for coordinators

func (*GetAgentActivityTool) Description added in v1.38.0

func (t *GetAgentActivityTool) Description() string

func (*GetAgentActivityTool) DisplayName added in v1.38.0

func (t *GetAgentActivityTool) DisplayName() string

func (*GetAgentActivityTool) DynamicDisplayName added in v1.38.0

func (t *GetAgentActivityTool) DynamicDisplayName(params map[string]interface{}) string

DynamicDisplayName returns "Checking activity of {AgentName}" based on the input params

func (*GetAgentActivityTool) Execute added in v1.38.0

func (t *GetAgentActivityTool) Execute(params map[string]interface{}) (interface{}, error)

func (*GetAgentActivityTool) InputSchema added in v1.38.0

func (t *GetAgentActivityTool) InputSchema() map[string]interface{}

func (*GetAgentActivityTool) Name added in v1.38.0

func (t *GetAgentActivityTool) Name() string

type GuardRailError

type GuardRailError struct {
	Type    string      `json:"type"`    // "circular_call", "max_depth", "circuit_open", "ttl_expired"
	Message string      `json:"message"` // Human-readable message
	Details interface{} `json:"details"` // Additional context
}

GuardRailError represents a blocked call due to guard rails

func (*GuardRailError) Error

func (e *GuardRailError) Error() string

type GuardRails

type GuardRails struct {
	// contains filtered or unexported fields
}

GuardRails provides protection against loops and cascading failures

func NewGuardRails

func NewGuardRails(selfAgentID string) *GuardRails

NewGuardRails creates a new guard rails instance

func (*GuardRails) Configure

func (g *GuardRails) Configure(maxDepth int, allowRecursive, trackChain bool, ttl time.Duration)

Configure updates guard rails settings

func (*GuardRails) ConfigureCircuitBreaker

func (g *GuardRails) ConfigureCircuitBreaker(threshold int, resetAfter, openFor time.Duration)

ConfigureCircuitBreaker updates circuit breaker settings

func (*GuardRails) RecordFailure

func (g *GuardRails) RecordFailure(agentID string)

RecordFailure records a failed call

func (*GuardRails) RecordSuccess

func (g *GuardRails) RecordSuccess(agentID string)

RecordSuccess records a successful call

func (*GuardRails) ValidateCall

func (g *GuardRails) ValidateCall(ctx *CallContext, targetAgentID string) *GuardRailError

ValidateCall checks if a call to targetAgentID is allowed

type ListAvailableAgentsTool

type ListAvailableAgentsTool struct {
	Client *AgentClient
}

ListAvailableAgentsTool implements list_available_agents

func (*ListAvailableAgentsTool) Description

func (t *ListAvailableAgentsTool) Description() string

func (*ListAvailableAgentsTool) DisplayName

func (t *ListAvailableAgentsTool) DisplayName() string

func (*ListAvailableAgentsTool) Execute

func (t *ListAvailableAgentsTool) Execute(params map[string]interface{}) (interface{}, error)

func (*ListAvailableAgentsTool) InputSchema

func (t *ListAvailableAgentsTool) InputSchema() map[string]interface{}

func (*ListAvailableAgentsTool) Name

func (t *ListAvailableAgentsTool) Name() string

type RequestStatus

type RequestStatus string

RequestStatus represents the current state of a request

const (
	RequestStatusProcessing RequestStatus = "processing"
	RequestStatusCancelled  RequestStatus = "cancelled"
	RequestStatusCompleted  RequestStatus = "completed"
)

type RequestTracker

type RequestTracker struct {
	// contains filtered or unexported fields
}

RequestTracker manages active requests and their cancellation

func NewRequestTracker

func NewRequestTracker() *RequestTracker

NewRequestTracker creates a new request tracker

func (*RequestTracker) AddCancelCallback

func (rt *RequestTracker) AddCancelCallback(requestID, callbackURL string)

AddCancelCallback adds a callback URL to be called when request is cancelled

func (*RequestTracker) Cancel

func (rt *RequestTracker) Cancel(requestID string) error

Cancel cancels an active request

func (*RequestTracker) Complete

func (rt *RequestTracker) Complete(requestID string)

Complete marks a request as completed

func (*RequestTracker) Configure

func (rt *RequestTracker) Configure(maxActive int, cleanupAfter time.Duration)

Configure updates tracker settings

func (*RequestTracker) Get

func (rt *RequestTracker) Get(requestID string) *ActiveRequest

Get returns an active request by ID

func (*RequestTracker) GetActiveCount

func (rt *RequestTracker) GetActiveCount() int

GetActiveCount returns the number of active (processing) requests

func (*RequestTracker) GetAll

func (rt *RequestTracker) GetAll() []*ActiveRequest

GetAll returns all requests (for debugging)

func (*RequestTracker) Register

func (rt *RequestTracker) Register(requestID, threadID string) (context.Context, *ActiveRequest, error)

Register creates a new active request and returns its context The returned context will be cancelled when Cancel is called

func (*RequestTracker) RegisterWithParentContext

func (rt *RequestTracker) RegisterWithParentContext(parentCtx context.Context, requestID, threadID string) (context.Context, *ActiveRequest, error)

RegisterWithParentContext creates a request with a parent context Cancellation of either parent or explicit cancel will cancel the request

type SSECancelledEvent

type SSECancelledEvent struct {
	RequestID string `json:"request_id"`
	Reason    string `json:"reason"`
}

SSECancelledEvent returns the SSE event data for a cancelled request

type SSEStartEvent

type SSEStartEvent struct {
	RequestID string `json:"request_id"`
	ThreadID  string `json:"thread_id,omitempty"`
}

SSEStartEvent returns the SSE event data for request start

Jump to

Keyboard shortcuts

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