diff --git a/agent/commands.go b/agent/commands.go index 3f59448..e7322bb 100644 --- a/agent/commands.go +++ b/agent/commands.go @@ -12,7 +12,7 @@ import ( ) -func (a *sessionHost) handleCommand(ctx context.Context, input string) bool { +func (a *harness) handleCommand(ctx context.Context, input string) bool { if !strings.HasPrefix(input, "/") { return false } diff --git a/agent/core_test.go b/agent/core_test.go index 26fb052..b4a9d28 100644 --- a/agent/core_test.go +++ b/agent/core_test.go @@ -83,7 +83,7 @@ func defaultBE() *mockBackend { // newCore builds a minimal *agent for tests, bypassing loadSystemPrompt // by directly setting Preamble on Runtime. -func newCore(t *testing.T, be backend.Backend, hooks Hooks) *sessionHost { +func newCore(t *testing.T, be backend.Backend, hooks Hooks) *harness { t.Helper() t.Setenv("OLLIE", "") if be == nil { @@ -107,7 +107,7 @@ func newCore(t *testing.T, be backend.Backend, hooks Hooks) *sessionHost { NewDispatcher: tools.NewDispatcher, }) t.Cleanup(c.Close) - return c.(*sessionHost) + return c.(*harness) } // collectEvents runs Submit synchronously and returns all emitted events. @@ -1788,7 +1788,7 @@ func TestSetSessionID_UpdatesPreamble(t *testing.T) { CWD: t.TempDir(), Runtime: env, NewDispatcher: tools.NewDispatcher, - }).(*sessionHost) + }).(*harness) t.Cleanup(c.Close) newID := NewSessionID() @@ -2270,7 +2270,7 @@ func (m *mockEnvServer) get(k string) string { return m.env[k] } -func newCoreWithExecServer(t *testing.T, srv *mockEnvServer) *sessionHost { +func newCoreWithExecServer(t *testing.T, srv *mockEnvServer) *harness { t.Helper() t.Setenv("OLLIE", "") d := tools.NewDispatcher() @@ -2291,7 +2291,7 @@ func newCoreWithExecServer(t *testing.T, srv *mockEnvServer) *sessionHost { NewDispatcher: tools.NewDispatcher, }) t.Cleanup(c.Close) - return c.(*sessionHost) + return c.(*harness) } func TestSetEnv_PropagatestoExecuteServer(t *testing.T) { diff --git a/agent/host.go b/agent/harness.go similarity index 93% rename from agent/host.go rename to agent/harness.go index 20c76c4..8a64658 100644 --- a/agent/host.go +++ b/agent/harness.go @@ -307,9 +307,9 @@ type AgentCoreConfig struct { BaseLayers []string } -// sessionHost is the Core implementation. It owns all sessionHost and session state +// harness is the Core implementation. It owns all harness and session state // but has no knowledge of how output is rendered. -type sessionHost struct { +type harness struct { sess *session.Session r *Agent log *olog.Logger @@ -337,12 +337,12 @@ type sessionHost 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 (a *sessionHost) ToolCallCount() int64 { +func (a *harness) ToolCallCount() int64 { return a.toolCallCount.Load() } // SetEnv stores a session-scoped variable and propagates it to the execute server. -func (a *sessionHost) SetEnv(key, value string) { +func (a *harness) SetEnv(key, value string) { a.sess.SetEnv(key, value) if a.r.runtime == nil || a.r.runtime.Dispatcher == nil { return @@ -355,7 +355,7 @@ func (a *sessionHost) SetEnv(key, value string) { } // pushSessionEnv injects OLLIE_SESSION_ID into the execute server subprocess env. -func (a *sessionHost) pushSessionEnv() { +func (a *harness) pushSessionEnv() { if a.r.runtime == nil || a.r.runtime.Dispatcher == nil || a.sess.ID() == "" { return } @@ -370,14 +370,14 @@ func (a *sessionHost) pushSessionEnv() { } // pushLockDir sets the flock directory on the execute server to the session tmpdir. -func (a *sessionHost) pushLockDir() { +func (a *harness) pushLockDir() { if a.r.runtime == nil || a.r.runtime.Dispatcher == nil || a.sess.ID() == "" { return } } -var _ Core = (*sessionHost)(nil) // compile-time interface check +var _ Core = (*harness)(nil) // compile-time interface check var sweepTmpOnce sync.Once @@ -453,7 +453,7 @@ func NewAgentCore(cfg AgentCoreConfig) Core { log = olog.NewWriter("core", olog.LevelError+1, io.Discard, io.Discard) } - a := &sessionHost{ + a := &harness{ sess: session.New(session.Config{ ID: cfg.SessionID, Uname: cfg.Uname, @@ -495,7 +495,7 @@ func NewAgentCore(cfg AgentCoreConfig) Core { } // Close releases resources for this session, including its tmpdir. -func (a *sessionHost) Close() { +func (a *harness) Close() { a.log.Debug("Close() session=%q", a.sess.ID()) a.flushSave() if a.r.runtime != nil && a.r.runtime.Dispatcher != nil { @@ -512,7 +512,7 @@ func (a *sessionHost) Close() { } // execServer returns the execute server if available, or nil. -func (a *sessionHost) execServer() interface{} { +func (a *harness) execServer() interface{} { if a.r.runtime == nil || a.r.runtime.Dispatcher == nil { return nil } @@ -520,7 +520,7 @@ func (a *sessionHost) execServer() interface{} { return srv } -func (a *sessionHost) Detach() bool { +func (a *harness) Detach() bool { if srv := a.execServer(); srv != nil { if d, ok := srv.(interface{ Detach() bool }); ok { return d.Detach() @@ -529,7 +529,7 @@ func (a *sessionHost) Detach() bool { return false } -func (a *sessionHost) ListDetached() []DetachedInfo { +func (a *harness) ListDetached() []DetachedInfo { if srv := a.execServer(); srv != nil { type listDetacher interface { ListDetachedRaw() []any @@ -564,7 +564,7 @@ func (a *sessionHost) ListDetached() []DetachedInfo { return nil } -func (a *sessionHost) SignalDetached(pid, signal int) error { +func (a *harness) SignalDetached(pid, signal int) error { if srv := a.execServer(); srv != nil { type signaler interface { SignalDetached(int, syscall.Signal) error @@ -576,7 +576,7 @@ func (a *sessionHost) SignalDetached(pid, signal int) error { return fmt.Errorf("no execute server available") } -func (a *sessionHost) GetDetachedOutput(pid int) (string, error) { +func (a *harness) GetDetachedOutput(pid int) (string, error) { if srv := a.execServer(); srv != nil { type outputGetter interface { GetDetachedOutput(int) (string, error) @@ -588,7 +588,7 @@ func (a *sessionHost) GetDetachedOutput(pid int) (string, error) { return "", fmt.Errorf("no execute server available") } -func (a *sessionHost) DismissDetached(pid int) bool { +func (a *harness) DismissDetached(pid int) bool { if srv := a.execServer(); srv != nil { type dismisser interface { DismissDetached(int) bool @@ -600,7 +600,7 @@ func (a *sessionHost) DismissDetached(pid int) bool { return false } -func (a *sessionHost) InjectSystemEvent(content string) { +func (a *harness) InjectSystemEvent(content string) { a.Queue("\n" + content + "\n") } @@ -622,7 +622,7 @@ func classifyReaction(emoji string) (category, description string, positive bool } } -func (a *sessionHost) Reactions() map[string]string { +func (a *harness) Reactions() map[string]string { result := make(map[string]string) if a.r.history == nil { return result @@ -633,11 +633,11 @@ func (a *sessionHost) Reactions() map[string]string { return result } -func (a *sessionHost) React(emoji string) { +func (a *harness) React(emoji string) { _ = a.ReactTo("", emoji) } -func (a *sessionHost) ReactTo(responseID, emoji string) error { +func (a *harness) ReactTo(responseID, emoji string) error { if a.r.history == nil { return fmt.Errorf("no active session") } @@ -687,33 +687,33 @@ func (a *sessionHost) ReactTo(responseID, emoji string) error { return nil } -func (a *sessionHost) AgentName() string { +func (a *harness) AgentName() string { v := a.r.agentName a.log.Debug("AgentName() = %q", v) return v } -func (a *sessionHost) BackendName() string { +func (a *harness) BackendName() string { v := a.r.runtime.Backend.Name() a.log.Debug("BackendName() = %q", v) return v } -func (a *sessionHost) ModelName() string { +func (a *harness) ModelName() string { v := a.r.runtime.Backend.Model() a.log.Debug("ModelName() = %q", v) return v } -func (a *sessionHost) State() string { +func (a *harness) State() string { return a.sess.State() } -func (a *sessionHost) notifyChange() { +func (a *harness) notifyChange() { a.changeMu.Lock() a.changeCond.Broadcast() a.changeMu.Unlock() } -func (a *sessionHost) setState(state string) { +func (a *harness) setState(state string) { a.sess.SetState(state) a.log.Debug("state -> %q", state) a.notifyChange() @@ -721,7 +721,7 @@ func (a *sessionHost) setState(state string) { // WaitChange blocks until the named field changes from current, then returns // the new value. Returns ("", false) if ctx is cancelled. -func (a *sessionHost) WaitChange(ctx context.Context, field, current string) (string, bool) { +func (a *harness) WaitChange(ctx context.Context, field, current string) (string, bool) { read := func() string { switch field { case WatchState: @@ -759,14 +759,14 @@ func (a *sessionHost) WaitChange(ctx context.Context, field, current string) (st } -func (a *sessionHost) Reply() string { +func (a *harness) 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 (a *sessionHost) CWD() string { +func (a *harness) CWD() string { if a.sess.CWD() != "" { a.log.Debug("CWD() = %q", a.sess.CWD()) return a.sess.CWD() @@ -778,7 +778,7 @@ func (a *sessionHost) CWD() string { // SetCWD changes the working directory for tool execution and updates the // system prompt. Returns an error if the path does not exist. -func (a *sessionHost) SetCWD(dir string) error { +func (a *harness) SetCWD(dir string) error { a.log.Debug("SetCWD(%q)", dir) dir = paths.ExpandHome(dir) if dir != "" { @@ -806,7 +806,7 @@ func (a *sessionHost) SetCWD(dir string) error { // SetSessionID renames the session. It updates the in-memory ID, renames // persisted files on disk, and propagates to the execute server env. -func (a *sessionHost) SetSessionID(newID string) error { +func (a *harness) SetSessionID(newID string) error { a.log.Debug("SetSessionID(%q) old=%q", newID, a.sess.ID()) oldID := a.sess.ID() if oldID == newID { @@ -843,7 +843,7 @@ const defaultContextLength = 128000 const defaultToolResultMaxBytes = 131072 // autoCompactLimit returns the token threshold for auto-compaction (75%). -func (a *sessionHost) autoCompactLimit(ctx context.Context) int { +func (a *harness) autoCompactLimit(ctx context.Context) int { ctxLen := a.r.runtime.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength @@ -852,7 +852,7 @@ func (a *sessionHost) autoCompactLimit(ctx context.Context) int { } // autoWarnLimit returns the token threshold for a context-usage warning (60%). -func (a *sessionHost) autoWarnLimit(ctx context.Context) int { +func (a *harness) autoWarnLimit(ctx context.Context) int { ctxLen := a.r.runtime.Backend.ContextLength(ctx) if ctxLen <= 0 { ctxLen = defaultContextLength @@ -863,7 +863,7 @@ func (a *sessionHost) 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 (a *sessionHost) spawnContext(ctx context.Context) string { +func (a *harness) spawnContext(ctx context.Context) string { result := a.r.runtime.Hooks.Run(ctx, HookAgentSpawn, map[string]string{ "session_id": a.sess.ID(), "agent": a.r.agentName, @@ -886,7 +886,7 @@ func (a *sessionHost) 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 (a *sessionHost) runCompact(ctx context.Context, trigger string) (int, error) { +func (a *harness) runCompact(ctx context.Context, trigger string) (int, error) { payload := map[string]string{"session_id": a.sess.ID(), "trigger": trigger, "cwd": a.CWD()} pre := a.r.runtime.Hooks.Run(ctx, HookPreCompact, payload, a.log) if pre.Warning != "" { @@ -933,11 +933,11 @@ func (a *sessionHost) runCompact(ctx context.Context, trigger string) (int, erro return n, nil } -func (a *sessionHost) activeSessionPath(id, suffix string) string { +func (a *harness) activeSessionPath(id, suffix string) string { return filepath.Join(a.sessionsDir, "active", id+suffix) } -func (a *sessionHost) saveSession() { +func (a *harness) saveSession() { a.saveMu.Lock() a.saveDirty = true if a.saveTimer == nil { @@ -947,7 +947,7 @@ func (a *sessionHost) saveSession() { } // flushSave immediately persists the session if dirty. -func (a *sessionHost) flushSave() { +func (a *harness) flushSave() { a.saveMu.Lock() dirty := a.saveDirty a.saveDirty = false @@ -975,7 +975,7 @@ func (a *sessionHost) flushSave() { // SaveSession writes the current session state to the given path, including // backend and model metadata for external restore. -func (a *sessionHost) SaveSession(path string) error { +func (a *harness) SaveSession(path string) error { a.mu.RLock() defer a.mu.RUnlock() if a.r.history == nil { @@ -985,7 +985,7 @@ func (a *sessionHost) SaveSession(path string) error { a.r.runtime.Backend.Name(), a.r.runtime.Backend.Model(), a.CWD(), a.remote) } -func (a *sessionHost) getActionCancel() context.CancelCauseFunc { +func (a *harness) getActionCancel() context.CancelCauseFunc { if a := a.r.currentAction.Load(); a != nil { return a.cancel } @@ -994,7 +994,7 @@ func (a *sessionHost) getActionCancel() context.CancelCauseFunc { // Interrupt cancels the current in-progress agent turn. // Returns true if an action was running and was cancelled. -func (a *sessionHost) Interrupt(cause error) bool { +func (a *harness) Interrupt(cause error) bool { a.log.Debug("Interrupt() cause=%v", cause) if cancel := a.getActionCancel(); cancel != nil { cancel(cause) @@ -1003,7 +1003,7 @@ func (a *sessionHost) Interrupt(cause error) bool { return false } -func (a *sessionHost) Inject(prompt string) { +func (a *harness) 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 !a.pendingInject.CompareAndSwap(nil, &prompt) { @@ -1014,42 +1014,42 @@ func (a *sessionHost) Inject(prompt string) { a.emit(Event{Role: "user", Content: prompt}) } -func (a *sessionHost) injectRewrite(prompt string) { +func (a *harness) injectRewrite(prompt string) { a.pendingInject.Store(&prompt) a.emit(Event{Role: "info", Content: "\n"}) a.emit(Event{Role: "user", Content: prompt}) } -func (a *sessionHost) Queue(prompt string) { +func (a *harness) Queue(prompt string) { a.sess.Queue(prompt) a.sess.Bus().Publish("queued", prompt) } -func (a *sessionHost) drainQueue() { +func (a *harness) drainQueue() { if prompt, ok := a.sess.PopQueue(); ok { a.Submit(context.Background(), prompt) } } -func (a *sessionHost) Bus() *pubsub.Bus { +func (a *harness) Bus() *pubsub.Bus { return a.sess.Bus() } -func (a *sessionHost) emit(ev Event) { +func (a *harness) emit(ev Event) { a.sess.Bus().Publish("event", ev) } -func (a *sessionHost) PopQueue() (string, bool) { +func (a *harness) PopQueue() (string, bool) { return a.sess.PopQueue() } -func (a *sessionHost) IsRunning() bool { +func (a *harness) IsRunning() bool { v := a.r.currentAction.Load() != nil a.log.Debug("IsRunning() = %v", v) return v } -func (a *sessionHost) CtxSz() string { +func (a *harness) CtxSz() string { if a.r.history == nil { a.log.Debug("CtxSz() no session") return "no active session" @@ -1065,7 +1065,7 @@ func (a *sessionHost) CtxSz() string { return v } -func (a *sessionHost) Cost() string { +func (a *harness) Cost() string { if a.r.history == nil { return "no active session" } @@ -1073,7 +1073,7 @@ func (a *sessionHost) Cost() string { a.r.history.LastTurnCostUSD, a.r.history.SessionCostUSD) } -func (a *sessionHost) Usage() string { +func (a *harness) Usage() string { if a.r.history == nil { a.log.Debug("Usage() no session") return "no active session" @@ -1091,7 +1091,7 @@ func (a *sessionHost) Usage() string { return str } -func (a *sessionHost) Context() []backend.Message { +func (a *harness) Context() []backend.Message { a.mu.RLock() var msgs []backend.Message if a.r.history != nil { @@ -1104,30 +1104,30 @@ func (a *sessionHost) Context() []backend.Message { return msgs } -func (a *sessionHost) SystemPrompt() string { +func (a *harness) SystemPrompt() string { a.log.Debug("SystemPrompt() len=%d", len(a.r.runtime.Preamble)) return a.r.runtime.Preamble } -func (a *sessionHost) GenerationParams() backend.GenerationParams { +func (a *harness) GenerationParams() backend.GenerationParams { a.mu.RLock() defer a.mu.RUnlock() return a.r.runtime.GenParams } -func (a *sessionHost) CompactionModel() string { +func (a *harness) CompactionModel() string { a.mu.RLock() defer a.mu.RUnlock() return a.r.runtime.CompactionModel } -func (a *sessionHost) SetCompactionModel(model string) { +func (a *harness) SetCompactionModel(model string) { a.mu.Lock() defer a.mu.Unlock() a.r.runtime.CompactionModel = model } -func (a *sessionHost) SetGenerationParams(params backend.GenerationParams) error { +func (a *harness) SetGenerationParams(params backend.GenerationParams) error { if a.IsRunning() { return fmt.Errorf("cannot change params while agent is running") } @@ -1137,7 +1137,7 @@ func (a *sessionHost) SetGenerationParams(params backend.GenerationParams) error return nil } -func (a *sessionHost) ListModels() string { +func (a *harness) ListModels() string { a.log.Debug("ListModels()") models := a.r.runtime.Backend.Models(context.Background()) slices.Sort(models) @@ -1166,7 +1166,7 @@ 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 (a *sessionHost) Submit(ctx context.Context, input string) { +func (a *harness) Submit(ctx context.Context, input string) { defer func() { if r := recover(); r != nil { a.log.Error("panic: %v\n%s", r, debug.Stack()) @@ -1213,7 +1213,7 @@ func (a *sessionHost) Submit(ctx context.Context, input string) { // executeTurn runs a single agent turn and returns the next prompt to execute, // or "" if there is nothing more to do. -func (a *sessionHost) executeTurn(ctx context.Context, input string) string { +func (a *harness) executeTurn(ctx context.Context, input string) string { a.emit(Event{Role: "user", Content: input}) hookResult := a.r.runtime.Hooks.Run(ctx, HookPreTurn, map[string]string{ diff --git a/agent/waitchange_test.go b/agent/waitchange_test.go index f0af55e..98c2175 100644 --- a/agent/waitchange_test.go +++ b/agent/waitchange_test.go @@ -12,10 +12,10 @@ import ( ) // newTestCore returns a minimal agent with session wired up. -func newTestCore(initialState string) *sessionHost { +func newTestCore(initialState string) *harness { sess := session.New(session.Config{ID: "test"}) sess.SetState(initialState) - a := &sessionHost{ + a := &harness{ sess: sess, log: olog.NewWriter("test", olog.LevelError+1, io.Discard, io.Discard), }