ollie/cmd/olliesrv/internal/agent/agent.go

661 lines
21 KiB
Go

package agent
import (
"context"
"fmt"
"os"
"runtime"
"slices"
"strings"
"sync"
"sync/atomic"
"time"
"ollie/cmd/olliesrv/internal/backend"
toolclient "ollie/cmd/olliesrv/internal/toolclient"
olog "ollie/log"
"ollie/util"
)
// Agent holds the state of the current Agent entity (the "agent"
// in the traditional sense). It is swappable: when the user runs /agent,
// a new Agent is built from the new agent config while the session
// host remains stable.
type Agent struct {
history *History
runtime *Runtime
name string // display name (defaults to id)
profile string // config profile name (e.g. "default")
agentsDir string
systemPrompt string // system prompt for /agent reloads
envBlock string // environment block for /agent reloads
newToolServer func() *toolclient.ToolsrvConn
newBackend func(string) (backend.Backend, error)
currentAction atomic.Pointer[actionHandle]
closed atomic.Bool
actionMu sync.Mutex // serializes action registration with shutdown
warnedContext bool // true after warning about context window limits
resultCache resultCache // per-agent cache of tool results for read-safe tools
// Execution state — owned by the agent, protected by stateMu.
state string // "idle", "thinking", "calling: <tool>"
stateSince time.Time // when the current state began
reply string // last assistant response
getCwd func() string // returns session working directory
cwdOverride string // per-agent cwd override (empty = inherit session cwd)
id string // agent identity (unique principal)
parentID string // immutable ID of the agent that spawned this agent
depth int // sub-agent depth (0=top-level, 1=sub-agent, 2=sub-sub-agent)
activeChildren atomic.Int32 // number of currently-running child sub-agents
peers map[string]struct{} // peer agent names (same session)
peerMu sync.RWMutex
fifo Fifo // prompt queue
toolCallCount atomic.Int64
pendingInject atomic.Pointer[string]
submitMu sync.Mutex // serializes Submit calls (commands + turns)
stateMu sync.RWMutex
signalMu sync.Mutex
signalCh chan struct{} // closed on state change; replaced with a fresh channel
// Event handler — receives events from this agent.
output EventHandler
log *olog.Logger
sessionID string // the owning session's ID
memoryWakePending atomic.Bool
save func() // trigger debounced persistence
flush func() // immediately flush persistence
onStateChange func(agentID, state string) // optional callback for state changes
onProcStart func(agentID string, pid int, tool, cmd string) // optional callback for proc start
onClear func(agentID string) // optional callback for clear
// Chat logs — dual format: JSONL for programmatic access, plain text for humans.
chatMu sync.RWMutex
rawLog []byte // JSONL format (one Block per line)
rawStart int // byte offset of first byte in rawLog (for truncation accounting)
textLog []byte // Rendered plain text
textStart int // byte offset of first byte in textLog
chatVers uint32 // incremented on each append for change detection
chatCond *sync.Cond // signaled on new chat data
chatSignalMu sync.Mutex
chatSignalCh chan struct{} // closed on new chat data; replaced with fresh channel
partialLine []byte // current in-flight (partial) block as a JSONL line; never persisted to rawLog
partialVers uint64 // bumped whenever partialLine changes, for stream change detection
plan []byte // agent's plan file contents
blockCounter uint64 // monotonic counter for block ID generation
}
// Backend returns the active backend from the runtime.
func (ag *Agent) Backend() backend.Backend {
return ag.runtime.Backend
}
// Name returns the agent's display name.
func (ag *Agent) Name() string { return ag.name }
// Profile returns the config profile name (e.g. "default").
func (ag *Agent) Profile() string { return ag.profile }
// Messages returns the conversation history messages.
func (ag *Agent) Messages() []backend.Message {
if ag.history == nil {
return nil
}
return ag.history.history()
}
// UsageStats holds token usage and cost information.
type UsageStats struct {
TotalInputTokens int
TotalCachedInputTokens int
TotalCacheCreationTokens int
TotalOutputTokens int
TotalRequests int
Estimated bool
LastTurnCostUSD float64
SessionCostUSD float64
CacheHitRatio float64
}
// Usage returns the agent's usage statistics.
func (ag *Agent) Usage() *UsageStats {
if ag.history == nil {
return nil
}
return &UsageStats{
TotalInputTokens: ag.history.TotalInputTokens,
TotalCachedInputTokens: ag.history.TotalCachedInputTokens,
TotalCacheCreationTokens: ag.history.TotalCacheCreationTokens,
TotalOutputTokens: ag.history.TotalOutputTokens,
TotalRequests: ag.history.TotalRequests,
Estimated: ag.history.Estimated,
LastTurnCostUSD: ag.history.LastTurnCostUSD,
SessionCostUSD: ag.history.SessionCostUSD,
CacheHitRatio: ag.history.cacheHitRatio(),
}
}
// SetName changes the agent's display name without reloading config.
func (ag *Agent) SetName(name string) {
ag.stateMu.Lock()
ag.name = name
ag.stateMu.Unlock()
}
// ID returns the agent's unique identity.
func (ag *Agent) ID() string { return ag.id }
// ParentID returns the immutable ID of the agent that spawned this agent.
func (ag *Agent) ParentID() string { return ag.parentID }
// BackendName returns the name of the active backend.
func (ag *Agent) BackendName() string {
if ag.runtime.Backend == nil {
return ""
}
return ag.runtime.Backend.Name()
}
// ModelName returns the name of the active model.
func (ag *Agent) ModelName() string {
if ag.runtime.Backend == nil {
return ""
}
return ag.runtime.Backend.Model()
}
// SwitchBackend replaces the active backend by name. The new backend is fully
// constructed before the runtime is changed, so a failed switch leaves the
// current backend usable. Backend changes are only safe between turns.
func (ag *Agent) SwitchBackend(name string) error {
name = strings.TrimSpace(name)
if name == "" {
return fmt.Errorf("backend requires a name")
}
if ag.IsRunning() {
return fmt.Errorf("cannot switch backend while agent is running")
}
if ag.newBackend == nil {
return fmt.Errorf("backend switching unavailable")
}
be, err := ag.newBackend(name)
if err != nil {
return fmt.Errorf("backend %q: %w", name, err)
}
ag.runtime.Backend = be
ag.save()
ag.flush()
ag.notifyChange()
return nil
}
// SetToolServer updates the tool server connection and factory.
// Used when resuming a paused session that was restored without infra.
func (ag *Agent) SetToolServer(newToolServer func() *toolclient.ToolsrvConn, conn *toolclient.ToolsrvConn) {
ag.newToolServer = newToolServer
if ag.runtime != nil {
ag.runtime.ToolServer = conn
}
}
// Cwd returns the agent's effective working directory: the per-agent override
// if one is set, otherwise the inherited session working directory.
func (ag *Agent) Cwd() string {
ag.stateMu.RLock()
override := ag.cwdOverride
ag.stateMu.RUnlock()
if override != "" {
return override
}
if ag.getCwd != nil {
return ag.getCwd()
}
wd, _ := os.Getwd()
return wd
}
// SessionCwd returns the inherited session working directory, ignoring any
// per-agent override.
func (ag *Agent) SessionCwd() string {
if ag.getCwd != nil {
return ag.getCwd()
}
wd, _ := os.Getwd()
return wd
}
// CwdOverride returns the per-agent cwd override, or "" if the agent inherits
// the session working directory.
func (ag *Agent) CwdOverride() string {
ag.stateMu.RLock()
defer ag.stateMu.RUnlock()
return ag.cwdOverride
}
// SetCwdOverride sets (or, with an empty dir, clears) the per-agent working
// directory override and propagates it to the tool server so that this agent's
// tools execute in the chosen directory. Other agents in the session are
// unaffected. The effective directory is also reflected in the preamble
// environment section.
func (ag *Agent) SetCwdOverride(dir string) {
dir = strings.TrimSpace(dir)
ag.stateMu.Lock()
ag.cwdOverride = dir
ag.stateMu.Unlock()
// Refresh the preamble environment section to reflect the effective cwd.
eff := ag.Cwd()
if ag.runtime != nil {
ag.envBlock = EnvironmentBlock(eff, runtime.GOOS, util.IsGitRepo(eff), "")
if ag.runtime.Preamble != nil {
ag.runtime.Preamble.Set(SectionEnv, ag.envBlock)
}
}
ag.SyncCwdToToolServer()
ag.notifyChange()
}
// SyncCwdToToolServer pushes this agent's per-agent cwd override to its current
// tool server. The toolsrv keeps per-agent overrides in process memory, so this
// must be re-run whenever the tool server connection is (re)established — on
// profile switch, resume, or respawn — otherwise the agent silently falls back
// to the session cwd. Pushing an empty override clears any stale entry.
func (ag *Agent) SyncCwdToToolServer() {
ag.stateMu.RLock()
override := ag.cwdOverride
ag.stateMu.RUnlock()
if ts := ag.ToolServer(); ts != nil && ag.id != "" {
ts.SetAgentCWD(ag.id, override)
}
}
// SetGetCwd sets the CWD getter callback.
func (ag *Agent) SetGetCwd(fn func() string) {
ag.getCwd = fn
}
// IsRunning returns true if the agent has an active turn in progress.
func (ag *Agent) IsRunning() bool {
return ag.currentAction.Load() != nil
}
// Interrupt cancels the current in-progress agent turn.
// Returns true if an action was running and was cancelled.
func (ag *Agent) Interrupt(cause error) bool {
ag.actionMu.Lock()
defer ag.actionMu.Unlock()
if h := ag.currentAction.Load(); h != nil {
h.cancel(cause)
return true
}
return false
}
// Compact runs manual context compaction. Returns an error if the agent is
// running or compaction fails.
func (ag *Agent) Compact(ctx context.Context) error {
if ag.IsRunning() {
return fmt.Errorf("cannot compact while agent is running")
}
if ag.history == nil {
return nil
}
ag.SetState("compacting")
_, err := ag.runCompact(ctx, "manual")
ag.SetState("idle")
if err != nil {
return err
}
ag.save()
return nil
}
// Clear resets the agent history and chat log. Returns an error if the agent is running.
func (ag *Agent) Clear() error {
if ag.IsRunning() {
return fmt.Errorf("cannot clear while agent is running")
}
ag.history = nil
ag.chatMu.Lock()
ag.rawLog = nil
ag.rawStart = 0
ag.textLog = nil
ag.textStart = 0
ag.partialLine = nil
ag.partialVers++
ag.chatVers++
ag.chatMu.Unlock()
ag.chatCond.Broadcast()
if ag.onClear != nil {
ag.onClear(ag.id)
}
return nil
}
// SwitchProfile loads a new agent profile by name, replacing the runtime
// and clearing history. Returns the loaded config on success for the caller
// to reload tools. Returns an error if the agent is running or the profile
// cannot be loaded.
func (ag *Agent) SwitchProfile(name string) (*AgentConfig, error) {
if ag.IsRunning() {
return nil, fmt.Errorf("cannot switch agent while agent is running")
}
cfgPath := AgentConfigPath(ag.agentsDir, name)
f, err := os.Open(cfgPath)
if err != nil {
return nil, fmt.Errorf("agent %q: %w", name, err)
}
cfg, err := Load(f)
f.Close()
if err != nil {
return nil, fmt.Errorf("agent %q: %w", name, err)
}
disp := ag.newToolServer()
if disp != nil && ag.id != "" {
disp.SetAgentID(ag.id)
}
env := []string{"OLLIE_SESSION_ID=" + ag.sessionID, "OLLIE_UNAME=" + ag.id}
rt := BuildRuntime(cfg, disp, ag.Cwd(), env, ag.systemPrompt, ag.envBlock)
if cfg.Backend != "" {
newBe, err := ag.newBackend(cfg.Backend)
if err != nil {
return nil, fmt.Errorf("backend %q: %w", cfg.Backend, err)
}
if cfg.Model != "" {
newBe.SetModel(cfg.Model)
}
rt.Backend = newBe
} else {
rt.Backend = ag.runtime.Backend
if cfg.Model != "" {
rt.Backend.SetModel(cfg.Model)
}
}
ag.runtime = rt
ag.profile = name
ag.history = nil
// The profile switch dialed a fresh tool server connection; re-apply the
// per-agent cwd override so it survives the reconnect.
ag.SyncCwdToToolServer()
ag.save()
ag.flush()
ag.notifyChange()
return cfg, nil
}
// Event is a typed output event emitted during an agent turn or in response
// to a command.
// Event carries a single piece of output from the agent loop to consumers
// (frontends, loggers, the turn orchestrator in turn.go).
//
// The Role field determines the event semantics:
//
// Role Name Content Emitted by
// ──────────── ───────── ────────────────────────── ──────────
// "user" — user input text turn.go (before run)
// "assistant" — streamed LLM text chunk loop.go (streamResponse)
// "reasoning" — <think>…</think> chunks loop.go (streamResponse)
// "call" tool name JSON args loop.go (execOne)
// "tool" tool name result text (may stream) loop.go (execOne)
// "usage" — "in out est cost cached creation" loop.go (streamResponse)
// "state" — "thinking"|"compacting"|… loop.go / turn.go
// "limitretry" — — loop.go (rate limit hit)
// "retry" — "HH:MM:SS" countdown loop.go (retryCountdown)
// "maxsteps" — step count as string loop.go (budget exhausted)
// "error" — error message turn.go / loop.go
// "info" — informational text turn.go (compaction, cost)
//
// turn.go intercepts events before forwarding to the external handler:
// - "assistant" → accumulates reply text
// - "call" → updates agent state to "calling: <name>"
// - "tool" → logs result
// - "state" → calls ag.SetState()
// - "limitretry"→ sets state to "limitretry"
// - "usage" → parses token counts, updates history
// - "error" → logs
type Event struct {
Role string
Name string
Content string
ResponseID string
OutputFormat string
}
// EventHandler receives events from the agent.
type EventHandler func(Event)
// actionHandle holds the cancel function for the current agent turn.
type actionHandle struct {
cancel context.CancelCauseFunc
done chan struct{}
}
// SetSessionEnv injects session env vars into the execute server.
func (ag *Agent) SetSessionEnv(sessionID string) {
if ag.runtime == nil || ag.runtime.ToolServer == nil {
return
}
ag.runtime.ToolServer.SetEnv("OLLIE_SESSION_ID", sessionID)
if ag.id != "" {
ag.runtime.ToolServer.SetAgentID(ag.id)
}
if u := os.Getenv("USER"); u != "" {
ag.runtime.ToolServer.SetEnv("USER", u)
}
}
// SetEnv stores an environment variable on the execute server.
func (ag *Agent) SetEnv(key, value string) {
ag.runtime.ToolServer.SetEnv(key, value)
}
// Close releases agent resources (dispatcher, execute server). It first prevents
// new work, cancels the active turn, and waits for that turn to release its
// references to the runtime before closing the tools connection.
func (ag *Agent) Close() {
ag.closed.Store(true)
ag.actionMu.Lock()
h := ag.currentAction.Load()
if h != nil {
h.cancel(context.Canceled)
}
ag.actionMu.Unlock()
if h != nil {
<-h.done
}
if ag.runtime != nil && ag.runtime.ToolServer != nil {
ag.runtime.ToolServer.Close()
}
}
// SetCWD updates the preamble environment section and toolserver CWD.
// Note: This updates local state only. To change session CWD, use session.SetCwd().
func (ag *Agent) SetCWD(dir string) {
isGitRepo := util.IsGitRepo(dir)
ag.envBlock = EnvironmentBlock(dir, runtime.GOOS, isGitRepo, "")
ag.runtime.Preamble.Set(SectionEnv, ag.envBlock)
if ag.runtime.ToolServer != nil {
ag.runtime.ToolServer.SetCWD(dir)
}
ag.notifyChange()
}
// RefreshTools forces a refresh of the tool listing in the agent's preamble.
func (ag *Agent) RefreshTools() {
if infos, err := ag.runtime.ToolServer.ListTools(); err == nil {
ag.runtime.Preamble.Set(SectionTools, RenderTools(infos))
}
}
// SetToolsPreamble updates the tools section of the preamble.
func (ag *Agent) SetToolsPreamble(content string) {
ag.runtime.Preamble.Set(SectionTools, content)
}
// CtxSz returns a human-readable context size string.
func (ag *Agent) CtxSz() string {
if ag.history == nil {
return "no active session"
}
ctxLen := ag.runtime.Backend.ContextLength(context.Background())
if ctxLen <= 0 {
ctxLen = defaultContextLength
}
estimated := ag.history.estimateTokens()
pct := estimated * 100 / ctxLen
return fmt.Sprintf("%d / %d (%d%%)", estimated, ctxLen, pct)
}
// CostStr returns formatted cost information.
func (ag *Agent) CostStr() string {
if ag.history == nil {
return "no active session"
}
return fmt.Sprintf("costLast=$%.4f\ncostSession=$%.4f\n",
ag.history.LastTurnCostUSD, ag.history.SessionCostUSD)
}
// UsageStr returns formatted usage information.
func (ag *Agent) UsageStr() string {
if ag.history == nil {
return "no active session"
}
str := fmt.Sprintf("%d in, %d out, %d requests",
ag.history.TotalInputTokens, ag.history.TotalOutputTokens,
ag.history.TotalRequests)
if ag.history.TotalCachedInputTokens > 0 {
str += fmt.Sprintf(", %d cached", ag.history.TotalCachedInputTokens)
}
if ag.history.Estimated {
str += " [estimated]"
}
return str
}
// Context returns the full message context (system prompt + history).
func (ag *Agent) Context() []backend.Message {
var msgs []backend.Message
if ag.history != nil {
msgs = slices.Clone(ag.history.history())
}
if preamble := ag.runtime.PreambleString(); preamble != "" {
msgs = append([]backend.Message{{Role: "system", Content: preamble}}, msgs...)
}
return msgs
}
// SystemPrompt returns the rendered system prompt.
func (ag *Agent) SystemPrompt() string {
return ag.runtime.PreambleString()
}
// GenParams returns the current generation parameters.
func (ag *Agent) GenParams() backend.GenerationParams {
return ag.runtime.GenParams
}
// SetGenParams sets the generation parameters.
func (ag *Agent) SetGenParams(params backend.GenerationParams) {
ag.runtime.GenParams = params
}
// ListModels returns available models from the backend.
func (ag *Agent) ListModels() []string {
return ag.runtime.Backend.Models(context.Background())
}
// ToolServer returns the tool execution server, or nil if unavailable.
// Exported for use by the 9P filesystem layer to sync tool registries.
func (ag *Agent) ToolServer() *toolclient.ToolsrvConn {
return ag.runtime.ToolServer
}
// SetRuntime replaces the agent's runtime (used on session resume).
func (ag *Agent) SetRuntime(rt *Runtime) {
ag.runtime = rt
}
// SetCompactionModel sets the model used for context compaction.
func (ag *Agent) SetCompactionModel(model string) {
ag.runtime.CompactionModel = model
}
// CompactionModel returns the configured context compaction model.
func (ag *Agent) CompactionModel() string {
return ag.runtime.CompactionModel
}
// Queue pushes a prompt onto the agent's FIFO.
func (ag *Agent) Queue(prompt string) error {
return ag.fifo.Push(prompt)
}
// RequireMemoryWake makes the next top-level turn load persistent memory first.
func (ag *Agent) RequireMemoryWake() {
if ag.depth == 0 {
ag.memoryWakePending.Store(true)
}
}
func (ag *Agent) consumeMemoryWake() bool {
return ag.memoryWakePending.CompareAndSwap(true, false)
}
// PopQueue pops the next prompt from the FIFO.
func (ag *Agent) PopQueue() (string, bool) {
return ag.fifo.Pop()
}
// AgentParams holds the runtime dependencies for constructing a new Agent.
// Not to be confused with AgentConfig, which is the on-disk JSON schema.
type AgentParams struct {
ID string // unique agent identity (uname)
SessionID string // id of session the agent belongs to
Profile string // config profile name (e.g. "default" → agents/default.json)
ParentID string // immutable ID of the agent that spawned this agent
History *History
Runtime *Runtime
AgentsDir string
GetCwd func() string // returns session working directory
SystemPrompt string
EnvBlock string
NewToolServer func() *toolclient.ToolsrvConn
NewBackend func(string) (backend.Backend, error)
Log *olog.Logger
Save func()
Flush func()
}
// NewAgent constructs an Agent from the given configuration.
// Post-construction wiring (SetSessionEnv) is handled internally —
// no additional calls are required after construction.
func NewAgent(cfg AgentParams) *Agent {
ag := &Agent{
history: cfg.History,
runtime: cfg.Runtime,
profile: cfg.Profile,
name: cfg.ID, // display name defaults to uname
agentsDir: cfg.AgentsDir,
id: cfg.ID,
parentID: cfg.ParentID,
getCwd: cfg.GetCwd,
systemPrompt: cfg.SystemPrompt,
envBlock: cfg.EnvBlock,
newToolServer: cfg.NewToolServer,
newBackend: cfg.NewBackend,
log: cfg.Log,
sessionID: cfg.SessionID,
save: cfg.Save,
flush: cfg.Flush,
state: "idle",
}
ag.signalCh = make(chan struct{})
ag.chatSignalCh = make(chan struct{})
ag.chatCond = sync.NewCond(ag.chatMu.RLocker())
ag.initChatHandler()
ag.SetSessionEnv(cfg.SessionID)
return ag
}