diff --git a/agent/commands.go b/agent/commands.go index 106b004..c558f17 100644 --- a/agent/commands.go +++ b/agent/commands.go @@ -12,7 +12,7 @@ import ( ) -func (s *agent) handleCommand(ctx context.Context, input string) bool { +func (a *agent) handleCommand(ctx context.Context, input string) bool { if !strings.HasPrefix(input, "/") { return false } @@ -25,9 +25,9 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { args := parts[1:] listFromHandler := func(name string) { - if h := s.listHandlers[name]; h != nil { + if h := a.listHandlers[name]; h != nil { for _, item := range h() { - s.emit(infoEvent(" " + item)) + a.emit(infoEvent(" " + item)) } } } @@ -37,88 +37,88 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { "/i": func(args []string) { prompt := strings.Join(args, " ") if prompt == "" { - s.emit(infoEvent("error: /i requires a prompt")) + a.emit(infoEvent("error: /i requires a prompt")) return } - if s.IsRunning() { - s.Inject(prompt) + if a.IsRunning() { + a.Inject(prompt) } else { - go s.Submit(context.Background(), prompt) + go a.Submit(context.Background(), prompt) } }, "/irw": func(args []string) { prompt := strings.Join(args, " ") if prompt == "" { - s.emit(infoEvent("error: /irw requires a prompt")) + a.emit(infoEvent("error: /irw requires a prompt")) return } - s.injectRewrite(prompt) + a.injectRewrite(prompt) }, "/backend": func(args []string) { if len(args) == 0 { - s.emit(infoEvent(s.runtime.Backend.Name())) + a.emit(infoEvent(a.runtime.Backend.Name())) return } - if s.IsRunning() { - s.emit(infoEvent("error: cannot switch backend while agent is running")) + if a.IsRunning() { + a.emit(infoEvent("error: cannot switch backend while agent is running")) return } - be, err := s.newBackend(args[0]) + be, err := a.newBackend(args[0]) if err != nil { - s.emit(infoEvent(fmt.Sprintf("error: failed to switch backend: %v", err))) + a.emit(infoEvent(fmt.Sprintf("error: failed to switch backend: %v", err))) return } - s.runtime.Backend = be - s.emit(infoEvent(fmt.Sprintf("switched backend to: %s (model: %s)", be.Name(), be.Model()))) + a.runtime.Backend = be + a.emit(infoEvent(fmt.Sprintf("switched backend to: %s (model: %s)", be.Name(), be.Model()))) }, "/models": func(args []string) { - models := s.runtime.Backend.Models(ctx) + models := a.runtime.Backend.Models(ctx) if len(models) == 0 { - s.emit(infoEvent("no models available")) + a.emit(infoEvent("no models available")) return } slices.Sort(models) - current := s.runtime.Backend.Model() + current := a.runtime.Backend.Model() for _, m := range models { marker := " " if m == current { marker = "* " } - s.emit(infoEvent(marker + m)) + a.emit(infoEvent(marker + m)) } }, "/model": func(args []string) { if len(args) == 0 { - s.emit(infoEvent(s.runtime.Backend.Model())) + a.emit(infoEvent(a.runtime.Backend.Model())) return } - s.runtime.Backend.SetModel(args[0]) - s.emit(infoEvent("switched model to: " + args[0])) + a.runtime.Backend.SetModel(args[0]) + a.emit(infoEvent("switched model to: " + args[0])) }, "/maxsteps": func(args []string) { if len(args) == 0 { - if s.runtime.MaxSteps == 0 { - s.emit(infoEvent("maxsteps: unlimited")) + if a.runtime.MaxSteps == 0 { + a.emit(infoEvent("maxsteps: unlimited")) } else { - s.emit(infoEvent(fmt.Sprintf("maxsteps: %d", s.runtime.MaxSteps))) + a.emit(infoEvent(fmt.Sprintf("maxsteps: %d", a.runtime.MaxSteps))) } return } n, err := strconv.Atoi(args[0]) if err != nil || n < 0 { - s.emit(infoEvent("error: maxsteps requires a non-negative integer (0 = unlimited)")) + a.emit(infoEvent("error: maxsteps requires a non-negative integer (0 = unlimited)")) return } - s.runtime.MaxSteps = n + a.runtime.MaxSteps = n if n == 0 { - s.emit(infoEvent("maxsteps: unlimited")) + a.emit(infoEvent("maxsteps: unlimited")) } else { - s.emit(infoEvent(fmt.Sprintf("maxsteps: %d", n))) + a.emit(infoEvent(fmt.Sprintf("maxsteps: %d", n))) } }, @@ -140,48 +140,48 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { } seen[name] = true marker := " " - if name == s.agentName { + if name == a.agentName { marker = "* " } - s.emit(infoEvent(marker + name)) + a.emit(infoEvent(marker + name)) found = true } } if !found { - s.emit(infoEvent("no agents found")) + a.emit(infoEvent("no agents found")) } }, "/agent": func(args []string) { if len(args) == 0 { - s.emit(infoEvent("active agent: " + s.agentName)) + a.emit(infoEvent("active agent: " + a.agentName)) return } - if s.IsRunning() { - s.emit(infoEvent("error: cannot switch agent while agent is running")) + if a.IsRunning() { + a.emit(infoEvent("error: cannot switch agent while agent is running")) return } name := args[0] - cfgPath := AgentConfigPath(s.agentsDir, name) + cfgPath := AgentConfigPath(a.agentsDir, name) f, err := os.Open(cfgPath) if err != nil { - s.emit(infoEvent(fmt.Sprintf("error: agent %q: %v", name, err))) + a.emit(infoEvent(fmt.Sprintf("error: agent %q: %v", name, err))) return } cfg, err := Load(f) f.Close() if err != nil { - s.emit(infoEvent(fmt.Sprintf("error: agent %q: %v", name, err))) + a.emit(infoEvent(fmt.Sprintf("error: agent %q: %v", name, err))) return } - d := s.newDispatcher() - env := []string{"OLLIE_SESSION_ID=" + s.sess.ID(), "OLLIE_UNAME=" + s.sess.Uname()} - env = append(env, s.promptEnvExtra...) - rt := BuildRuntime(cfg, d, s.sess.CWD(), env, s.baseLayers...) + d := a.newDispatcher() + env := []string{"OLLIE_SESSION_ID=" + a.sess.ID(), "OLLIE_UNAME=" + a.sess.Uname()} + env = append(env, a.promptEnvExtra...) + rt := BuildRuntime(cfg, d, a.sess.CWD(), env, a.baseLayers...) if rt.CfgBackend != "" { - newBe, err := s.newBackend(rt.CfgBackend) + newBe, err := a.newBackend(rt.CfgBackend) if err != nil { - s.emit(infoEvent(fmt.Sprintf("error: backend %q: %v", rt.CfgBackend, err))) + a.emit(infoEvent(fmt.Sprintf("error: backend %q: %v", rt.CfgBackend, err))) return } if rt.CfgModel != "" { @@ -189,47 +189,47 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { } rt.Backend = newBe } else { - rt.Backend = s.runtime.Backend + rt.Backend = a.runtime.Backend if rt.CfgModel != "" { rt.Backend.SetModel(rt.CfgModel) } } - s.runtime = rt - s.agentName = name - s.history = nil - s.pushSessionEnv() - s.notifyChange() + a.runtime = rt + a.agentName = name + a.history = nil + a.pushSessionEnv() + a.notifyChange() for _, msg := range rt.Messages { - s.emit(infoEvent(msg)) + a.emit(infoEvent(msg)) } - s.emit(infoEvent("agent: " + name)) + a.emit(infoEvent("agent: " + name)) }, "/compact": func(args []string) { - if s.IsRunning() { - s.emit(infoEvent("error: cannot compact while agent is running")) + if a.IsRunning() { + a.emit(infoEvent("error: cannot compact while agent is running")) return } - if s.history == nil { - s.emit(infoEvent("nothing to compact")) + if a.history == nil { + a.emit(infoEvent("nothing to compact")) return } - snapshot := s.history.PreCompactionSnapshot() - s.setState("compacting") - n, err := s.runCompact(ctx, "manual") - s.setState("idle") + snapshot := a.history.PreCompactionSnapshot() + a.setState("compacting") + n, err := a.runCompact(ctx, "manual") + a.setState("idle") if err != nil { - s.emit(infoEvent("compact error: " + err.Error())) + a.emit(infoEvent("compact error: " + err.Error())) return } if n == 0 { - s.emit(infoEvent("nothing to compact")) + a.emit(infoEvent("nothing to compact")) return } - if s.sessionsDir != "" && s.sess.ID() != "" { - histPath := s.activeSessionPath(s.sess.ID(), ".compaction.jsonl") + if a.sessionsDir != "" && a.sess.ID() != "" { + histPath := a.activeSessionPath(a.sess.ID(), ".compaction.jsonl") if err := os.MkdirAll(filepath.Dir(histPath), 0700); err != nil { - s.emit(infoEvent("compaction history save: " + err.Error())) + a.emit(infoEvent("compaction history save: " + err.Error())) } else if data, err := json.Marshal(snapshot); err == nil { f, err := os.OpenFile(histPath, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0600) if err == nil { @@ -238,76 +238,76 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { } } } - s.emit(infoEvent(fmt.Sprintf("compacted %d messages", n))) - s.saveSession() + a.emit(infoEvent(fmt.Sprintf("compacted %d messages", n))) + a.saveSession() }, "/context": func(args []string) { - if s.history == nil { - s.emit(infoEvent("no active session")) + if a.history == nil { + a.emit(infoEvent("no active session")) return } - ctxLen := s.runtime.Backend.ContextLength(ctx) + ctxLen := a.runtime.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength } - estimated := s.history.estimateTokens() + estimated := a.history.estimateTokens() pct := estimated * 100 / ctxLen - s.emit(infoEvent(fmt.Sprintf("~%d / %d tokens (%d%%)", estimated, ctxLen, pct))) - s.emit(infoEvent(strings.TrimRight(s.history.contextDebug(), "\n"))) + a.emit(infoEvent(fmt.Sprintf("~%d / %d tokens (%d%%)", estimated, ctxLen, pct))) + a.emit(infoEvent(strings.TrimRight(a.history.contextDebug(), "\n"))) }, "/cost": func(args []string) { - if s.history == nil { - s.emit(infoEvent("no active session")) + if a.history == nil { + a.emit(infoEvent("no active session")) return } - s.emit(infoEvent(fmt.Sprintf("last=$%.4f session=$%.4f", - s.history.LastTurnCostUSD, s.history.SessionCostUSD))) + a.emit(infoEvent(fmt.Sprintf("last=$%.4f session=$%.4f", + a.history.LastTurnCostUSD, a.history.SessionCostUSD))) }, "/usage": func(args []string) { - if s.history == nil { - s.emit(infoEvent("no active session")) + if a.history == nil { + a.emit(infoEvent("no active session")) return } - ctxLen := s.runtime.Backend.ContextLength(ctx) + ctxLen := a.runtime.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength } - estimated := s.history.estimateTokens() + estimated := a.history.estimateTokens() pct := estimated * 100 / ctxLen usageStr := fmt.Sprintf("~%d / %d tokens (%d%%) | %d in, %d out, %d requests", estimated, ctxLen, pct, - s.history.TotalInputTokens, s.history.TotalOutputTokens, - s.history.TotalRequests) - if s.history.Estimated { + a.history.TotalInputTokens, a.history.TotalOutputTokens, + a.history.TotalRequests) + if a.history.Estimated { usageStr += " [estimated]" } - s.emit(infoEvent(usageStr)) + a.emit(infoEvent(usageStr)) }, "/history": func(args []string) { - if s.history == nil { - s.emit(infoEvent("no active session")) + if a.history == nil { + a.emit(infoEvent("no active session")) return } - for _, msg := range s.history.history() { + for _, msg := range a.history.history() { preview := msg.Content if len(preview) > 200 { preview = preview[:200] + "..." } - s.emit(infoEvent(fmt.Sprintf("[%s] %s", msg.Role, preview))) + a.emit(infoEvent(fmt.Sprintf("[%s] %s", msg.Role, preview))) } }, "/clear": func(args []string) { - if s.IsRunning() { - s.emit(infoEvent("error: cannot clear while agent is running")) + if a.IsRunning() { + a.emit(infoEvent("error: cannot clear while agent is running")) return } - s.history = nil - s.emit(infoEvent("cleared")) + a.history = nil + a.emit(infoEvent("cleared")) }, "/sessions": func(args []string) { @@ -333,10 +333,10 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { }) } } - appendSessionFiles(s.sessionsDir) - appendSessionFiles(filepath.Join(s.sessionsDir, "active")) + appendSessionFiles(a.sessionsDir) + appendSessionFiles(filepath.Join(a.sessionsDir, "active")) - cwd := s.CWD() + cwd := a.CWD() found := false for _, file := range files { data, readErr := os.ReadFile(file.path) @@ -351,7 +351,7 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { continue } marker := " " - if file.id == s.sess.ID() { + if file.id == a.sess.ID() { marker = "* " } goal := "" @@ -364,75 +364,75 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { if len(goal) > 60 { goal = goal[:60] + "..." } - s.emit(infoEvent(marker + fmt.Sprintf("%-24s [%s] %q", file.id, ps.Agent, goal))) + a.emit(infoEvent(marker + fmt.Sprintf("%-24s [%s] %q", file.id, ps.Agent, goal))) found = true } if !found { - s.emit(infoEvent("no sessions for " + cwd)) + a.emit(infoEvent("no sessions for " + cwd)) } }, "/save": func(args []string) { - if s.history == nil { - s.emit(infoEvent("error: no active session")) + if a.history == nil { + a.emit(infoEvent("error: no active session")) return } if len(args) == 0 { - s.emit(infoEvent("error: /save requires a name")) + a.emit(infoEvent("error: /save requires a name")) return } name := args[0] - path := s.sessionsDir + "/" + name + ".json" - if err := s.history.saveTo(path, name, s.agentName, s.CWD()); err != nil { - s.emit(infoEvent("error: " + err.Error())) + path := a.sessionsDir + "/" + name + ".json" + if err := a.history.saveTo(path, name, a.agentName, a.CWD()); err != nil { + a.emit(infoEvent("error: " + err.Error())) return } - s.emit(infoEvent("saved: " + path)) + a.emit(infoEvent("saved: " + path)) }, "/resume": func(args []string) { if len(args) == 0 { - s.emit(infoEvent("error: /resume requires a session id or name")) + a.emit(infoEvent("error: /resume requires a session id or name")) return } - if s.IsRunning() { - s.emit(infoEvent("error: cannot resume while agent is running")) + if a.IsRunning() { + a.emit(infoEvent("error: cannot resume while agent is running")) return } name := args[0] - path := s.sessionsDir + "/" + name + ".json" + path := a.sessionsDir + "/" + name + ".json" data, err := os.ReadFile(path) if err != nil { - s.emit(infoEvent(fmt.Sprintf("error: %v", err))) + a.emit(infoEvent(fmt.Sprintf("error: %v", err))) return } var ps PersistedSession if err := json.Unmarshal(data, &ps); err != nil { - s.emit(infoEvent(fmt.Sprintf("error: %v", err))) + a.emit(infoEvent(fmt.Sprintf("error: %v", err))) return } - s.history = RestoreHistory(&ps) - s.emit(infoEvent(fmt.Sprintf("resumed session %s (%d messages)", name, len(ps.Messages)))) + a.history = RestoreHistory(&ps) + a.emit(infoEvent(fmt.Sprintf("resumed session %s (%d messages)", name, len(ps.Messages)))) }, "/cwd": func(args []string) { if len(args) == 0 { - s.emit(infoEvent("cwd: " + s.CWD())) + a.emit(infoEvent("cwd: " + a.CWD())) return } dir := strings.Join(args, " ") - if err := s.SetCWD(dir); err != nil { - s.emit(infoEvent("error: " + err.Error())) + if err := a.SetCWD(dir); err != nil { + a.emit(infoEvent("error: " + err.Error())) return } - s.emit(infoEvent("cwd: " + dir)) + a.emit(infoEvent("cwd: " + dir)) }, "/skills": func(args []string) { listFromHandler("skills") }, "/tools": func(args []string) { listFromHandler("tools") }, "/sp": func(args []string) { - s.emit(infoEvent(s.runtime.Preamble)) + a.emit(infoEvent(a.runtime.Preamble)) }, "/help": func(args []string) { @@ -466,7 +466,7 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { " ! - run shell command", } for _, l := range lines { - s.emit(infoEvent(l)) + a.emit(infoEvent(l)) } }, } @@ -475,7 +475,7 @@ func (s *agent) handleCommand(ctx context.Context, input string) bool { if !ok { return false } - s.emit(infoEvent("")) + a.emit(infoEvent("")) fn(args) return true } diff --git a/agent/core.go b/agent/core.go index b16069a..5b8584a 100644 --- a/agent/core.go +++ b/agent/core.go @@ -349,17 +349,17 @@ type agent struct { // ToolCallCount returns the total number of tool calls executed in this // session. The counter is monotonically increasing and never resets. // Blocked calls (pre-tool hook exit 2) are not counted. -func (s *agent) ToolCallCount() int64 { - return s.toolCallCount.Load() +func (a *agent) ToolCallCount() int64 { + return a.toolCallCount.Load() } // SetEnv stores a session-scoped variable and propagates it to the execute server. -func (s *agent) SetEnv(key, value string) { - s.sess.SetEnv(key, value) - if s.runtime == nil || s.runtime.Dispatcher == nil { +func (a *agent) SetEnv(key, value string) { + a.sess.SetEnv(key, value) + if a.runtime == nil || a.runtime.Dispatcher == nil { return } - if srv, ok := s.runtime.Dispatcher.GetServer("execute"); ok { + if srv, ok := a.runtime.Dispatcher.GetServer("execute"); ok { if es, ok := srv.(tools.EnvSetter); ok { es.SetEnv(key, value) } @@ -367,23 +367,23 @@ func (s *agent) SetEnv(key, value string) { } // pushSessionEnv injects OLLIE_SESSION_ID into the execute server subprocess env. -func (s *agent) pushSessionEnv() { - if s.runtime == nil || s.runtime.Dispatcher == nil || s.sess.ID() == "" { +func (a *agent) pushSessionEnv() { + if a.runtime == nil || a.runtime.Dispatcher == nil || a.sess.ID() == "" { return } - if srv, ok := s.runtime.Dispatcher.GetServer("execute"); ok { + if srv, ok := a.runtime.Dispatcher.GetServer("execute"); ok { if es, ok := srv.(tools.EnvSetter); ok { - es.SetEnv("OLLIE_SESSION_ID", s.sess.ID()) - if s.sess.Uname() != "" { - es.SetEnv("OLLIE_UNAME", s.sess.Uname()) + es.SetEnv("OLLIE_SESSION_ID", a.sess.ID()) + if a.sess.Uname() != "" { + es.SetEnv("OLLIE_UNAME", a.sess.Uname()) } } } } // pushLockDir sets the flock directory on the execute server to the session tmpdir. -func (s *agent) pushLockDir() { - if s.runtime == nil || s.runtime.Dispatcher == nil || s.sess.ID() == "" { +func (a *agent) pushLockDir() { + if a.runtime == nil || a.runtime.Dispatcher == nil || a.sess.ID() == "" { return } @@ -505,33 +505,33 @@ func NewAgentCore(cfg AgentCoreConfig) Core { } // Close releases resources for this session, including its tmpdir. -func (s *agent) Close() { - s.log.Debug("Close() session=%q", s.sess.ID()) - s.flushSave() - if s.runtime != nil && s.runtime.Dispatcher != nil { - if srv, ok := s.runtime.Dispatcher.GetServer("execute"); ok { +func (a *agent) Close() { + a.log.Debug("Close() session=%q", a.sess.ID()) + a.flushSave() + if a.runtime != nil && a.runtime.Dispatcher != nil { + if srv, ok := a.runtime.Dispatcher.GetServer("execute"); ok { if c, ok := srv.(interface{ Close() }); ok { - s.log.Debug("Close() calling execute.Close()") + a.log.Debug("Close() calling execute.Close()") c.Close() } } } - if s.sess.ID() != "" { - os.RemoveAll(filepath.Join(ollieTmpDir(), s.sess.ID())) //nolint:errcheck + if a.sess.ID() != "" { + os.RemoveAll(filepath.Join(ollieTmpDir(), a.sess.ID())) //nolint:errcheck } } // execServer returns the execute server if available, or nil. -func (s *agent) execServer() interface{} { - if s.runtime == nil || s.runtime.Dispatcher == nil { +func (a *agent) execServer() interface{} { + if a.runtime == nil || a.runtime.Dispatcher == nil { return nil } - srv, _ := s.runtime.Dispatcher.GetServer("execute") + srv, _ := a.runtime.Dispatcher.GetServer("execute") return srv } -func (s *agent) Detach() bool { - if srv := s.execServer(); srv != nil { +func (a *agent) Detach() bool { + if srv := a.execServer(); srv != nil { if d, ok := srv.(interface{ Detach() bool }); ok { return d.Detach() } @@ -539,8 +539,8 @@ func (s *agent) Detach() bool { return false } -func (s *agent) ListDetached() []DetachedInfo { - if srv := s.execServer(); srv != nil { +func (a *agent) ListDetached() []DetachedInfo { + if srv := a.execServer(); srv != nil { type listDetacher interface { ListDetachedRaw() []any } @@ -574,8 +574,8 @@ func (s *agent) ListDetached() []DetachedInfo { return nil } -func (s *agent) SignalDetached(pid, signal int) error { - if srv := s.execServer(); srv != nil { +func (a *agent) SignalDetached(pid, signal int) error { + if srv := a.execServer(); srv != nil { type signaler interface { SignalDetached(int, syscall.Signal) error } @@ -586,8 +586,8 @@ func (s *agent) SignalDetached(pid, signal int) error { return fmt.Errorf("no execute server available") } -func (s *agent) GetDetachedOutput(pid int) (string, error) { - if srv := s.execServer(); srv != nil { +func (a *agent) GetDetachedOutput(pid int) (string, error) { + if srv := a.execServer(); srv != nil { type outputGetter interface { GetDetachedOutput(int) (string, error) } @@ -598,8 +598,8 @@ func (s *agent) GetDetachedOutput(pid int) (string, error) { return "", fmt.Errorf("no execute server available") } -func (s *agent) DismissDetached(pid int) bool { - if srv := s.execServer(); srv != nil { +func (a *agent) DismissDetached(pid int) bool { + if srv := a.execServer(); srv != nil { type dismisser interface { DismissDetached(int) bool } @@ -610,8 +610,8 @@ func (s *agent) DismissDetached(pid int) bool { return false } -func (s *agent) InjectSystemEvent(content string) { - s.Queue("\n" + content + "\n") +func (a *agent) InjectSystemEvent(content string) { + a.Queue("\n" + content + "\n") } // classifyReaction returns a category and description for a reaction emoji. @@ -632,23 +632,23 @@ func classifyReaction(emoji string) (category, description string, positive bool } } -func (s *agent) Reactions() map[string]string { +func (a *agent) Reactions() map[string]string { result := make(map[string]string) - if s.history == nil { + if a.history == nil { return result } - for _, reaction := range s.history.Reactions { + for _, reaction := range a.history.Reactions { result[reaction.ResponseID] = reaction.Emoji } return result } -func (s *agent) React(emoji string) { - _ = s.ReactTo("", emoji) +func (a *agent) React(emoji string) { + _ = a.ReactTo("", emoji) } -func (s *agent) ReactTo(responseID, emoji string) error { - if s.history == nil { +func (a *agent) ReactTo(responseID, emoji string) error { + if a.history == nil { return fmt.Errorf("no active session") } category, _, _ := classifyReaction(emoji) @@ -656,9 +656,9 @@ func (s *agent) ReactTo(responseID, emoji string) error { return fmt.Errorf("unsupported reaction: %s", emoji) } if responseID == "" { - for i := len(s.history.messages) - 1; i >= 0; i-- { - if s.history.messages[i].Role == "assistant" { - responseID = s.history.messages[i].ID + for i := len(a.history.messages) - 1; i >= 0; i-- { + if a.history.messages[i].Role == "assistant" { + responseID = a.history.messages[i].ID break } } @@ -667,8 +667,8 @@ func (s *agent) ReactTo(responseID, emoji string) error { return fmt.Errorf("no assistant response to react to") } found := false - for i := range s.history.messages { - if s.history.messages[i].Role == "assistant" && s.history.messages[i].ID == responseID { + for i := range a.history.messages { + if a.history.messages[i].Role == "assistant" && a.history.messages[i].ID == responseID { found = true break } @@ -679,71 +679,71 @@ func (s *agent) ReactTo(responseID, emoji string) error { reaction := Reaction{ID: NewReactionID(), ResponseID: responseID, Emoji: emoji, Category: category, CreatedAt: time.Now()} replaced := false - for i := range s.history.Reactions { - if s.history.Reactions[i].ResponseID == responseID { - if s.history.Reactions[i].Emoji == emoji { + for i := range a.history.Reactions { + if a.history.Reactions[i].ResponseID == responseID { + if a.history.Reactions[i].Emoji == emoji { return nil } - s.history.Reactions[i] = reaction + a.history.Reactions[i] = reaction replaced = true break } } if !replaced { - s.history.Reactions = append(s.history.Reactions, reaction) + a.history.Reactions = append(a.history.Reactions, reaction) } - s.history.recomputeReactionCounts() - s.saveSession() + a.history.recomputeReactionCounts() + a.saveSession() return nil } -func (s *agent) AgentName() string { - v := s.agentName - s.log.Debug("AgentName() = %q", v) +func (a *agent) AgentName() string { + v := a.agentName + a.log.Debug("AgentName() = %q", v) return v } -func (s *agent) BackendName() string { - v := s.runtime.Backend.Name() - s.log.Debug("BackendName() = %q", v) +func (a *agent) BackendName() string { + v := a.runtime.Backend.Name() + a.log.Debug("BackendName() = %q", v) return v } -func (s *agent) ModelName() string { - v := s.runtime.Backend.Model() - s.log.Debug("ModelName() = %q", v) +func (a *agent) ModelName() string { + v := a.runtime.Backend.Model() + a.log.Debug("ModelName() = %q", v) return v } -func (s *agent) State() string { - return s.sess.State() +func (a *agent) State() string { + return a.sess.State() } -func (s *agent) notifyChange() { - s.changeMu.Lock() - s.changeCond.Broadcast() - s.changeMu.Unlock() +func (a *agent) notifyChange() { + a.changeMu.Lock() + a.changeCond.Broadcast() + a.changeMu.Unlock() } -func (s *agent) setState(state string) { - s.sess.SetState(state) - s.log.Debug("state -> %q", state) - s.notifyChange() +func (a *agent) setState(state string) { + a.sess.SetState(state) + a.log.Debug("state -> %q", state) + a.notifyChange() } // WaitChange blocks until the named field changes from current, then returns // the new value. Returns ("", false) if ctx is cancelled. -func (s *agent) WaitChange(ctx context.Context, field, current string) (string, bool) { +func (a *agent) WaitChange(ctx context.Context, field, current string) (string, bool) { read := func() string { switch field { case WatchState: - return s.State() + return a.State() case WatchUsage: - return s.Usage() + return a.Usage() case WatchCtxSz: - return s.CtxSz() + return a.CtxSz() case WatchCWD: - return s.CWD() + return a.CWD() case WatchAgent: - return s.AgentName() + return a.AgentName() } return "" } @@ -751,98 +751,98 @@ func (s *agent) WaitChange(ctx context.Context, field, current string) (string, // context.AfterFunc fires in a separate goroutine when ctx is done, // broadcasting to unblock any waiters. stop := context.AfterFunc(ctx, func() { - s.changeMu.Lock() - s.changeCond.Broadcast() - s.changeMu.Unlock() + a.changeMu.Lock() + a.changeCond.Broadcast() + a.changeMu.Unlock() }) defer stop() - s.changeMu.Lock() - defer s.changeMu.Unlock() + a.changeMu.Lock() + defer a.changeMu.Unlock() for ctx.Err() == nil { if v := read(); v != current { return v, true } - s.changeCond.Wait() + a.changeCond.Wait() } return "", false } -func (s *agent) Reply() string { - r := s.sess.Reply() - s.log.Debug("Reply() len=%d", len(r)) +func (a *agent) Reply() string { + r := a.sess.Reply() + a.log.Debug("Reply() len=%d", len(r)) return r } // CWD returns the current working directory for tool execution. -func (s *agent) CWD() string { - if s.sess.CWD() != "" { - s.log.Debug("CWD() = %q", s.sess.CWD()) - return s.sess.CWD() +func (a *agent) CWD() string { + if a.sess.CWD() != "" { + a.log.Debug("CWD() = %q", a.sess.CWD()) + return a.sess.CWD() } wd, _ := os.Getwd() - s.log.Debug("CWD() = %q (from getwd)", wd) + a.log.Debug("CWD() = %q (from getwd)", wd) return wd } // SetCWD changes the working directory for tool execution and updates the // system prompt. Returns an error if the path does not exist. -func (s *agent) SetCWD(dir string) error { - s.log.Debug("SetCWD(%q)", dir) +func (a *agent) SetCWD(dir string) error { + a.log.Debug("SetCWD(%q)", dir) dir = paths.ExpandHome(dir) if dir != "" { if _, err := os.Stat(dir); err != nil { return fmt.Errorf("cwd: %w", err) } } - oldCwd := s.sess.CWD() - s.sess.SetCWD(dir) + oldCwd := a.sess.CWD() + a.sess.SetCWD(dir) // Update cwd references in the system prompt. if oldCwd != "" && dir != "" && oldCwd != dir { - s.runtime.Preamble = strings.ReplaceAll(s.runtime.Preamble, oldCwd, dir) + a.runtime.Preamble = strings.ReplaceAll(a.runtime.Preamble, oldCwd, dir) } // Propagate to any tool server that knows how to handle it (e.g. execute). - if s.runtime != nil && s.runtime.Dispatcher != nil { - if srv, ok := s.runtime.Dispatcher.GetServer("execute"); ok { + if a.runtime != nil && a.runtime.Dispatcher != nil { + if srv, ok := a.runtime.Dispatcher.GetServer("execute"); ok { if ws, ok := srv.(tools.CWDSetter); ok { ws.SetCWD(dir) } } } - s.notifyChange() + a.notifyChange() return nil } // SetSessionID renames the session. It updates the in-memory ID, renames // persisted files on disk, and propagates to the execute server env. -func (s *agent) SetSessionID(newID string) error { - s.log.Debug("SetSessionID(%q) old=%q", newID, s.sess.ID()) - oldID := s.sess.ID() +func (a *agent) SetSessionID(newID string) error { + a.log.Debug("SetSessionID(%q) old=%q", newID, a.sess.ID()) + oldID := a.sess.ID() if oldID == newID { return nil } // Rename active persisted files on disk. - if s.sessionsDir != "" && oldID != "" { + if a.sessionsDir != "" && oldID != "" { for _, suffix := range []string{".json", ".compaction.jsonl"} { - oldPath := s.activeSessionPath(oldID, suffix) + oldPath := a.activeSessionPath(oldID, suffix) if _, err := os.Stat(oldPath); err == nil { - if err := os.Rename(oldPath, s.activeSessionPath(newID, suffix)); err != nil { + if err := os.Rename(oldPath, a.activeSessionPath(newID, suffix)); err != nil { return fmt.Errorf("rename %s: %w", suffix, err) } } } } - s.sess.SetID(newID) + a.sess.SetID(newID) // Update session ID references in the system prompt. - s.runtime.Preamble = strings.ReplaceAll(s.runtime.Preamble, oldID, newID) + a.runtime.Preamble = strings.ReplaceAll(a.runtime.Preamble, oldID, newID) // Rename tmpdir so isread markers remain valid after rename. oldTemp := filepath.Join(ollieTmpDir(), oldID) newTemp := filepath.Join(ollieTmpDir(), newID) if _, err := os.Stat(oldTemp); err == nil { os.Rename(oldTemp, newTemp) //nolint:errcheck } - s.pushSessionEnv() + a.pushSessionEnv() return nil } @@ -853,8 +853,8 @@ const defaultContextLength = 128000 const defaultToolResultMaxBytes = 131072 // autoCompactLimit returns the token threshold for auto-compaction (75%). -func (s *agent) autoCompactLimit(ctx context.Context) int { - ctxLen := s.runtime.Backend.ContextLength(ctx) +func (a *agent) autoCompactLimit(ctx context.Context) int { + ctxLen := a.runtime.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength } @@ -862,8 +862,8 @@ func (s *agent) autoCompactLimit(ctx context.Context) int { } // autoWarnLimit returns the token threshold for a context-usage warning (60%). -func (s *agent) autoWarnLimit(ctx context.Context) int { - ctxLen := s.runtime.Backend.ContextLength(ctx) +func (a *agent) autoWarnLimit(ctx context.Context) int { + ctxLen := a.runtime.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength } @@ -873,18 +873,18 @@ func (s *agent) autoWarnLimit(ctx context.Context) int { // spawnContext assembles the agent context injected at each session refresh // point (session start, post-clear, post-compaction). It combines the // agent-specific prompt with any agentSpawn hook output. -func (s *agent) spawnContext(ctx context.Context) string { - result := s.runtime.Hooks.Run(ctx, HookAgentSpawn, map[string]string{ - "session_id": s.sess.ID(), - "agent": s.agentName, - "cwd": s.CWD(), - "model": s.runtime.Backend.Model(), - }, s.log) +func (a *agent) spawnContext(ctx context.Context) string { + result := a.runtime.Hooks.Run(ctx, HookAgentSpawn, map[string]string{ + "session_id": a.sess.ID(), + "agent": a.agentName, + "cwd": a.CWD(), + "model": a.runtime.Backend.Model(), + }, a.log) if result.Warning != "" { - s.emit(infoEvent(result.Warning)) + a.emit(infoEvent(result.Warning)) } if sum := result.Summary(); sum != "" { - s.emit(infoEvent("agentSpawn: " + sum)) + a.emit(infoEvent("agentSpawn: " + sum)) } var parts []string if result.Context != "" { @@ -896,107 +896,107 @@ func (s *agent) spawnContext(ctx context.Context) string { // runCompact executes a full compaction cycle: pre-hook, compact, spawn-context // re-injection, post-hook. Returns (n compacted, error). Returns (0, nil) if // the pre-hook blocked or there was nothing to compact. Caller manages setState. -func (s *agent) runCompact(ctx context.Context, trigger string) (int, error) { - payload := map[string]string{"session_id": s.sess.ID(), "trigger": trigger, "cwd": s.CWD()} - pre := s.runtime.Hooks.Run(ctx, HookPreCompact, payload, s.log) +func (a *agent) runCompact(ctx context.Context, trigger string) (int, error) { + payload := map[string]string{"session_id": a.sess.ID(), "trigger": trigger, "cwd": a.CWD()} + pre := a.runtime.Hooks.Run(ctx, HookPreCompact, payload, a.log) if pre.Warning != "" { - s.emit(infoEvent(pre.Warning)) + a.emit(infoEvent(pre.Warning)) } if sum := pre.Summary(); sum != "" { - s.emit(infoEvent("preCompact: " + sum)) + a.emit(infoEvent("preCompact: " + sum)) } if pre.Blocked { - s.emit(infoEvent("compact cancelled by hook")) + a.emit(infoEvent("compact cancelled by hook")) return 0, nil } if pre.Context != "" { - s.history.appendUserMessage(pre.Context) + a.history.appendUserMessage(pre.Context) } // Use a cheaper model for compaction if configured. - compactModel := resolveCompactionModel(s.runtime.CompactionModel, s.runtime.Backend) - origModel := s.runtime.Backend.Model() + compactModel := resolveCompactionModel(a.runtime.CompactionModel, a.runtime.Backend) + origModel := a.runtime.Backend.Model() if compactModel != "" && compactModel != origModel { - s.runtime.Backend.SetModel(compactModel) - defer s.runtime.Backend.SetModel(origModel) + a.runtime.Backend.SetModel(compactModel) + defer a.runtime.Backend.SetModel(origModel) } - n, _, err := s.history.compact(ctx, s.runtime.Backend) + n, _, err := a.history.compact(ctx, a.runtime.Backend) if err != nil { return 0, err } if n > 0 { - s.auditLog.Debug("compact: removed %d messages trigger=%s session=%s", n, trigger, s.sess.ID()) - s.warnedContext = false - if sc := s.spawnContext(ctx); sc != "" { - s.history.appendUserMessage(sc) + a.auditLog.Debug("compact: removed %d messages trigger=%s session=%s", n, trigger, a.sess.ID()) + a.warnedContext = false + if sc := a.spawnContext(ctx); sc != "" { + a.history.appendUserMessage(sc) } } - post := s.runtime.Hooks.Run(ctx, HookPostCompact, payload, s.log) + post := a.runtime.Hooks.Run(ctx, HookPostCompact, payload, a.log) if post.Warning != "" { - s.emit(infoEvent(post.Warning)) + a.emit(infoEvent(post.Warning)) } if sum := post.Summary(); sum != "" { - s.emit(infoEvent("postCompact: " + sum)) + a.emit(infoEvent("postCompact: " + sum)) } if post.Context != "" { - s.history.appendUserMessage(post.Context) + a.history.appendUserMessage(post.Context) } return n, nil } -func (s *agent) activeSessionPath(id, suffix string) string { - return filepath.Join(s.sessionsDir, "active", id+suffix) +func (a *agent) activeSessionPath(id, suffix string) string { + return filepath.Join(a.sessionsDir, "active", id+suffix) } -func (s *agent) saveSession() { - s.saveMu.Lock() - s.saveDirty = true - if s.saveTimer == nil { - s.saveTimer = time.AfterFunc(2*time.Second, s.flushSave) +func (a *agent) saveSession() { + a.saveMu.Lock() + a.saveDirty = true + if a.saveTimer == nil { + a.saveTimer = time.AfterFunc(2*time.Second, a.flushSave) } - s.saveMu.Unlock() + a.saveMu.Unlock() } // flushSave immediately persists the session if dirty. -func (s *agent) flushSave() { - s.saveMu.Lock() - dirty := s.saveDirty - s.saveDirty = false - if s.saveTimer != nil { - s.saveTimer.Stop() - s.saveTimer = nil +func (a *agent) flushSave() { + a.saveMu.Lock() + dirty := a.saveDirty + a.saveDirty = false + if a.saveTimer != nil { + a.saveTimer.Stop() + a.saveTimer = nil } - s.saveMu.Unlock() + a.saveMu.Unlock() if !dirty { return } - if s.history == nil || s.sess.ID() == "" || s.sessionsDir == "" { + if a.history == nil || a.sess.ID() == "" || a.sessionsDir == "" { return } - path := s.activeSessionPath(s.sess.ID(), ".json") + path := a.activeSessionPath(a.sess.ID(), ".json") if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { - s.log.Error("session save: %v", err) + a.log.Error("session save: %v", err) return } - if err := s.history.saveToFull(path, s.sess.ID(), s.agentName, - s.runtime.Backend.Name(), s.runtime.Backend.Model(), s.CWD(), s.remote); err != nil { - s.log.Error("session save: %v", err) + if err := a.history.saveToFull(path, a.sess.ID(), a.agentName, + a.runtime.Backend.Name(), a.runtime.Backend.Model(), a.CWD(), a.remote); err != nil { + a.log.Error("session save: %v", err) } } // SaveSession writes the current session state to the given path, including // backend and model metadata for external restore. -func (s *agent) SaveSession(path string) error { - s.mu.RLock() - defer s.mu.RUnlock() - if s.history == nil { +func (a *agent) SaveSession(path string) error { + a.mu.RLock() + defer a.mu.RUnlock() + if a.history == nil { return fmt.Errorf("no active session") } - return s.history.saveToFull(path, s.sess.ID(), s.agentName, - s.runtime.Backend.Name(), s.runtime.Backend.Model(), s.CWD(), s.remote) + return a.history.saveToFull(path, a.sess.ID(), a.agentName, + a.runtime.Backend.Name(), a.runtime.Backend.Model(), a.CWD(), a.remote) } -func (s *agent) getActionCancel() context.CancelCauseFunc { - if a := s.currentAction.Load(); a != nil { +func (a *agent) getActionCancel() context.CancelCauseFunc { + if a := a.currentAction.Load(); a != nil { return a.cancel } return nil @@ -1004,152 +1004,152 @@ func (s *agent) getActionCancel() context.CancelCauseFunc { // Interrupt cancels the current in-progress agent turn. // Returns true if an action was running and was cancelled. -func (s *agent) Interrupt(cause error) bool { - s.log.Debug("Interrupt() cause=%v", cause) - if cancel := s.getActionCancel(); cancel != nil { +func (a *agent) Interrupt(cause error) bool { + a.log.Debug("Interrupt() cause=%v", cause) + if cancel := a.getActionCancel(); cancel != nil { cancel(cause) return true } return false } -func (s *agent) Inject(prompt string) { +func (a *agent) Inject(prompt string) { // If an inject is already pending, fall back to the normal FIFO so nothing // is lost. Use CompareAndSwap to avoid a race between the nil check and store. - if !s.pendingInject.CompareAndSwap(nil, &prompt) { - s.sess.Queue(prompt) + if !a.pendingInject.CompareAndSwap(nil, &prompt) { + a.sess.Queue(prompt) return } - s.emit(Event{Role: "info", Content: "\n"}) - s.emit(Event{Role: "user", Content: prompt}) + a.emit(Event{Role: "info", Content: "\n"}) + a.emit(Event{Role: "user", Content: prompt}) } -func (s *agent) injectRewrite(prompt string) { - s.pendingInject.Store(&prompt) - s.emit(Event{Role: "info", Content: "\n"}) - s.emit(Event{Role: "user", Content: prompt}) +func (a *agent) injectRewrite(prompt string) { + a.pendingInject.Store(&prompt) + a.emit(Event{Role: "info", Content: "\n"}) + a.emit(Event{Role: "user", Content: prompt}) } -func (s *agent) Queue(prompt string) { - s.sess.Queue(prompt) - s.sess.Bus().Publish("queued", prompt) +func (a *agent) Queue(prompt string) { + a.sess.Queue(prompt) + a.sess.Bus().Publish("queued", prompt) } -func (s *agent) drainQueue() { - if prompt, ok := s.sess.PopQueue(); ok { - s.Submit(context.Background(), prompt) +func (a *agent) drainQueue() { + if prompt, ok := a.sess.PopQueue(); ok { + a.Submit(context.Background(), prompt) } } -func (s *agent) Bus() *pubsub.Bus { - return s.sess.Bus() +func (a *agent) Bus() *pubsub.Bus { + return a.sess.Bus() } -func (s *agent) emit(ev Event) { - s.sess.Bus().Publish("event", ev) +func (a *agent) emit(ev Event) { + a.sess.Bus().Publish("event", ev) } -func (s *agent) PopQueue() (string, bool) { - return s.sess.PopQueue() +func (a *agent) PopQueue() (string, bool) { + return a.sess.PopQueue() } -func (s *agent) IsRunning() bool { - v := s.currentAction.Load() != nil - s.log.Debug("IsRunning() = %v", v) +func (a *agent) IsRunning() bool { + v := a.currentAction.Load() != nil + a.log.Debug("IsRunning() = %v", v) return v } -func (s *agent) CtxSz() string { - if s.history == nil { - s.log.Debug("CtxSz() no session") +func (a *agent) CtxSz() string { + if a.history == nil { + a.log.Debug("CtxSz() no session") return "no active session" } - ctxLen := s.runtime.Backend.ContextLength(context.Background()) + ctxLen := a.runtime.Backend.ContextLength(context.Background()) if ctxLen <= 0 { ctxLen = defaultContextLength } - estimated := s.history.estimateTokens() + estimated := a.history.estimateTokens() pct := estimated * 100 / ctxLen v := fmt.Sprintf("%d / %d (%d%%)", estimated, ctxLen, pct) - s.log.Debug("CtxSz() = %q", v) + a.log.Debug("CtxSz() = %q", v) return v } -func (s *agent) Cost() string { - if s.history == nil { +func (a *agent) Cost() string { + if a.history == nil { return "no active session" } return fmt.Sprintf("costLast=$%.4f\ncostSession=$%.4f\n", - s.history.LastTurnCostUSD, s.history.SessionCostUSD) + a.history.LastTurnCostUSD, a.history.SessionCostUSD) } -func (s *agent) Usage() string { - if s.history == nil { - s.log.Debug("Usage() no session") +func (a *agent) Usage() string { + if a.history == nil { + a.log.Debug("Usage() no session") return "no active session" } str := fmt.Sprintf("%d in, %d out, %d requests", - s.history.TotalInputTokens, s.history.TotalOutputTokens, - s.history.TotalRequests) - if s.history.TotalCachedInputTokens > 0 { - str += fmt.Sprintf(", %d cached", s.history.TotalCachedInputTokens) + a.history.TotalInputTokens, a.history.TotalOutputTokens, + a.history.TotalRequests) + if a.history.TotalCachedInputTokens > 0 { + str += fmt.Sprintf(", %d cached", a.history.TotalCachedInputTokens) } - if s.history.Estimated { + if a.history.Estimated { str += " [estimated]" } - s.log.Debug("Usage() = %q", str) + a.log.Debug("Usage() = %q", str) return str } -func (s *agent) Context() []backend.Message { - s.mu.RLock() +func (a *agent) Context() []backend.Message { + a.mu.RLock() var msgs []backend.Message - if s.history != nil { - msgs = slices.Clone(s.history.history()) + if a.history != nil { + msgs = slices.Clone(a.history.history()) } - s.mu.RUnlock() - if s.runtime.Preamble != "" { - msgs = append([]backend.Message{{Role: "system", Content: s.runtime.Preamble}}, msgs...) + a.mu.RUnlock() + if a.runtime.Preamble != "" { + msgs = append([]backend.Message{{Role: "system", Content: a.runtime.Preamble}}, msgs...) } return msgs } -func (s *agent) SystemPrompt() string { - s.log.Debug("SystemPrompt() len=%d", len(s.runtime.Preamble)) - return s.runtime.Preamble +func (a *agent) SystemPrompt() string { + a.log.Debug("SystemPrompt() len=%d", len(a.runtime.Preamble)) + return a.runtime.Preamble } -func (s *agent) GenerationParams() backend.GenerationParams { - s.mu.RLock() - defer s.mu.RUnlock() - return s.runtime.GenParams +func (a *agent) GenerationParams() backend.GenerationParams { + a.mu.RLock() + defer a.mu.RUnlock() + return a.runtime.GenParams } -func (s *agent) CompactionModel() string { - s.mu.RLock() - defer s.mu.RUnlock() - return s.runtime.CompactionModel +func (a *agent) CompactionModel() string { + a.mu.RLock() + defer a.mu.RUnlock() + return a.runtime.CompactionModel } -func (s *agent) SetCompactionModel(model string) { - s.mu.Lock() - defer s.mu.Unlock() - s.runtime.CompactionModel = model +func (a *agent) SetCompactionModel(model string) { + a.mu.Lock() + defer a.mu.Unlock() + a.runtime.CompactionModel = model } -func (s *agent) SetGenerationParams(params backend.GenerationParams) error { - if s.IsRunning() { +func (a *agent) SetGenerationParams(params backend.GenerationParams) error { + if a.IsRunning() { return fmt.Errorf("cannot change params while agent is running") } - s.mu.Lock() - s.runtime.GenParams = params - s.mu.Unlock() + a.mu.Lock() + a.runtime.GenParams = params + a.mu.Unlock() return nil } -func (s *agent) ListModels() string { - s.log.Debug("ListModels()") - models := s.runtime.Backend.Models(context.Background()) +func (a *agent) ListModels() string { + a.log.Debug("ListModels()") + models := a.runtime.Backend.Models(context.Background()) slices.Sort(models) return strings.Join(models, "\n") } @@ -1176,18 +1176,18 @@ func firstSentence(s string) string { // // Continuations (post-turn hook context, unconsumed inject, FIFO drain) are // handled via an explicit loop rather than recursion to avoid stack growth. -func (s *agent) Submit(ctx context.Context, input string) { +func (a *agent) Submit(ctx context.Context, input string) { defer func() { if r := recover(); r != nil { - s.log.Error("panic: %v\n%s", r, debug.Stack()) - if a := s.currentAction.Swap(nil); a != nil { + a.log.Error("panic: %v\n%s", r, debug.Stack()) + if a := a.currentAction.Swap(nil); a != nil { a.cancel(fmt.Errorf("%v", r)) } - s.setState("idle") - s.emit(Event{Role: "error", Content: fmt.Sprintf("%v", r)}) + a.setState("idle") + a.emit(Event{Role: "error", Content: fmt.Sprintf("%v", r)}) } }() - s.log.Debug("Submit() input_len=%d running=%v", len(input), s.IsRunning()) + a.log.Debug("Submit() input_len=%d running=%v", len(input), a.IsRunning()) if input == "" { return } @@ -1195,51 +1195,51 @@ func (s *agent) Submit(ctx context.Context, input string) { // Fast path: inject and FIFO push use atomics and are safe without // the submit lock. Handle them before acquiring submitMu so they // don't block behind a long-running turn or command. - if s.IsRunning() { - if s.handleCommand(ctx, input) { + if a.IsRunning() { + if a.handleCommand(ctx, input) { return } - s.sess.Queue(input) + a.sess.Queue(input) return } // Serialize commands and turns so that e.g. a /compact arriving via // ctl cannot race with an executeTurn arriving via prompt. - s.submitMu.Lock() - defer s.submitMu.Unlock() + a.submitMu.Lock() + defer a.submitMu.Unlock() - if s.handleCommand(ctx, input) { + if a.handleCommand(ctx, input) { return } - if s.IsRunning() { - s.sess.Queue(input) + if a.IsRunning() { + a.sess.Queue(input) return } for input != "" && ctx.Err() == nil { - input = s.executeTurn(ctx, input) + input = a.executeTurn(ctx, input) } } // executeTurn runs a single agent turn and returns the next prompt to execute, // or "" if there is nothing more to do. -func (s *agent) executeTurn(ctx context.Context, input string) string { - s.emit(Event{Role: "user", Content: input}) +func (a *agent) executeTurn(ctx context.Context, input string) string { + a.emit(Event{Role: "user", Content: input}) - hookResult := s.runtime.Hooks.Run(ctx, HookPreTurn, map[string]string{ - "session_id": s.sess.ID(), - "cwd": s.CWD(), + hookResult := a.runtime.Hooks.Run(ctx, HookPreTurn, map[string]string{ + "session_id": a.sess.ID(), + "cwd": a.CWD(), "prompt": input, - }, s.log) + }, a.log) if hookResult.Blocked { - s.emit(infoEvent("hook blocked prompt")) + a.emit(infoEvent("hook blocked prompt")) return "" } if hookResult.Warning != "" { - s.emit(infoEvent(hookResult.Warning)) + a.emit(infoEvent(hookResult.Warning)) } if sum := hookResult.Summary(); sum != "" { - s.emit(infoEvent("preTurn: " + sum)) + a.emit(infoEvent("preTurn: " + sum)) } if hookResult.Context != "" { input += "\n" + hookResult.Context @@ -1247,171 +1247,171 @@ func (s *agent) executeTurn(ctx context.Context, input string) string { // Snapshot session state before this turn modifies it. Restored on failure // so the session is clean for the next attempt. - snapSession := s.history + snapSession := a.history var snapMessages []backend.Message - if s.history != nil { - snapMessages = cloneMessages(s.history.messages) + if a.history != nil { + snapMessages = cloneMessages(a.history.messages) } - if s.history == nil { - for _, msg := range s.startupMessages { - s.log.Debug("startup: %s", msg) - s.emit(infoEvent(msg)) + if a.history == nil { + for _, msg := range a.startupMessages { + a.log.Debug("startup: %s", msg) + a.emit(infoEvent(msg)) } - s.startupMessages = nil - s.history = newHistory(input) - if sc := s.spawnContext(ctx); sc != "" { - s.history.appendUserMessage(sc) + a.startupMessages = nil + a.history = newHistory(input) + if sc := a.spawnContext(ctx); sc != "" { + a.history.appendUserMessage(sc) } - s.history.appendUserMessage(input) + a.history.appendUserMessage(input) } else { - s.history.appendUserMessage(input) + a.history.appendUserMessage(input) } actCtx, actCancel := context.WithCancelCause(ctx) handle := &actionHandle{cancel: actCancel} - s.currentAction.Store(handle) - s.setState("thinking") + a.currentAction.Store(handle) + a.setState("thinking") - s.auditLog.Debug("turn: start input=%s session=%s", auditTruncate(input), s.sess.ID()) + a.auditLog.Debug("turn: start input=%s session=%s", auditTruncate(input), a.sess.ID()) // Build per-turn agentConfig from the current runtime. - s.cfg = agentConfig{ - Backend: s.runtime.Backend, - preamble: s.runtime.Preamble, - Tools: s.runtime.Tools, - Exec: s.runtime.Exec, - ClassifyTool: s.runtime.ClassifyTool, - ClassifyTier: s.runtime.ClassifyTier, - GenerationParams: s.runtime.GenParams, - MaxSteps: s.runtime.MaxSteps, - ReadPlanStep: s.readPlanStep, - TurnError: s.turnError, + a.cfg = agentConfig{ + Backend: a.runtime.Backend, + preamble: a.runtime.Preamble, + Tools: a.runtime.Tools, + Exec: a.runtime.Exec, + ClassifyTool: a.runtime.ClassifyTool, + ClassifyTier: a.runtime.ClassifyTier, + GenerationParams: a.runtime.GenParams, + MaxSteps: a.runtime.MaxSteps, + ReadPlanStep: a.readPlanStep, + TurnError: a.turnError, } var replyBuf strings.Builder - s.cfg.Output = func(ev Event) { + a.cfg.Output = func(ev Event) { switch ev.Role { case "assistant": replyBuf.WriteString(ev.Content) case "call": - s.setState("calling: " + ev.Name) - s.auditLog.Debug("call: %s %s", ev.Name, auditTruncate(string(ev.Content))) + a.setState("calling: " + ev.Name) + a.auditLog.Debug("call: %s %s", ev.Name, auditTruncate(string(ev.Content))) case "tool": - s.auditLog.Debug("result: %s %s", ev.Name, auditTruncate(ev.Content)) + a.auditLog.Debug("result: %s %s", ev.Name, auditTruncate(ev.Content)) case "state": - s.setState(ev.Content) + a.setState(ev.Content) case "limitretry": - s.setState("limitretry") + a.setState("limitretry") case "error": - s.auditLog.Debug("error: %s", ev.Content) + a.auditLog.Debug("error: %s", ev.Content) } - if ev.Role == "usage" && s.history != nil { + if ev.Role == "usage" && a.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) - s.history.addUsage(backend.Usage{ + a.history.addUsage(backend.Usage{ InputTokens: in, CachedInputTokens: cached, CacheCreationTokens: creation, OutputTokens: out, CostUSD: costUSD, }, est != 0) - s.notifyChange() + a.notifyChange() if maxCostStr := os.Getenv("OLLIE_MAX_SESSION_COST"); maxCostStr != "" { - if limit, ferr := strconv.ParseFloat(maxCostStr, 64); ferr == nil && limit > 0 && s.history.SessionCostUSD >= limit { - s.emit(infoEvent(fmt.Sprintf("spending cap $%.2f reached — stopping", limit))) - s.Interrupt(ErrInterrupted) + if limit, ferr := strconv.ParseFloat(maxCostStr, 64); ferr == nil && limit > 0 && a.history.SessionCostUSD >= limit { + a.emit(infoEvent(fmt.Sprintf("spending cap $%.2f reached — stopping", limit))) + a.Interrupt(ErrInterrupted) } } } - s.emit(ev) + a.emit(ev) } - s.cfg.PopInject = func() string { - if p := s.pendingInject.Swap(nil); p != nil { + a.cfg.PopInject = func() string { + if p := a.pendingInject.Swap(nil); p != nil { return *p } return "" } - s.cfg.PreTool = func(ctx context.Context, name string, args json.RawMessage) HookResult { - return s.runtime.Hooks.Run(ctx, HookPreTool, map[string]string{ - "session_id": s.sess.ID(), - "cwd": s.CWD(), + a.cfg.PreTool = func(ctx context.Context, name string, args json.RawMessage) HookResult { + return a.runtime.Hooks.Run(ctx, HookPreTool, map[string]string{ + "session_id": a.sess.ID(), + "cwd": a.CWD(), "tool": name, "args": string(args), - }, s.log) + }, a.log) } - s.cfg.PostTool = func(ctx context.Context, name string, args json.RawMessage, result string) HookResult { - return s.runtime.Hooks.Run(ctx, HookPostTool, map[string]string{ - "session_id": s.sess.ID(), - "cwd": s.CWD(), + a.cfg.PostTool = func(ctx context.Context, name string, args json.RawMessage, result string) HookResult { + return a.runtime.Hooks.Run(ctx, HookPostTool, map[string]string{ + "session_id": a.sess.ID(), + "cwd": a.CWD(), "tool": name, "args": string(args), "result": result, - }, s.log) + }, a.log) } - s.cfg.IncrToolCallCount = func() int64 { - return s.toolCallCount.Add(1) + a.cfg.IncrToolCallCount = func() int64 { + return a.toolCallCount.Add(1) } - s.cfg.SaveSession = func() { s.saveSession() } - s.cfg.ResultCache = &s.resultCache - s.cfg.AutoCompact = func(ctx context.Context) { - if ctx.Err() != nil || s.history == nil { + a.cfg.SaveSession = func() { a.saveSession() } + a.cfg.ResultCache = &a.resultCache + a.cfg.AutoCompact = func(ctx context.Context) { + if ctx.Err() != nil || a.history == nil { return } - limit := s.autoCompactLimit(ctx) - if limit <= 0 || s.history.estimateTokens() < limit { + limit := a.autoCompactLimit(ctx) + if limit <= 0 || a.history.estimateTokens() < limit { return } - s.emit(Event{Role: "info", Content: "auto-compacting context...\n"}) - s.setState("compacting") - if _, err := s.runCompact(ctx, "auto"); err != nil { + a.emit(Event{Role: "info", Content: "auto-compacting context...\n"}) + a.setState("compacting") + if _, err := a.runCompact(ctx, "auto"); err != nil { panic(fmt.Sprintf("mid-turn auto-compact: %v", err)) } - s.setState("thinking") + a.setState("thinking") } // Warn once when context usage crosses 60%; compact at 75%. - if s.history != nil { - tokens := s.history.estimateTokens() - if compactLimit := s.autoCompactLimit(ctx); compactLimit > 0 && tokens >= compactLimit { - s.emit(Event{Role: "info", Content: "auto-compacting context...\n"}) - s.setState("compacting") - if _, err := s.runCompact(ctx, "auto"); err != nil { + if a.history != nil { + tokens := a.history.estimateTokens() + if compactLimit := a.autoCompactLimit(ctx); compactLimit > 0 && tokens >= compactLimit { + a.emit(Event{Role: "info", Content: "auto-compacting context...\n"}) + a.setState("compacting") + if _, err := a.runCompact(ctx, "auto"); err != nil { panic(fmt.Sprintf("auto-compact: %v", err)) } - s.setState("thinking") - } else if warnLimit := s.autoWarnLimit(ctx); warnLimit > 0 && tokens >= warnLimit && !s.warnedContext { - ctxLen := s.cfg.Backend.ContextLength(ctx) + a.setState("thinking") + } else if warnLimit := a.autoWarnLimit(ctx); warnLimit > 0 && tokens >= warnLimit && !a.warnedContext { + ctxLen := a.cfg.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength } pct := tokens * 100 / ctxLen - s.emit(Event{Role: "info", Content: fmt.Sprintf("context at %d%% — will auto-compact at 75%%\n", pct)}) - s.warnedContext = true + a.emit(Event{Role: "info", Content: fmt.Sprintf("context at %d%% — will auto-compact at 75%%\n", pct)}) + a.warnedContext = true } } // Spending cap: reject before spending more tokens. - if maxCostStr := os.Getenv("OLLIE_MAX_SESSION_COST"); maxCostStr != "" && s.history != nil { + if maxCostStr := os.Getenv("OLLIE_MAX_SESSION_COST"); maxCostStr != "" && a.history != nil { if limit, err := strconv.ParseFloat(maxCostStr, 64); err == nil && limit > 0 { - if s.history.SessionCostUSD >= limit { - s.emit(Event{Role: "error", Content: fmt.Sprintf("spending cap $%.2f reached (session total $%.4f)", limit, s.history.SessionCostUSD)}) - s.setState("idle") + if a.history.SessionCostUSD >= limit { + a.emit(Event{Role: "error", Content: fmt.Sprintf("spending cap $%.2f reached (session total $%.4f)", limit, a.history.SessionCostUSD)}) + a.setState("idle") actCancel(nil) - s.currentAction.CompareAndSwap(handle, nil) + a.currentAction.CompareAndSwap(handle, nil) if snapSession == nil { - s.history = nil + a.history = nil } else { - s.history.messages = snapMessages + a.history.messages = snapMessages } return "" } } } - if s.history != nil { - s.history.resetTurnAccumulators() + if a.history != nil { + a.history.resetTurnAccumulators() } // Run the turn, retrying once after compaction on context overflow. @@ -1420,9 +1420,9 @@ func (s *agent) executeTurn(ctx context.Context, input string) string { err error ) for { - err = run(actCtx, s.cfg, s.history) + err = run(actCtx, a.cfg, a.history) actCancel(nil) - s.currentAction.CompareAndSwap(handle, nil) + a.currentAction.CompareAndSwap(handle, nil) if err == nil { break @@ -1431,77 +1431,77 @@ func (s *agent) executeTurn(ctx context.Context, input string) string { break } var ctxErr *backend.ContextOverflowError - if !overflowRetried && errors.As(err, &ctxErr) && s.history != nil { + if !overflowRetried && errors.As(err, &ctxErr) && a.history != nil { overflowRetried = true - s.history.messages = snapMessages - s.emit(Event{Role: "info", Content: "context overflow — compacting and retrying...\n"}) - s.setState("compacting") - if _, cerr := s.runCompact(ctx, "overflow"); cerr != nil { + a.history.messages = snapMessages + a.emit(Event{Role: "info", Content: "context overflow — compacting and retrying...\n"}) + a.setState("compacting") + if _, cerr := a.runCompact(ctx, "overflow"); cerr != nil { break } - s.history.appendUserMessage(input) - s.setState("thinking") - s.history.resetTurnAccumulators() + a.history.appendUserMessage(input) + a.setState("thinking") + a.history.resetTurnAccumulators() replyBuf.Reset() actCtx, actCancel = context.WithCancelCause(ctx) handle = &actionHandle{cancel: actCancel} - s.currentAction.Store(handle) + a.currentAction.Store(handle) continue } break } - s.mu.Lock() - s.sess.SetReply(replyBuf.String()) - s.mu.Unlock() + a.mu.Lock() + a.sess.SetReply(replyBuf.String()) + a.mu.Unlock() replyBuf.Reset() - s.setState("idle") - s.flushSave() + a.setState("idle") + a.flushSave() if err != nil { // Keep completed work — only remove cancelled tool results. // Error results are valuable feedback for the agent. - if s.history != nil { - s.history.removeCancelledToolResults() + if a.history != nil { + a.history.removeCancelledToolResults() } if errors.Is(err, context.Canceled) || errors.Is(err, ErrInterrupted) { - s.auditLog.Debug("turn: interrupted session=%s", s.sess.ID()) - s.saveSession() + a.auditLog.Debug("turn: interrupted session=%s", a.sess.ID()) + a.saveSession() return "" } - s.emit(Event{Role: "error", Content: err.Error()}) + a.emit(Event{Role: "error", Content: err.Error()}) // Drain one FIFO item — the turnError hook may have queued a recovery prompt. - if next, ok := s.sess.PopQueue(); ok { + if next, ok := a.sess.PopQueue(); ok { return next } return "" } - stopResult := s.runtime.Hooks.Run(ctx, HookPostTurn, map[string]string{ - "session_id": s.sess.ID(), - "cwd": s.CWD(), - }, s.log) + stopResult := a.runtime.Hooks.Run(ctx, HookPostTurn, map[string]string{ + "session_id": a.sess.ID(), + "cwd": a.CWD(), + }, a.log) if stopResult.Warning != "" { - s.emit(infoEvent(stopResult.Warning)) + a.emit(infoEvent(stopResult.Warning)) } if sum := stopResult.Summary(); sum != "" { - s.emit(infoEvent("postTurn: " + sum)) + a.emit(infoEvent("postTurn: " + sum)) } - if !stopResult.Blocked && stopResult.Context != "" && s.history != nil { - s.history.appendUserMessage(stopResult.Context) + if !stopResult.Blocked && stopResult.Context != "" && a.history != nil { + a.history.appendUserMessage(stopResult.Context) } - if s.history != nil { - s.history.recordTurnCost(s.cfg.Backend.Model()) - appendUsageLog(s.sess.ID(), s.cfg.Backend.Name(), s.cfg.Backend.Model(), s.history) - if s.history.LastTurnCostUSD > 0 { - s.emit(Event{Role: "info", Content: fmt.Sprintf("costLast=$%.4f\n", s.history.LastTurnCostUSD)}) + if a.history != nil { + a.history.recordTurnCost(a.cfg.Backend.Model()) + appendUsageLog(a.sess.ID(), a.cfg.Backend.Name(), a.cfg.Backend.Model(), a.history) + if a.history.LastTurnCostUSD > 0 { + a.emit(Event{Role: "info", Content: fmt.Sprintf("costLast=$%.4f\n", a.history.LastTurnCostUSD)}) } - s.auditLog.Debug("turn: end reply=%s cost=$%.4f session_total=$%.4f session=%s", - auditTruncate(s.sess.Reply()), s.history.LastTurnCostUSD, s.history.SessionCostUSD, s.sess.ID()) - s.notifyChange() + a.auditLog.Debug("turn: end reply=%s cost=$%.4f session_total=$%.4f session=%s", + auditTruncate(a.sess.Reply()), a.history.LastTurnCostUSD, a.history.SessionCostUSD, a.sess.ID()) + a.notifyChange() } - s.saveSession() + a.saveSession() // Post-turn hook said "continue" — its context becomes the next prompt. if stopResult.Blocked && stopResult.Context != "" { @@ -1510,12 +1510,12 @@ func (s *agent) executeTurn(ctx context.Context, input string) string { // Inject that was pending but never consumed (text-only response with no // tool calls) — treat it as the next user message. - if p := s.pendingInject.Swap(nil); p != nil { + if p := a.pendingInject.Swap(nil); p != nil { return *p } // Drain one item from the FIFO; the outer loop handles the rest. - if next, ok := s.sess.PopQueue(); ok { + if next, ok := a.sess.PopQueue(); ok { return next }