480 lines
15 KiB
Go
480 lines
15 KiB
Go
// turn.go — Turn orchestration and main execution entry point.
|
|
//
|
|
// Submit queues a user prompt and starts a turn. executeTurn is the core
|
|
// execution path: resolve tools, call the LLM, handle tool calls, and loop
|
|
// until the model produces a final response or an error limit is reached.
|
|
|
|
package agent
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"runtime/debug"
|
|
"strings"
|
|
|
|
"ollie/cmd/olliesrv/internal/backend"
|
|
"ollie/cmd/olliesrv/internal/metrics"
|
|
"ollie/toolsrv/protocol"
|
|
)
|
|
|
|
// applyTools updates the model-facing tool schemas and textual tool section.
|
|
func (ag *Agent) applyTools(ti []protocol.ToolInfo, revision uint64) {
|
|
ag.runtime.Tools = toolInfosToBackend(ti)
|
|
meta := make(map[string]protocol.ToolInfo, len(ti))
|
|
for _, info := range ti {
|
|
meta[info.Name] = info
|
|
}
|
|
ag.runtime.ToolMeta = meta
|
|
ag.runtime.ToolRevision = revision
|
|
toolPreamble := RenderTools(ti)
|
|
if ag.runtime.Preamble.Get(SectionTools) != toolPreamble {
|
|
ag.runtime.Preamble.Set(SectionTools, toolPreamble)
|
|
}
|
|
}
|
|
|
|
// Submit processes one line of user input: it starts an agent turn that streams
|
|
// events to the bus. If a turn is already in progress the prompt is queued as
|
|
// an in-stream interruption instead.
|
|
//
|
|
// Slash commands are NOT handled here — the caller (Session.Submit) dispatches
|
|
// those before delegating to Agent.Submit.
|
|
//
|
|
// Continuations (post-turn hook context, unconsumed inject, FIFO drain) are
|
|
// handled via an explicit loop rather than recursion to avoid stack growth.
|
|
func (ag *Agent) Submit(ctx context.Context, input string) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
ag.log.Error("panic: %v\n%s", r, debug.Stack())
|
|
if a := ag.currentAction.Swap(nil); a != nil {
|
|
a.cancel(fmt.Errorf("%v", r))
|
|
}
|
|
ag.SetState("idle")
|
|
ag.emit(Event{Role: "error", Name: backend.SeverityFatal, Content: fmt.Sprintf("%v", r)})
|
|
// Re-submit next FIFO item in a new goroutine so queued
|
|
// prompts aren't orphaned by the panic.
|
|
if next, ok := ag.fifo.Pop(); ok {
|
|
go ag.Submit(ctx, next)
|
|
}
|
|
}
|
|
}()
|
|
ag.log.Debug("Agent.Submit() input_len=%d running=%v", len(input), ag.IsRunning())
|
|
if ag.closed.Load() {
|
|
return
|
|
}
|
|
if input == "" {
|
|
return
|
|
}
|
|
|
|
if ag.IsRunning() {
|
|
// While running, slash commands are still dispatched immediately.
|
|
if ag.HandleCommand(ctx, input) {
|
|
return
|
|
}
|
|
ag.fifo.Push(input)
|
|
return
|
|
}
|
|
|
|
// Serialize turns so that e.g. a /compact arriving via ctl cannot race
|
|
// with an executeTurn arriving via prompt.
|
|
ag.submitMu.Lock()
|
|
defer ag.submitMu.Unlock()
|
|
if ag.closed.Load() {
|
|
return
|
|
}
|
|
|
|
if ag.HandleCommand(ctx, input) {
|
|
return
|
|
}
|
|
|
|
if ag.IsRunning() {
|
|
ag.fifo.Push(input)
|
|
return
|
|
}
|
|
|
|
for input != "" && ctx.Err() == nil {
|
|
input = ag.executeTurn(ctx, input)
|
|
}
|
|
// Drain any remaining FIFO items (covers cases where executeTurn
|
|
// returned "" early: interrupt, toolsrv unavailable, etc.).
|
|
for ctx.Err() == nil {
|
|
next, ok := ag.fifo.Pop()
|
|
if !ok {
|
|
break
|
|
}
|
|
for next != "" && ctx.Err() == nil {
|
|
next = ag.executeTurn(ctx, next)
|
|
}
|
|
}
|
|
}
|
|
|
|
// executeTurn runs a single agent turn and returns the next prompt to execute,
|
|
// or "" if there is nothing more to do.
|
|
func (ag *Agent) executeTurn(ctx context.Context, input string) string {
|
|
if ag.consumeMemoryWake() {
|
|
input = "Run the memory_wake tool now. Follow its output completely before addressing my request.\n\n" + input
|
|
}
|
|
|
|
// Fetch tools from toolsrv before constructing the user message. Tool hints
|
|
// must be part of the message stored in history and sent to the backend.
|
|
if ag.runtime.ToolServer != nil {
|
|
ti, listErr := ag.runtime.ToolServer.ListTools()
|
|
if listErr != nil {
|
|
ag.log.Error("tool server unavailable: %v", listErr)
|
|
ag.emit(Event{Role: "error", Name: backend.SeverityFatal, Content: fmt.Sprintf("tool server unavailable: %v\nremediation: the tool server process is not responding; kill and recreate the session", listErr)})
|
|
ag.SetState("idle")
|
|
return ""
|
|
}
|
|
ag.applyTools(ti, ag.runtime.ToolServer.ToolRegistryRevision())
|
|
}
|
|
|
|
// Build context block: user prompts, matched skills, and tool hints.
|
|
var contextParts []string
|
|
if ag.runtime.UserPrompt != "" {
|
|
contextParts = append(contextParts, ag.runtime.UserPrompt)
|
|
}
|
|
if skillContent := matchSkills(input); skillContent != "" {
|
|
contextParts = append(contextParts, skillContent)
|
|
}
|
|
if len(ag.runtime.ToolMeta) > 0 {
|
|
// Convert map to slice for matching
|
|
tools := make([]protocol.ToolInfo, 0, len(ag.runtime.ToolMeta))
|
|
for _, ti := range ag.runtime.ToolMeta {
|
|
tools = append(tools, ti)
|
|
}
|
|
if toolHints := matchTools(input, tools); toolHints != "" {
|
|
contextParts = append(contextParts, toolHints)
|
|
}
|
|
}
|
|
|
|
// Emit context as separate block (filtered by GUI), then user message.
|
|
// History still gets the combined input for LLM context.
|
|
userInput := input
|
|
if len(contextParts) > 0 {
|
|
contextBlock := strings.Join(contextParts, "\n")
|
|
ag.emit(Event{Role: "context", Content: contextBlock})
|
|
input = "<context>\n" + contextBlock + "\n</context>\n\n" + userInput
|
|
}
|
|
ag.emit(Event{Role: "user", Content: userInput})
|
|
|
|
// Snapshot messages before this turn modifies them. Restored on context
|
|
// overflow so the retry starts from a clean state.
|
|
var preTurnLen int
|
|
if ag.history != nil {
|
|
preTurnLen = len(ag.history.messages)
|
|
}
|
|
|
|
if ag.history == nil {
|
|
ag.history = newHistory(input)
|
|
if sc := ag.spawnContext(ctx); sc != "" {
|
|
ag.history.appendUserMessage(sc)
|
|
}
|
|
ag.history.appendUserMessage(input)
|
|
} else {
|
|
ag.history.appendUserMessage(input)
|
|
}
|
|
|
|
actCtx, actCancel := context.WithCancelCause(ctx)
|
|
turnDone := make(chan struct{})
|
|
handle := &actionHandle{cancel: actCancel, done: turnDone}
|
|
ag.actionMu.Lock()
|
|
if ag.closed.Load() {
|
|
ag.actionMu.Unlock()
|
|
actCancel(context.Canceled)
|
|
return ""
|
|
}
|
|
ag.currentAction.Store(handle)
|
|
ag.actionMu.Unlock()
|
|
defer func() {
|
|
if ag.currentAction.CompareAndSwap(handle, nil) {
|
|
actCancel(context.Canceled)
|
|
close(handle.done)
|
|
}
|
|
}()
|
|
ag.SetState("thinking")
|
|
|
|
ag.log.Debug("turn: start input=%s session=%s", auditTruncate(input), ag.sessionID)
|
|
|
|
// Tool schemas and hints were prepared before the user message was built.
|
|
// Install turn-scoped output handler to intercept events for reply
|
|
// accumulation, state updates, and usage tracking.
|
|
var replyBuf strings.Builder
|
|
origOutput := ag.output
|
|
ag.output = func(ev Event) {
|
|
switch ev.Role {
|
|
case "assistant":
|
|
replyBuf.WriteString(ev.Content)
|
|
case "call":
|
|
ag.SetState("calling: " + ev.Name)
|
|
ag.log.Debug("call: %s %s", ev.Name, auditTruncate(string(ev.Content)))
|
|
case "tool":
|
|
ag.log.Debug("result: %s %s", ev.Name, auditTruncate(ev.Content))
|
|
case "state":
|
|
ag.SetState(ev.Content)
|
|
case "limitretry":
|
|
ag.SetState("limitretry")
|
|
case "error":
|
|
ag.log.Debug("error: %s", ev.Content)
|
|
}
|
|
if ev.Role == "usage" && ag.history != nil {
|
|
var in, out, est, cached, creation int
|
|
var costUSD float64
|
|
fmt.Sscanf(ev.Content, "%d %d %d %g %d %d", &in, &out, &est, &costUSD, &cached, &creation)
|
|
usage := backend.Usage{
|
|
InputTokens: in,
|
|
CachedInputTokens: cached,
|
|
CacheCreationTokens: creation,
|
|
OutputTokens: out,
|
|
CostUSD: costUSD,
|
|
}
|
|
ag.history.addUsage(usage, est != 0)
|
|
if err := metrics.RecordUsage(ag.sessionID, ag.id, ag.BackendName(), ag.ModelName(), usage, est != 0); err != nil {
|
|
ag.log.Error("record metrics: %v", err)
|
|
}
|
|
ag.log.Debug("usage: model=%s input=%d cached=%d cache_create=%d output=%d hit_ratio=%.3f cost=$%.6f", ag.runtime.Backend.Model(), in, cached, creation, out, ag.history.cacheHitRatio(), costUSD)
|
|
ag.notifyChange()
|
|
}
|
|
if origOutput != nil {
|
|
origOutput(ev)
|
|
}
|
|
}
|
|
defer func() { ag.output = origOutput }()
|
|
|
|
// Warn once when context usage crosses 60%; compact at 75%.
|
|
if ag.history != nil {
|
|
tokens := ag.history.estimateTokens()
|
|
if compactLimit := ag.autoCompactLimit(actCtx); compactLimit > 0 && tokens >= compactLimit {
|
|
ag.emit(Event{Role: "info", Content: "auto-compacting context...\n"})
|
|
ag.SetState("compacting")
|
|
if _, err := ag.runCompact(actCtx, "auto"); err != nil {
|
|
ag.log.Error("auto-compact failed: %v", err)
|
|
}
|
|
ag.SetState("thinking")
|
|
} else if warnLimit := ag.autoWarnLimit(actCtx); warnLimit > 0 && tokens >= warnLimit && !ag.warnedContext {
|
|
ctxLen := ag.runtime.Backend.ContextLength(ctx)
|
|
if ctxLen <= 0 {
|
|
ctxLen = defaultContextLength
|
|
}
|
|
pct := tokens * 100 / ctxLen
|
|
ag.emit(Event{Role: "info", Content: fmt.Sprintf("context at %d%% — will auto-compact at 75%%\n", pct)})
|
|
ag.warnedContext = true
|
|
}
|
|
}
|
|
|
|
if ag.history != nil {
|
|
ag.history.resetTurnAccumulators()
|
|
}
|
|
|
|
// Strip cold-zone tool results (one batched LLM call per turn).
|
|
if ag.history != nil {
|
|
pendingCount, pendingTokens := ag.history.pendingColdSummaryStats()
|
|
ctxLen := ag.runtime.Backend.ContextLength(actCtx)
|
|
if ctxLen <= 0 {
|
|
ctxLen = defaultContextLength
|
|
}
|
|
ctxTokens := ag.history.estimateTokens()
|
|
forceStripCold := ctxTokens*100 >= ctxLen*40
|
|
if (pendingCount >= minColdSummaries && pendingTokens >= minColdSummaryTokens) || forceStripCold {
|
|
usage, estimated := ag.history.stripCold(actCtx, ag.runtime.Backend)
|
|
if usage.InputTokens > 0 || usage.OutputTokens > 0 {
|
|
ag.history.addUsage(usage, estimated)
|
|
ag.notifyChange()
|
|
}
|
|
}
|
|
}
|
|
|
|
// Inject the preamble system message into history before running.
|
|
if ag.history != nil {
|
|
ag.history.injectPreamble(ag.runtime.PreambleString())
|
|
}
|
|
|
|
// Run the turn, retrying once after compaction on context overflow.
|
|
var (
|
|
overflowRetried bool
|
|
err error
|
|
)
|
|
for {
|
|
err = ag.run(actCtx)
|
|
if ag.currentAction.CompareAndSwap(handle, nil) {
|
|
actCancel(nil)
|
|
close(handle.done)
|
|
}
|
|
|
|
if err == nil {
|
|
break
|
|
}
|
|
if errors.Is(err, context.Canceled) || errors.Is(err, ErrInterrupted) {
|
|
break
|
|
}
|
|
var ctxErr *backend.ContextOverflowError
|
|
if !overflowRetried && errors.As(err, &ctxErr) && ag.history != nil {
|
|
overflowRetried = true
|
|
ag.history.messages = ag.history.messages[:preTurnLen]
|
|
ag.emit(Event{Role: "info", Content: "context overflow — compacting and retrying...\n"})
|
|
ag.SetState("compacting")
|
|
if _, cerr := ag.runCompact(context.WithoutCancel(actCtx), "overflow"); cerr != nil {
|
|
break
|
|
}
|
|
if actCtx.Err() != nil {
|
|
err = actCtx.Err()
|
|
break
|
|
}
|
|
ag.history.appendUserMessage(input)
|
|
ag.SetState("thinking")
|
|
ag.history.resetTurnAccumulators()
|
|
replyBuf.Reset()
|
|
actCtx, actCancel = context.WithCancelCause(ctx)
|
|
handle = &actionHandle{cancel: actCancel, done: make(chan struct{})}
|
|
ag.actionMu.Lock()
|
|
if ag.closed.Load() {
|
|
ag.actionMu.Unlock()
|
|
actCancel(context.Canceled)
|
|
break
|
|
}
|
|
ag.currentAction.Store(handle)
|
|
ag.actionMu.Unlock()
|
|
// Re-inject preamble — compaction strips it from history.
|
|
ag.history.injectPreamble(ag.runtime.PreambleString())
|
|
continue
|
|
}
|
|
break
|
|
}
|
|
|
|
ag.setReply(replyBuf.String())
|
|
replyBuf.Reset()
|
|
ag.SetState("idle")
|
|
ag.flush()
|
|
|
|
if err != nil {
|
|
if ag.history != nil {
|
|
ag.history.recordTurnCost(ag.runtime.Backend.Model())
|
|
}
|
|
// Keep completed work — only remove cancelled tool results.
|
|
if ag.history != nil {
|
|
ag.history.removeCancelledToolResults()
|
|
}
|
|
if errors.Is(err, context.Canceled) || errors.Is(err, ErrInterrupted) {
|
|
ag.log.Debug("turn: interrupted session=%s", ag.sessionID)
|
|
ag.save()
|
|
return ""
|
|
}
|
|
severity, remediation := backend.ClassifyError(err)
|
|
content := err.Error()
|
|
if remediation != "" {
|
|
content += "\nremediation: " + remediation
|
|
}
|
|
ag.emit(Event{Role: "error", Name: severity, Content: content})
|
|
// Drain one FIFO item.
|
|
if next, ok := ag.fifo.Pop(); ok {
|
|
return next
|
|
}
|
|
return ""
|
|
}
|
|
|
|
if ag.history != nil {
|
|
ag.history.recordTurnCost(ag.runtime.Backend.Model())
|
|
if ag.history.LastTurnCostUSD > 0 {
|
|
ag.emit(Event{Role: "info", Content: fmt.Sprintf("costLast=$%.4f\n", ag.history.LastTurnCostUSD)})
|
|
}
|
|
ag.log.Debug("turn: end reply=%s cost=$%.4f session_total=$%.4f session=%s",
|
|
auditTruncate(ag.Reply()), ag.history.LastTurnCostUSD, ag.history.SessionCostUSD, ag.sessionID)
|
|
ag.notifyChange()
|
|
}
|
|
ag.save()
|
|
|
|
// Inject that was pending but never consumed (text-only response with no
|
|
// tool calls) — treat it as the next user message.
|
|
if p := ag.pendingInject.Swap(nil); p != nil {
|
|
return *p
|
|
}
|
|
|
|
// Drain one item from the FIFO; the outer loop handles the rest.
|
|
if next, ok := ag.fifo.Pop(); ok {
|
|
return next
|
|
}
|
|
|
|
return ""
|
|
}
|
|
|
|
// autoCompactLimit returns the token threshold for auto-compaction (75% of
|
|
// the compaction model's context window). Compaction may use a different model
|
|
// from the active turn model, so querying the active model here can overestimate
|
|
// the request limit.
|
|
func (ag *Agent) autoCompactLimit(ctx context.Context) int {
|
|
b := ag.runtime.Backend
|
|
compactModel := resolveCompactionModel(ag.runtime.CompactionModel, b)
|
|
origModel := b.Model()
|
|
if compactModel != "" && compactModel != origModel {
|
|
b.SetModel(compactModel)
|
|
defer b.SetModel(origModel)
|
|
}
|
|
ctxLen := b.ContextLength(ctx)
|
|
if ctxLen <= 0 {
|
|
ctxLen = defaultContextLength
|
|
}
|
|
return ctxLen * 3 / 4
|
|
}
|
|
|
|
// autoWarnLimit returns the token threshold for a context-usage warning (60%).
|
|
func (ag *Agent) autoWarnLimit(ctx context.Context) int {
|
|
ctxLen := ag.runtime.Backend.ContextLength(ctx)
|
|
if ctxLen <= 0 {
|
|
ctxLen = defaultContextLength
|
|
}
|
|
return ctxLen * 3 / 5
|
|
}
|
|
|
|
// spawnContext assembles the agent context injected at each session refresh
|
|
// point (session start, post-clear, post-compaction).
|
|
func (ag *Agent) spawnContext(ctx context.Context) string {
|
|
var parts []string
|
|
// Inject AGENTS.md from the working directory if it exists.
|
|
// Repo-provided content is untrusted data: wrap it so consumers and the
|
|
// model can distinguish it from user/system instructions (KDE filters
|
|
// <context> blocks in chat rendering).
|
|
if cwd := ag.Cwd(); cwd != "" {
|
|
if data, err := os.ReadFile(cwd + "/AGENTS.md"); err == nil && len(data) > 0 {
|
|
parts = append(parts, "<context>\n"+string(data)+"\n</context>")
|
|
}
|
|
}
|
|
return strings.Join(parts, "\n\n---\n\n")
|
|
}
|
|
|
|
// runCompact executes a full compaction cycle: compact, spawn-context
|
|
// re-injection. Returns (n compacted, error). Returns (0, nil) if
|
|
// there was nothing to compact. Caller manages setState.
|
|
func (ag *Agent) runCompact(ctx context.Context, trigger string) (int, error) {
|
|
// Use a cheaper model for compaction if configured.
|
|
compactModel := resolveCompactionModel(ag.runtime.CompactionModel, ag.runtime.Backend)
|
|
origModel := ag.runtime.Backend.Model()
|
|
if compactModel != "" && compactModel != origModel {
|
|
ag.runtime.Backend.SetModel(compactModel)
|
|
defer ag.runtime.Backend.SetModel(origModel)
|
|
}
|
|
n, _, err := ag.history.compact(ctx, ag.runtime.Backend)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
if n > 0 {
|
|
ag.log.Debug("compact: removed %d messages trigger=%s session=%s", n, trigger, ag.sessionID)
|
|
ag.warnedContext = false
|
|
if sc := ag.spawnContext(ctx); sc != "" {
|
|
ag.history.appendUserMessage(sc)
|
|
}
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
// ErrInterrupted is returned when the user cancels an agent turn (Ctrl-C).
|
|
var ErrInterrupted = errors.New("interrupted")
|
|
|
|
// defaultContextLength is used when the backend cannot report the model's
|
|
// actual context window. 128k tokens is a safe default for modern models.
|
|
const defaultContextLength = 128000
|
|
|
|
// infoEvent wraps a plain-text message as an info Event.
|
|
func infoEvent(text string) Event {
|
|
return Event{Role: "info", Content: text + "\n"}
|
|
}
|