Documentation
¶
Overview ¶
Package core provides a conversation engine for LLM-powered agents.
Core concepts:
- Agent: the main entry point. Create with NewAgent(), run queries with Run().
- Provider: interface for LLM backends. Implement Chat() and ChatStream().
- ToolExecutor: interface for tool execution. Implement GetTools() and Execute().
- UI: optional interface for interactive output. Use NoopUI for headless mode.
- EventPublisher: optional interface for event-driven output.
Quick start:
agent, err := core.NewAgent(core.Options{
Provider: &myProvider{},
Executor: core.NoopExecutor,
})
result, err := agent.Run(ctx, "Hello")
The core package has zero external dependencies. The events and internal/test packages are optional.
Index ¶
- Constants
- Variables
- func BuildCheckpointCompactedMessages(messages []Message, checkpoints []TurnCheckpoint) ([]Message, []TurnCheckpoint)
- func ClassifyError(err error, provider string) error
- func IsAuthError(err error) bool
- func IsBlankResponse(err error) bool
- func IsContentFiltered(err error) bool
- func IsContextOverflow(err error) bool
- func IsRateLimit(err error) bool
- func IsTransient(err error) bool
- func RecordTurnCheckpointAsync(state *State, messages []Message, startIndex, endIndex int, ...)
- type Agent
- func (a *Agent) ExportState() ([]byte, error)
- func (a *Agent) ImportState(data []byte) error
- func (a *Agent) InjectInput(input string) bool
- func (a *Agent) IsPaused() bool
- func (a *Agent) Pause()
- func (a *Agent) Provider() Provider
- func (a *Agent) ReasoningBuffer() *StreamingBuffer
- func (a *Agent) Resume()
- func (a *Agent) Run(ctx context.Context, query string) (string, error)
- func (a *Agent) RunStream(ctx context.Context, query string) (string, error)
- func (a *Agent) SetFlushCallback(fn func())
- func (a *Agent) SetProvider(p Provider)
- func (a *Agent) SetSystemPrompt(prompt string)
- func (a *Agent) State() *State
- func (a *Agent) Steer(msg Message)
- func (a *Agent) SteerSystem(content string)
- func (a *Agent) StreamingBuffer() *StreamingBuffer
- type AgentState
- type AgentStreamHandler
- type AuthError
- type BlankResponseError
- type ChatChoice
- type ChatRequest
- type ChatResponse
- type ChatUsage
- type CompactionResult
- type Compactor
- type ContentFilteredError
- type ContextOverflowError
- type ConversationHandler
- type ConversationOptimizer
- type ConversationOptimizerOptions
- type EventPublisher
- type ExponentialBackoff
- type FallbackParseResult
- type FallbackParser
- type FallbackParserOptions
- type ImageData
- type Message
- type NormalizedToolCalls
- type Options
- type OutputEvent
- type OutputManager
- type Provider
- type ProviderInfo
- type RateLimitError
- type ResponseValidator
- type ResponseValidatorOptions
- type RetryConfig
- type State
- func (s *State) AddCheckpoint(cp TurnCheckpoint)
- func (s *State) AddCost(cost float64)
- func (s *State) AddMessage(msg Message)
- func (s *State) AddTokens(prompt, completion, total int)
- func (s *State) ClearCheckpoints()
- func (s *State) EnsureSessionID()
- func (s *State) ExportState() ([]byte, error)
- func (s *State) GetCheckpoints() []TurnCheckpoint
- func (s *State) ImportState(data []byte) error
- func (s *State) LastAssistantMessage() *Message
- func (s *State) Len() int
- func (s *State) Messages() []Message
- func (s *State) SessionID() string
- func (s *State) SetCheckpoints(cps []TurnCheckpoint)
- func (s *State) SetMessages(msgs []Message)
- func (s *State) SetSessionID(id string)
- func (s *State) TotalCost() float64
- func (s *State) TotalTokens() int
- type StreamHandler
- type StreamingBuffer
- type Tool
- type ToolCall
- type ToolCallFunction
- type ToolCallNormalizer
- type ToolCategory
- type ToolExecutor
- type ToolFunction
- type ToolParameter
- type ToolParameters
- type TransientError
- type TurnCheckpoint
- type TurnSummaryBuilder
- type UI
Constants ¶
const ( // EventTypeQueryStarted is published when a new query begins. // Data map keys: "query" (string), "model" (string) EventTypeQueryStarted = "query_started" // EventTypeQueryCompleted is published when a query finishes. // Data map keys: "query", "response", "tokens", "cost", "duration_ms" EventTypeQueryCompleted = "query_completed" // EventTypeError is published on errors. // Data map keys: "message" (string), "error" (string) EventTypeError = "error" // EventTypeToolStart is published when a tool call begins. // Data map keys: "tool_name", "tool_call_id", "arguments", "tool_index" EventTypeToolStart = "tool_start" // EventTypeToolEnd is published when a tool call ends. // Data map keys: "tool_call_id", "tool_name", "status", "result", "duration_ms" EventTypeToolEnd = "tool_end" // EventTypeStreamChunk is published for each streaming content chunk. // Data map keys: "chunk" (string), "content_type" (string: "text"|"reasoning") EventTypeStreamChunk = "stream_chunk" // EventTypeMetricsUpdate is published with token usage updates. // Data map keys: "total_tokens", "context_tokens", "max_context_tokens", "iteration", "total_cost" EventTypeMetricsUpdate = "metrics_update" // EventTypeCompaction is published when context compaction occurs. // Data map keys: "strategy", "messages_before", "messages_after", "message_count_delta", "tokens_saved" EventTypeCompaction = "compaction" )
Generic event type constants used by the core package. These are the only event types that any consumer of core needs to handle.
const DefaultSystemPrompt = "You are a helpful assistant that can execute tools to complete tasks."
DefaultSystemPrompt is the minimal system prompt used when none is provided.
Variables ¶
var ( // ErrNoProvider is returned when NewAgent is called without a Provider. ErrNoProvider = errors.New("no provider configured") // ErrNoExecutor is returned when NewAgent is called without a ToolExecutor. ErrNoExecutor = errors.New("no tool executor configured") // ErrInterrupted is returned when the conversation is interrupted by the user. ErrInterrupted = errors.New("conversation interrupted") // ErrMaxIterations is returned when the maximum iteration count is exceeded. ErrMaxIterations = errors.New("maximum iterations exceeded") // ErrPaused is returned when Run is called while the agent is paused. ErrPaused = errors.New("agent is paused") )
Sentinel errors returned by the agent lifecycle.
Functions ¶
func BuildCheckpointCompactedMessages ¶
func BuildCheckpointCompactedMessages(messages []Message, checkpoints []TurnCheckpoint) ([]Message, []TurnCheckpoint)
BuildCheckpointCompactedMessages replaces consumed checkpoint ranges with summary messages and returns the compacted message list and updated checkpoints.
A checkpoint is "consumable" when: - Its StartIndex and EndIndex are valid (within the messages slice) - The messages in the range [StartIndex, EndIndex] exist - The checkpoint has a non-empty Summary
All checkpoint indices reference the raw state.Messages() slice. Because state.Messages() only appends (never inserts or deletes in the middle), checkpoint indices remain stable across calls. Consumed checkpoints are kept in the returned list so that every call to prepareMessages() can re-apply their summaries against the growing message history.
The function works from oldest to newest checkpoint (by StartIndex):
- For each consumable checkpoint, replace messages[StartIndex:EndIndex+1] with a single summary message (role "user", content from checkpoint.Summary)
- Handle consecutive-assistant boundaries: if the inserted summary message would create two consecutive assistant messages with no tool calls between them, merge or adjust.
The summary message role is "user" to maintain proper conversation flow (system -> user -> assistant pattern).
Return:
- compactedMessages: the new message slice with checkpoints applied
- updatedCheckpoints: all checkpoints preserved with their original indices (consumed checkpoints are kept so future calls re-apply their summaries)
func ClassifyError ¶
ClassifyError wraps a raw provider error in a typed error based on error message patterns. If err is nil or already a typed error, it is returned unchanged.
func IsAuthError ¶
IsAuthError reports whether err is an AuthError.
func IsBlankResponse ¶ added in v0.1.1
IsBlankResponse reports whether err is a BlankResponseError.
func IsContentFiltered ¶ added in v0.1.1
IsContentFiltered reports whether err is a ContentFilteredError.
func IsContextOverflow ¶
IsContextOverflow reports whether err is a ContextOverflowError.
func IsRateLimit ¶
IsRateLimit reports whether err is a RateLimitError.
func IsTransient ¶
IsTransient reports whether err is a TransientError.
func RecordTurnCheckpointAsync ¶
func RecordTurnCheckpointAsync(state *State, messages []Message, startIndex, endIndex int, timeout time.Duration)
RecordTurnCheckpointAsync asynchronously builds a checkpoint from the given messages and stores it in state. It spawns a goroutine to compute the summary so it doesn't block the conversation loop.
The message slice is snapshotted immediately (before the goroutine starts) so the background computation sees a consistent view even if the caller mutates the original slice. If the summary computation takes longer than timeout, a minimal checkpoint is stored instead.
Types ¶
type Agent ¶
type Agent struct {
// contains filtered or unexported fields
}
Agent is the main entry point for the conversation engine.
func NewAgent ¶
NewAgent creates a new Agent from the given options. Returns an error if required options (Provider or ToolExecutor) are not provided.
func (*Agent) ExportState ¶
ExportState serializes the current state to JSON.
func (*Agent) ImportState ¶
ImportState deserializes state from JSON.
func (*Agent) InjectInput ¶
InjectInput injects a user message into the conversation via a buffered channel. Returns true if the input was accepted (queued for the next loop iteration), false if a prior injection is still pending. The injection is fire-and-forget with backpressure: a true return means the input is queued, but the caller has no guarantee it was consumed yet.
func (*Agent) ReasoningBuffer ¶
func (a *Agent) ReasoningBuffer() *StreamingBuffer
ReasoningBuffer returns the reasoning streaming buffer.
func (*Agent) RunStream ¶
RunStream executes a single query through the streaming conversation loop. It uses provider.ChatStream() instead of provider.Chat(), so content is delivered incrementally via the StreamHandler callbacks. The streaming buffer is populated with content as it arrives, and the return value is empty when the buffer contains streamed content (the caller should read from StreamingBuffer() instead).
func (*Agent) SetFlushCallback ¶
func (a *Agent) SetFlushCallback(fn func())
SetFlushCallback sets a callback to flush streaming output.
func (*Agent) SetProvider ¶
SetProvider swaps the provider at runtime. The new provider will be used for all subsequent calls. It is safe to call between queries; calling during an active query may cause undefined behavior. Panics if p is nil — use NewAgent to create an agent without a provider and then call SetProvider with a non-nil provider.
func (*Agent) SetSystemPrompt ¶
SetSystemPrompt updates the system prompt for future queries.
func (*Agent) Steer ¶
Steer queues a transient message that will be appended to the next API call made by this agent. The message is consumed once and is not persisted in the conversation state. Use Steer to inject temporary guidance (e.g., "Focus on security concerns" or "Respond in JSON format") that should influence only the next model response.
Steer must be called between Run() calls. If called during an active Run(), the message is queued and will not be consumed until the next Run().
Example:
agent.Steer(core.Message{Role: "user", Content: "Focus on performance."})
func (*Agent) SteerSystem ¶
SteerSystem is a convenience for injecting a system-level steering message. It creates a Message with Role "system" and queues it via Steer, so the guidance takes effect on the next API call and is consumed once (like all steer messages). The message is not persisted in the conversation state.
Example:
agent.SteerSystem("Focus on performance, not correctness.")
func (*Agent) StreamingBuffer ¶
func (a *Agent) StreamingBuffer() *StreamingBuffer
StreamingBuffer returns the content streaming buffer.
type AgentState ¶
type AgentState struct {
Messages []Message `json:"messages"`
SessionID string `json:"session_id"`
TotalTokens int `json:"total_tokens"`
TotalCost float64 `json:"total_cost"`
PromptTokens int `json:"prompt_tokens"`
CompletionTokens int `json:"completion_tokens"`
Checkpoints []TurnCheckpoint `json:"checkpoints,omitempty"`
}
AgentState tracks the state of an agent's conversation.
type AgentStreamHandler ¶
type AgentStreamHandler struct {
// contains filtered or unexported fields
}
AgentStreamHandler is a concrete StreamHandler implementation that writes streamed content to the agent's buffers and publishes events via the agent's eventPublisher.
func NewAgentStreamHandler ¶
func NewAgentStreamHandler(agent *Agent, state *State) *AgentStreamHandler
NewAgentStreamHandler creates a new StreamHandler for the given agent.
func (*AgentStreamHandler) OnContent ¶
func (h *AgentStreamHandler) OnContent(content string)
OnContent handles a streaming content chunk: writes to the stream buffer, publishes stream_chunk and agent_message events, and invokes the flush callback. Empty content chunks are silently ignored.
func (*AgentStreamHandler) OnDone ¶
func (h *AgentStreamHandler) OnDone(resp *ChatResponse)
OnDone handles the end of streaming: records token usage, publishes metrics, and appends the final assistant message to the agent's state.
func (*AgentStreamHandler) OnError ¶
func (h *AgentStreamHandler) OnError(err error)
OnError handles streaming errors: publishes an error event via the bus.
func (*AgentStreamHandler) OnReasoning ¶
func (h *AgentStreamHandler) OnReasoning(reasoning string)
OnReasoning handles a streaming reasoning chunk: writes to the reasoning buffer, publishes a stream_chunk event, and invokes the flush callback. Empty reasoning chunks are silently ignored.
func (*AgentStreamHandler) Response ¶
func (h *AgentStreamHandler) Response() *ChatResponse
Response returns the final ChatResponse from OnDone, or nil if not yet set.
type AuthError ¶
type AuthError struct {
// Provider is the provider that returned the auth error.
Provider string
// Wrapped is the underlying error.
Wrapped error
}
AuthError indicates authentication failure with a provider.
type BlankResponseError ¶ added in v0.1.1
type BlankResponseError struct {
// Provider is the provider that produced the blank responses.
Provider string
// Count is the number of consecutive blank/repetitive responses.
Count int
}
BlankResponseError indicates the model produced consecutive blank or repetitive responses and the conversation was force-finalized.
func (*BlankResponseError) Error ¶ added in v0.1.1
func (e *BlankResponseError) Error() string
type ChatChoice ¶
type ChatChoice struct {
Index int `json:"index"`
Message Message `json:"message"`
FinishReason string `json:"finish_reason"`
}
ChatChoice represents one possible completion from the model.
type ChatRequest ¶
type ChatRequest struct {
Model string `json:"model,omitempty"`
Messages []Message `json:"messages"`
Tools []Tool `json:"tools,omitempty"`
ToolChoice string `json:"tool_choice,omitempty"`
MaxTokens int `json:"max_tokens,omitempty"`
Reasoning string `json:"reasoning,omitempty"`
Stream bool `json:"stream,omitempty"`
}
ChatRequest is a request to the chat completion endpoint.
type ChatResponse ¶
type ChatResponse struct {
ID string `json:"id"`
Model string `json:"model"`
Choices []ChatChoice `json:"choices"`
Usage ChatUsage `json:"usage"`
}
ChatResponse is a response from the chat completion endpoint.
func (ChatResponse) ToMessage ¶
func (r ChatResponse) ToMessage() Message
ToMessage returns the first choice's message from the response, or an empty message.
type ChatUsage ¶
type ChatUsage struct {
PromptTokens int `json:"prompt_tokens"`
CompletionTokens int `json:"completion_tokens"`
TotalTokens int `json:"total_tokens"`
EstimatedCost float64 `json:"estimated_cost,omitempty"`
Cost float64 `json:"cost,omitempty"`
}
ChatUsage tracks token usage and cost for a request/response.
type CompactionResult ¶
type CompactionResult struct {
Messages []Message
Strategy string // "none", "checkpoint", "structural", "emergency"
TokensBefore int
TokensAfter int
}
CompactionResult holds the output of a compaction operation along with metadata about what strategy was used and how much was saved.
func (CompactionResult) MessageCountDelta ¶
func (r CompactionResult) MessageCountDelta(before int) int
MessageCountDelta returns how many messages were removed.
func (CompactionResult) TokensSaved ¶
func (r CompactionResult) TokensSaved() int
TokensSaved returns the estimated tokens saved by compaction.
type Compactor ¶
type Compactor struct {
// contains filtered or unexported fields
}
Compactor reduces a message list to fit within a context window. It implements the strategy from sprout: checkpoint compaction, structural compaction, and emergency truncation.
This is a concrete type now. When different strategies are needed, extract an interface — the loop code doesn't change.
func NewCompactor ¶
func NewCompactor() *Compactor
NewCompactor creates a compactor with default settings.
type ContentFilteredError ¶ added in v0.1.1
type ContentFilteredError struct {
// Provider is the provider that returned the content filter.
Provider string
}
ContentFilteredError indicates the provider's content filter blocked the response after a retry attempt was already made.
func (*ContentFilteredError) Error ¶ added in v0.1.1
func (e *ContentFilteredError) Error() string
type ContextOverflowError ¶
type ContextOverflowError struct {
// TokensUsed is the number of tokens that were estimated or used.
TokensUsed int
// TokensLimit is the provider's context window limit.
TokensLimit int
// Wrapped is the underlying error.
Wrapped error
}
ContextOverflowError indicates the context window is exceeded.
func (*ContextOverflowError) Error ¶
func (e *ContextOverflowError) Error() string
func (*ContextOverflowError) Unwrap ¶
func (e *ContextOverflowError) Unwrap() error
type ConversationHandler ¶
type ConversationHandler struct {
// contains filtered or unexported fields
}
ConversationHandler manages the high-level conversation flow.
func (*ConversationHandler) ProcessQuery ¶
ProcessQuery handles a user query through the complete conversation flow.
func (*ConversationHandler) ProcessQueryStream ¶
func (ch *ConversationHandler) ProcessQueryStream(ctx context.Context, query string) (string, error)
ProcessQueryStream handles a user query through the streaming conversation flow. It uses provider.ChatStream() instead of provider.Chat(), so content is delivered incrementally via StreamHandler callbacks. The streaming buffer is populated as content arrives, and the return value is empty when the buffer contains streamed content (the caller should read from agent.StreamingBuffer() instead).
type ConversationOptimizer ¶
type ConversationOptimizer struct {
// contains filtered or unexported fields
}
ConversationOptimizer reduces redundant conversation history by deduplicating repeated file reads and transient shell command outputs. It mutates the provided message slice in place (replacing tool-result Content with placeholders) but is safe to use because prepareMessages() only ever passes it ephemeral copies of state — the stored conversation is never modified.
func NewConversationOptimizer ¶
func NewConversationOptimizer(opts ConversationOptimizerOptions) *ConversationOptimizer
NewConversationOptimizer creates a new optimizer from the given options.
func (*ConversationOptimizer) OptimizeConversation ¶
func (opt *ConversationOptimizer) OptimizeConversation(messages []Message) []Message
OptimizeConversation processes the message list, replacing redundant file reads and shell command outputs with compact placeholders. Returns the (possibly modified) message slice.
type ConversationOptimizerOptions ¶
type ConversationOptimizerOptions struct {
// Enabled enables optimization. When false, OptimizeConversation is a no-op.
Enabled bool
// KnownToolFn classifies tool names. Return ToolCategoryUnknown to skip a tool.
// If nil, the optimizer treats all tools as unknown (no optimization).
// Must be deterministic: the same tool name should always return the same category.
KnownToolFn func(name string) ToolCategory
}
ConversationOptimizerOptions configures the optimizer.
type EventPublisher ¶
EventPublisher is the interface for publishing events. It is implemented by events.EventBus and can be satisfied by any custom event system, allowing seed to be used without the events package.
type ExponentialBackoff ¶
type ExponentialBackoff struct {
InitialDelay time.Duration // delay for the first retry
MaxDelay time.Duration // cap on delay growth (0 = no cap)
Multiplier float64 // exponential growth factor (typically > 1)
MaxAttempts int // total number of attempts (0 = unlimited)
Jitter float64 // jitter: 0 = none, (0,1) = partial [base, base*(1+jitter)], >= 1 = full jitter [0, base]
// contains filtered or unexported fields
}
ExponentialBackoff implements exponential backoff with jitter for retry logic. It is NOT safe for concurrent use by multiple goroutines.
Usage:
backoff := &ExponentialBackoff{InitialDelay: 100 * time.Millisecond, MaxDelay: 10 * time.Second, MaxAttempts: 5}
for backoff.NextAttempt() {
time.Sleep(backoff.Delay())
if ok := doWork(); ok {
break
}
}
func NewExponentialBackoff ¶
func NewExponentialBackoff(initialDelay, maxDelay time.Duration, multiplier float64, maxAttempts int, jitter float64) *ExponentialBackoff
NewExponentialBackoff creates a configured ExponentialBackoff with a global random source. Pass a custom randSource to make jitter deterministic in tests.
func (*ExponentialBackoff) Delay ¶
func (b *ExponentialBackoff) Delay() time.Duration
Delay returns the current delay duration with jitter applied. Jitter adds a random percentage of the base delay to prevent thundering herd. If Jitter is in (0, 1), the delay is [base, base * (1 + Jitter)]. If Jitter is >= 1, full jitter mode replaces delay with [0, base]. Panics if called before NextAttempt() or if randSrc is nil.
func (*ExponentialBackoff) NextAttempt ¶
func (b *ExponentialBackoff) NextAttempt() bool
NextAttempt advances to the next retry and returns true if the attempt count has not exceeded MaxAttempts. Returns false immediately if MaxAttempts is already reached or exceeded.
func (*ExponentialBackoff) Reset ¶
func (b *ExponentialBackoff) Reset()
Reset resets the backoff state so it can be reused for a new retry sequence.
func (*ExponentialBackoff) WithRand ¶
func (b *ExponentialBackoff) WithRand(rs randSource) *ExponentialBackoff
WithRand sets a custom random source (useful for deterministic tests).
type FallbackParseResult ¶
FallbackParseResult contains extracted tool calls and cleaned content.
type FallbackParser ¶
type FallbackParser struct {
// contains filtered or unexported fields
}
FallbackParser extracts tool calls from malformed LLM response content.
func NewFallbackParser ¶
func NewFallbackParser(opts FallbackParserOptions) *FallbackParser
NewFallbackParser creates a new FallbackParser with the given options.
func (*FallbackParser) Parse ¶
func (fp *FallbackParser) Parse(content string) *FallbackParseResult
Parse extracts tool calls from malformed LLM response content.
func (*FallbackParser) ShouldUseFallback ¶
func (fp *FallbackParser) ShouldUseFallback(content string, hasStructuredToolCalls bool) bool
ShouldUseFallback returns true when structured tool_calls are missing and the content contains patterns suggestive of tool calls.
It uses a three-tier confidence model:
- Tier 1: Strong patterns (code fences, XML tags, quoted JSON keys) trigger immediately with a single match.
- Tier 2a: Weak patterns (bare keywords like "name:" or "arguments:") require at least two independent matches to avoid false positives on normal conversational text.
- Tier 2b: One weak pattern plus a JSON structure marker ({" or [") triggers a single weak match, catching patterns like: "name: search\n{\"query\": \"hello\"}"
type FallbackParserOptions ¶
type FallbackParserOptions struct {
// KnownToolNames returns true if the given name is a registered tool.
// When nil, all extracted tool names are accepted.
KnownToolNames func(string) bool
// Debug enables verbose logging (printed to stderr).
Debug bool
}
FallbackParserOptions configures the parser.
type ImageData ¶
type ImageData struct {
URL string `json:"url,omitempty"`
Base64 string `json:"base64,omitempty"`
Type string `json:"type,omitempty"` // MIME type (image/jpeg, image/png, etc.)
}
ImageData represents an image to include in a message. Images are typically provided as base64-encoded strings.
type Message ¶
type Message struct {
Role string `json:"role"`
Content string `json:"content"`
ReasoningContent string `json:"reasoning_content,omitempty"`
ToolCallID string `json:"tool_call_id,omitempty"`
ToolCalls []ToolCall `json:"tool_calls,omitempty"`
Images []ImageData `json:"images,omitempty"`
}
Message represents a single message in a conversation.
type NormalizedToolCalls ¶
type NormalizedToolCalls []ToolCall
NormalizedToolCalls is a slice of ToolCall values that have been validated and cleaned by the ToolCallNormalizer. It is a distinct type to make it clear in API signatures which calls have been normalized.
type Options ¶
type Options struct {
Provider Provider // required — LLM communication
Executor ToolExecutor // required — tool execution
UI UI // nil = headless
SystemPrompt string // empty = minimal default for testing
MaxIterations int // 0 = unlimited
Debug bool
EventPublisher EventPublisher // nil = no events
// OnIteration is an optional callback invoked synchronously at the start of
// each conversation-loop iteration. It receives the iteration number (0-based)
// and the current message count in state (includes prior messages but not
// the LLM response for the current iteration). The agent does not await a
// result or handle errors from this callback; if the callback panics, the
// panic is caught and logged (the agent continues).
OnIteration func(iteration int, messages int)
// Optimizer is used to optimize conversation history across iterations.
Optimizer *ConversationOptimizer
RetryConfig RetryConfig // retry behavior for transient errors; zero values use defaults
// DisableFallbackParser disables the fallback tool-call parser. When
// disabled, malformed tool calls in model responses will not be recovered.
// Default (false): fallback parser is enabled when tools are configured.
DisableFallbackParser bool
// DisableValidator disables the response validator (truncation/tentative
// detection). When disabled, incomplete responses will not trigger
// automatic continuation. Default (false): validator is enabled.
DisableValidator bool
// DisableNormalizer disables the tool call normalizer. When disabled,
// structured tool calls will not be cleaned before execution. Default
// (false): normalizer is enabled.
DisableNormalizer bool
}
Options configures an Agent.
type OutputEvent ¶
type OutputEvent struct {
Type string // "content", "reasoning", "tool_result", "agent_message", "error"
Content string
Source string // origin of the output (e.g., "stream", "provider", "tool")
Timestamp time.Time
Metadata map[string]string
}
OutputEvent represents a generic output event emitted through the async output channel.
type OutputManager ¶
type OutputManager interface {
// Buffer access
ContentBuffer() *StreamingBuffer
ReasoningBuffer() *StreamingBuffer
// Flush callback management
SetFlushCallback(fn func())
Flush()
// Async output channel for goroutine-safe background output delivery
AsyncOutput() <-chan OutputEvent
PublishOutput(event OutputEvent)
// Event metadata (session ID, model, etc.) attached to output events
SetEventMetadata(key string, value string)
GetEventMetadata(key string) string
// Reset clears all buffers and pending async output
Reset()
// Close shuts down the async output channel
Close()
}
OutputManager manages all output streams from the agent, including content and reasoning buffers, async output delivery, flush callbacks, and event metadata. The eventPublisher parameter may be nil (no events) or any EventPublisher implementation.
func NewOutputManager ¶
func NewOutputManager(eventBus EventPublisher) OutputManager
NewOutputManager creates a new OutputManager with the given optional event publisher.
type Provider ¶
type Provider interface {
// Chat sends a chat request and returns the response.
Chat(ctx context.Context, req *ChatRequest) (*ChatResponse, error)
// ChatStream sends a chat request and streams the response via the handler.
ChatStream(ctx context.Context, req *ChatRequest, handler StreamHandler) error
// Info returns metadata about the provider and its model.
Info() ProviderInfo
// EstimateTokens returns an approximate token count for the request.
EstimateTokens(req *ChatRequest) int
}
Provider represents an LLM provider that can be used for chat completions.
type ProviderInfo ¶
type ProviderInfo struct {
Model string `json:"model"`
ContextSize int `json:"context_size"`
HasVision bool `json:"has_vision"`
}
ProviderInfo contains metadata about a provider and its model.
type RateLimitError ¶
type RateLimitError struct {
// Provider is the provider that returned the rate limit error.
Provider string
// RetryAfter is a suggested delay before retrying (from Retry-After header or similar).
RetryAfter time.Duration
// Attempt is the request attempt number when the limit was hit.
Attempt int
// Wrapped is the underlying error.
Wrapped error
}
RateLimitError indicates the provider has rate-limited requests.
func (*RateLimitError) Error ¶
func (e *RateLimitError) Error() string
func (*RateLimitError) Unwrap ¶
func (e *RateLimitError) Unwrap() error
type ResponseValidator ¶
type ResponseValidator struct {
// contains filtered or unexported fields
}
ResponseValidator inspects LLM response content for quality issues like truncation, tentativeness, or other patterns that suggest the response should not be finalized yet.
It has zero dependencies on Agent or concrete types — all input is passed explicitly and the DebugLog callback is optional.
func NewResponseValidator ¶
func NewResponseValidator(opts ResponseValidatorOptions) *ResponseValidator
NewResponseValidator creates a new ResponseValidator with the given options.
func (*ResponseValidator) IsIncomplete ¶
func (rv *ResponseValidator) IsIncomplete(content string) bool
IsIncomplete checks if a response appears to be incomplete or truncated. It returns true if any of these conditions are detected:
- Trailing "..." (ellipsis at end) - Abrupt ending (ends with comma, hyphen, or no punctuation on non-code/URL text) - Unusually short (<10 words and not a known complete-short answer) - Unclosed code blocks (odd number of ``` markers)
func (*ResponseValidator) LooksLikeTentativePostToolResponse ¶
func (rv *ResponseValidator) LooksLikeTentativePostToolResponse(content string) bool
LooksLikeTentativePostToolResponse detects when the LLM has run tools but instead of giving a real response, it's just planning what to do next.
These "tentative" responses should trigger another loop iteration. A response is considered tentative when:
- It is under 40 words (longer responses are considered substantive even if they start with planning language)
- It starts with a planning prefix (case-insensitive), such as "Let me...", "I'll...", "I need to...", "I'm going to...", etc.
func (*ResponseValidator) LooksTruncated ¶
func (rv *ResponseValidator) LooksTruncated(content string) bool
LooksTruncated checks if a response appears structurally truncated. It is a subset of IsIncomplete that excludes the shortness heuristic, making it safe for use in the conversation continuation loop where short but complete answers (e.g., "Done.") should not trigger a retry.
It returns true if any of these conditions are detected:
- Trailing "..." (ellipsis at end) - Abrupt ending (ends with comma or hyphen) - Unclosed code blocks (odd number of ``` markers)
type ResponseValidatorOptions ¶
type ResponseValidatorOptions struct {
// DebugLog is an optional callback for debug output. When nil,
// debug logging is disabled.
DebugLog func(format string, args ...interface{})
}
ResponseValidatorOptions configures a ResponseValidator.
type RetryConfig ¶
type RetryConfig struct {
// MaxAttempts is the total number of attempts (initial + retries).
// Zero means use the default of 3. Setting to 1 means no retries
// (only the initial attempt).
MaxAttempts int
// InitialDelay is the delay before the first retry.
// Zero means use the default of 100ms.
InitialDelay time.Duration
// MaxDelay caps the exponential backoff growth.
// Zero means use the default of 5s.
MaxDelay time.Duration
// Multiplier is the exponential growth factor.
// Zero means use the default of 2.0.
Multiplier float64
// Jitter adds randomness to delays. 0 = none, (0,1) = partial, >=1 = full.
// Zero means use the default of 0.0 (no jitter).
Jitter float64
}
RetryConfig configures retry behavior for transient provider errors. Zero values use sensible defaults.
func (RetryConfig) InitialDelayOrDefault ¶
func (rc RetryConfig) InitialDelayOrDefault() time.Duration
InitialDelayOrDefault returns InitialDelay or the default (100ms).
func (RetryConfig) JitterOrDefault ¶
func (rc RetryConfig) JitterOrDefault() float64
JitterOrDefault returns Jitter or the default (0.0). Unlike other *OrDefault methods, zero is a meaningful value here: Jitter 0.0 means "no jitter" (deterministic retries), so it is not replaced by a default. Negative values are not expected and are returned as-is.
func (RetryConfig) MaxAttemptsOrDefault ¶
func (rc RetryConfig) MaxAttemptsOrDefault() int
MaxAttemptsOrDefault returns MaxAttempts or the default (3).
func (RetryConfig) MaxDelayOrDefault ¶
func (rc RetryConfig) MaxDelayOrDefault() time.Duration
MaxDelayOrDefault returns MaxDelay or the default (5s).
func (RetryConfig) MultiplierOrDefault ¶
func (rc RetryConfig) MultiplierOrDefault() float64
MultiplierOrDefault returns Multiplier or the default (2.0).
type State ¶
type State struct {
// contains filtered or unexported fields
}
State holds the conversation state for an agent.
func (*State) AddCheckpoint ¶
func (s *State) AddCheckpoint(cp TurnCheckpoint)
AddCheckpoint appends a turn checkpoint to the state.
func (*State) AddMessage ¶
AddMessage appends a message to the conversation.
func (*State) ClearCheckpoints ¶
func (s *State) ClearCheckpoints()
ClearCheckpoints removes all checkpoints.
func (*State) EnsureSessionID ¶
func (s *State) EnsureSessionID()
EnsureSessionID generates a session ID if not already set.
func (*State) ExportState ¶
ExportState serializes the state to JSON.
func (*State) GetCheckpoints ¶
func (s *State) GetCheckpoints() []TurnCheckpoint
GetCheckpoints returns a copy of the checkpoint list.
func (*State) ImportState ¶
ImportState deserializes state from JSON.
func (*State) LastAssistantMessage ¶
LastAssistantMessage returns the most recent assistant message, or nil if none exists.
func (*State) SetCheckpoints ¶
func (s *State) SetCheckpoints(cps []TurnCheckpoint)
SetCheckpoints replaces the checkpoint list.
func (*State) SetMessages ¶
SetMessages replaces the message list.
func (*State) SetSessionID ¶
SetSessionID sets the session ID.
func (*State) TotalTokens ¶
TotalTokens returns the total token count.
type StreamHandler ¶
type StreamHandler interface {
OnContent(content string)
OnReasoning(reasoning string)
OnDone(resp *ChatResponse)
OnError(err error)
}
StreamHandler handles streamed responses from a Provider.
type StreamingBuffer ¶
type StreamingBuffer struct {
// contains filtered or unexported fields
}
StreamingBuffer captures streamed output for controlled display.
func NewStreamingBuffer ¶
func NewStreamingBuffer() *StreamingBuffer
NewStreamingBuffer creates a new streaming buffer.
func (*StreamingBuffer) Len ¶
func (b *StreamingBuffer) Len() int
Len returns the current buffer length.
func (*StreamingBuffer) String ¶
func (b *StreamingBuffer) String() string
String returns the current buffer contents.
type Tool ¶
type Tool struct {
Type string `json:"type"`
Function ToolFunction `json:"function"`
}
Tool represents a tool definition that can be provided to the model. The structure matches the OpenAI function-calling wire format where tool definitions are nested under a "function" key.
type ToolCall ¶
type ToolCall struct {
ID string `json:"id"`
Type string `json:"type"`
Function ToolCallFunction `json:"function"`
}
ToolCall represents a function call requested by the model.
type ToolCallFunction ¶
ToolCallFunction represents the function details of a tool call.
type ToolCallNormalizer ¶
type ToolCallNormalizer struct {
// contains filtered or unexported fields
}
ToolCallNormalizer cleans up structured tool calls returned by the model before they are executed. It handles common model output irregularities:
- Strips <|channel|> suffix from tool names
- Generates synthetic IDs for tool calls missing one
- Deduplicates by ID+arguments (first occurrence wins)
- Repairs malformed JSON arguments
- Normalizes Type field to "function"
func NewToolCallNormalizer ¶
func NewToolCallNormalizer() *ToolCallNormalizer
NewToolCallNormalizer creates a new ToolCallNormalizer.
func (*ToolCallNormalizer) Normalize ¶
func (n *ToolCallNormalizer) Normalize(calls []ToolCall) NormalizedToolCalls
Normalize processes a slice of ToolCall values and returns a cleaned, deduplicated NormalizedToolCalls slice. Calls with empty names after normalization or unrepairable JSON arguments are dropped.
type ToolCategory ¶
type ToolCategory int
ToolCategory classifies a tool for optimization purposes.
const ( // ToolCategoryUnknown means the optimizer should skip this tool. ToolCategoryUnknown ToolCategory = iota // ToolCategoryFileRead indicates a tool that reads file contents. ToolCategoryFileRead // ToolCategoryShellCommand indicates a tool that runs shell commands. ToolCategoryShellCommand )
type ToolExecutor ¶
type ToolExecutor interface {
// GetTools returns the list of available tools.
GetTools() []Tool
// Execute runs the given tool calls and returns the resulting messages.
Execute(ctx context.Context, calls []ToolCall) []Message
}
ToolExecutor represents a system that can execute tool calls.
var NoopExecutor ToolExecutor = &noopExecutor{}
NoopExecutor is a ToolExecutor with no tools. Use it when the agent only needs to produce text responses without tool execution.
type ToolFunction ¶ added in v0.1.1
type ToolFunction struct {
Name string `json:"name"`
Description string `json:"description"`
Parameters ToolParameters `json:"parameters"`
}
ToolFunction describes a tool's identity and parameter schema.
type ToolParameter ¶
type ToolParameter struct {
Type string `json:"type"`
Description string `json:"description,omitempty"`
}
ToolParameter defines a single parameter within a tool's schema.
type ToolParameters ¶
type ToolParameters struct {
Type string `json:"type"`
Properties map[string]ToolParameter `json:"properties"`
Required []string `json:"required,omitempty"`
}
ToolParameters defines the schema for a tool's arguments.
type TransientError ¶
type TransientError struct {
// Op is the operation that failed (e.g. "chat", "stream").
Op string
// Provider is the provider that returned the error.
Provider string
// RetryAfter is a suggested delay before retrying (zero means use default backoff).
RetryAfter time.Duration
// Wrapped is the underlying error.
Wrapped error
}
TransientError indicates a temporary failure that may succeed on retry.
func (*TransientError) Error ¶
func (e *TransientError) Error() string
func (*TransientError) Unwrap ¶
func (e *TransientError) Unwrap() error
type TurnCheckpoint ¶
type TurnCheckpoint struct {
// StartIndex is the index of the first message in the turn (the user query).
StartIndex int `json:"start_index"`
// EndIndex is the index of the last message in the turn (the final assistant response).
EndIndex int `json:"end_index"`
// Summary is a concise description of what happened in the turn.
Summary string `json:"summary"`
// ActionableSummary is a bullet-list of accomplishments with file paths,
// commands run, and other concrete details useful for continued context.
ActionableSummary string `json:"actionable_summary"`
}
TurnCheckpoint captures a summary of a completed conversation turn. It records the message range consumed by the turn and a compact summary that can replace the original messages during context compaction.
func BuildCheckpointSummary ¶
func BuildCheckpointSummary(messages []Message) TurnCheckpoint
BuildCheckpointSummary is a convenience function that creates a checkpoint summary from messages without requiring a builder instance.
func ShiftCheckpointIndices ¶
func ShiftCheckpointIndices(oldMessages, newMessages []Message, checkpoints []TurnCheckpoint) []TurnCheckpoint
ShiftCheckpointIndices updates checkpoint StartIndex/EndIndex values after compaction has removed or merged messages. It takes the original message list (before compaction), the compacted message list (after compaction), and the current checkpoints, and returns a new slice of checkpoints with corrected indices.
For each checkpoint:
- If both its StartIndex and EndIndex messages survived compaction, shift them to their new positions in the compacted array.
- If the checkpoint's range was partially consumed (some messages removed), trim the range to only include surviving messages.
- If the entire range was consumed, mark the checkpoint as invalid by setting StartIndex and EndIndex to -1.
The function uses a greedy position-matching algorithm: 1. Walk through old and new message arrays simultaneously 2. Match messages by role, content, and tool_call_id 3. Build a mapping from old index → new index (or -1 if removed) 4. Apply the mapping to each checkpoint
Parameters:
- oldMessages: message list before compaction
- newMessages: message list after compaction
- checkpoints: current checkpoints with old indices
Returns:
- Updated checkpoints with corrected indices
type TurnSummaryBuilder ¶
type TurnSummaryBuilder struct {
// KnownFileTools is a set of tool names that operate on files.
// If nil, the default set is used.
KnownFileTools map[string]bool
// KnownShellTools is a set of tool names that execute shell commands.
// If nil, the default set is used.
KnownShellTools map[string]bool
// KnownErrorPatterns are substrings that indicate a tool result is an error.
// If nil, the default set is used.
KnownErrorPatterns []string
}
TurnSummaryBuilder builds a TurnCheckpoint from a slice of messages representing a single conversation turn. It extracts the user question, tool calls, errors, files modified, and final status to produce both a narrative summary and an actionable bullet list.
func NewTurnSummaryBuilder ¶
func NewTurnSummaryBuilder() *TurnSummaryBuilder
NewTurnSummaryBuilder creates a new builder with default configuration.
func (*TurnSummaryBuilder) Build ¶
func (b *TurnSummaryBuilder) Build(messages []Message) TurnCheckpoint
Build constructs a TurnCheckpoint from the given messages. The messages should represent a single turn: starting with a user query, followed by any number of assistant/tool-call/tool-result cycles, and ending with the final assistant response. Returns a checkpoint with StartIndex=0 and EndIndex=len(messages)-1 since the caller is responsible for setting the actual indices in state.
type UI ¶
type UI interface {
// Prompt displays a prompt and returns the user's input.
Prompt(message string) (string, error)
// Confirm displays a confirmation message and returns the user's choice.
Confirm(message string) (bool, error)
// Print writes output without a trailing newline.
Print(message string)
// PrintLine writes output with a trailing newline.
PrintLine(message string)
}
UI represents a user interface for prompting and output.
var NoopUI UI = &noopUI{}
NoopUI is a headless UI implementation that discards all output and never prompts. Use it when the agent runs without a terminal or interactive layer.
Source Files
¶
- agent.go
- backoff.go
- checkpoint_compaction.go
- checkpoint_shifting.go
- compaction.go
- conversation.go
- conversation_optimizer.go
- error_classifier.go
- errors.go
- fallback_parser.go
- finalize.go
- fp_bare_json.go
- fp_json.go
- fp_named_blocks.go
- fp_tool_blocks.go
- fp_xml.go
- interfaces.go
- message_pipeline.go
- noop.go
- output_manager.go
- response_validator.go
- retry.go
- state.go
- streaming.go
- tool_call_normalizer.go
- turn_checkpoints.go
- turn_summary.go
- types.go