session: remove Session interface, export concrete struct

Session is now an exported struct with unexported fields.
Methods ARE the API — no interface declaration in the package.
Consumers that need mockability define their own interface
(implicit satisfaction, standard Go pattern).
This commit is contained in:
Levi Neely 2026-07-29 19:25:14 +02:00
parent 1664644023
commit 5db3eb5138
5 changed files with 70 additions and 220 deletions

View File

@ -12,7 +12,7 @@ import (
) )
func (a *harness) handleCommand(ctx context.Context, input string) bool { func (a *Session) handleCommand(ctx context.Context, input string) bool {
if !strings.HasPrefix(input, "/") { if !strings.HasPrefix(input, "/") {
return false return false
} }

View File

@ -83,7 +83,7 @@ func defaultBE() *mockBackend {
// newCore builds a minimal *agent for tests, bypassing loadSystemPrompt // newCore builds a minimal *agent for tests, bypassing loadSystemPrompt
// by directly setting Preamble on Runtime. // by directly setting Preamble on Runtime.
func newCore(t *testing.T, be backend.Backend, hooks Hooks) *harness { func newCore(t *testing.T, be backend.Backend, hooks Hooks) *Session {
t.Helper() t.Helper()
t.Setenv("OLLIE", "") t.Setenv("OLLIE", "")
if be == nil { if be == nil {
@ -107,11 +107,11 @@ func newCore(t *testing.T, be backend.Backend, hooks Hooks) *harness {
NewDispatcher: tools.NewDispatcher, NewDispatcher: tools.NewDispatcher,
}) })
t.Cleanup(c.Close) t.Cleanup(c.Close)
return c.(*harness) return c
} }
// collectEvents runs Submit synchronously and returns all emitted events. // collectEvents runs Submit synchronously and returns all emitted events.
func collectEvents(ctx context.Context, c Session, input string) []Event { func collectEvents(ctx context.Context, c *Session, input string) []Event {
var mu sync.Mutex var mu sync.Mutex
var evs []Event var evs []Event
sub := c.Bus().Subscribe("event", func(ev Event) { sub := c.Bus().Subscribe("event", func(ev Event) {
@ -138,7 +138,7 @@ func byRole(evs []Event, role string) []string {
} }
// waitState blocks until c.State() == want, failing after 2 s. // waitState blocks until c.State() == want, failing after 2 s.
func waitState(t *testing.T, c Session, want string) { func waitState(t *testing.T, c *Session, want string) {
t.Helper() t.Helper()
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel() defer cancel()
@ -1788,7 +1788,7 @@ func TestSetSessionID_UpdatesPreamble(t *testing.T) {
CWD: t.TempDir(), CWD: t.TempDir(),
Runtime: env, Runtime: env,
NewDispatcher: tools.NewDispatcher, NewDispatcher: tools.NewDispatcher,
}).(*harness) })
t.Cleanup(c.Close) t.Cleanup(c.Close)
newID := NewSessionID() newID := NewSessionID()
@ -2270,7 +2270,7 @@ func (m *mockEnvServer) get(k string) string {
return m.env[k] return m.env[k]
} }
func newCoreWithExecServer(t *testing.T, srv *mockEnvServer) *harness { func newCoreWithExecServer(t *testing.T, srv *mockEnvServer) *Session {
t.Helper() t.Helper()
t.Setenv("OLLIE", "") t.Setenv("OLLIE", "")
d := tools.NewDispatcher() d := tools.NewDispatcher()
@ -2291,7 +2291,7 @@ func newCoreWithExecServer(t *testing.T, srv *mockEnvServer) *harness {
NewDispatcher: tools.NewDispatcher, NewDispatcher: tools.NewDispatcher,
}) })
t.Cleanup(c.Close) t.Cleanup(c.Close)
return c.(*harness) return c
} }
func TestSetEnv_PropagatestoExecuteServer(t *testing.T) { func TestSetEnv_PropagatestoExecuteServer(t *testing.T) {

View File

@ -308,7 +308,7 @@ type Config struct {
// harness is the Core implementation. It owns all harness and session state // harness is the Core implementation. It owns all harness and session state
// but has no knowledge of how output is rendered. // but has no knowledge of how output is rendered.
type harness struct { type Session struct {
// Session-level state // Session-level state
id string id string
uname string uname string
@ -350,12 +350,12 @@ type harness struct {
// ToolCallCount returns the total number of tool calls executed in this // ToolCallCount returns the total number of tool calls executed in this
// session. The counter is monotonically increasing and never resets. // session. The counter is monotonically increasing and never resets.
// Blocked calls (pre-tool hook exit 2) are not counted. // Blocked calls (pre-tool hook exit 2) are not counted.
func (a *harness) ToolCallCount() int64 { func (a *Session) ToolCallCount() int64 {
return a.toolCallCount.Load() return a.toolCallCount.Load()
} }
// SetEnv stores a session-scoped variable and propagates it to the execute server. // SetEnv stores a session-scoped variable and propagates it to the execute server.
func (a *harness) SetEnv(key, value string) { func (a *Session) SetEnv(key, value string) {
a.envMu.Lock() a.envMu.Lock()
a.env[key] = value a.env[key] = value
a.envMu.Unlock() a.envMu.Unlock()
@ -370,7 +370,7 @@ func (a *harness) SetEnv(key, value string) {
} }
// pushSessionEnv injects OLLIE_SESSION_ID into the execute server subprocess env. // pushSessionEnv injects OLLIE_SESSION_ID into the execute server subprocess env.
func (a *harness) pushSessionEnv() { func (a *Session) pushSessionEnv() {
if a.r.runtime == nil || a.r.runtime.Dispatcher == nil || a.id == "" { if a.r.runtime == nil || a.r.runtime.Dispatcher == nil || a.id == "" {
return return
} }
@ -385,14 +385,13 @@ func (a *harness) pushSessionEnv() {
} }
// pushLockDir sets the flock directory on the execute server to the session tmpdir. // pushLockDir sets the flock directory on the execute server to the session tmpdir.
func (a *harness) pushLockDir() { func (a *Session) pushLockDir() {
if a.r.runtime == nil || a.r.runtime.Dispatcher == nil || a.id == "" { if a.r.runtime == nil || a.r.runtime.Dispatcher == nil || a.id == "" {
return return
} }
} }
var _ Session = (*harness)(nil) // compile-time interface check
var sweepTmpOnce sync.Once var sweepTmpOnce sync.Once
@ -436,7 +435,7 @@ func sweepStaleTmpDirs() {
} }
// New creates an agent from the given configuration. // New creates an agent from the given configuration.
func New(cfg Config) Session { func New(cfg Config) *Session {
sweepStaleTmpDirs() sweepStaleTmpDirs()
if cfg.ModelName != "" { if cfg.ModelName != "" {
cfg.Backend.SetModel(cfg.ModelName) cfg.Backend.SetModel(cfg.ModelName)
@ -468,7 +467,7 @@ func New(cfg Config) Session {
log = olog.NewWriter("core", olog.LevelError+1, io.Discard, io.Discard) log = olog.NewWriter("core", olog.LevelError+1, io.Discard, io.Discard)
} }
a := &harness{ a := &Session{
id: cfg.SessionID, id: cfg.SessionID,
uname: cfg.Uname, uname: cfg.Uname,
cwd: paths.ExpandHome(cfg.CWD), cwd: paths.ExpandHome(cfg.CWD),
@ -511,7 +510,7 @@ func New(cfg Config) Session {
} }
// Close releases resources for this session, including its tmpdir. // Close releases resources for this session, including its tmpdir.
func (a *harness) Close() { func (a *Session) Close() {
a.log.Debug("Close() session=%q", a.id) a.log.Debug("Close() session=%q", a.id)
a.flushSave() a.flushSave()
if a.r.runtime != nil && a.r.runtime.Dispatcher != nil { if a.r.runtime != nil && a.r.runtime.Dispatcher != nil {
@ -528,7 +527,7 @@ func (a *harness) Close() {
} }
// execServer returns the execute server if available, or nil. // execServer returns the execute server if available, or nil.
func (a *harness) execServer() interface{} { func (a *Session) execServer() interface{} {
if a.r.runtime == nil || a.r.runtime.Dispatcher == nil { if a.r.runtime == nil || a.r.runtime.Dispatcher == nil {
return nil return nil
} }
@ -536,7 +535,7 @@ func (a *harness) execServer() interface{} {
return srv return srv
} }
func (a *harness) Detach() bool { func (a *Session) Detach() bool {
if srv := a.execServer(); srv != nil { if srv := a.execServer(); srv != nil {
if d, ok := srv.(interface{ Detach() bool }); ok { if d, ok := srv.(interface{ Detach() bool }); ok {
return d.Detach() return d.Detach()
@ -545,7 +544,7 @@ func (a *harness) Detach() bool {
return false return false
} }
func (a *harness) ListDetached() []DetachedInfo { func (a *Session) ListDetached() []DetachedInfo {
if srv := a.execServer(); srv != nil { if srv := a.execServer(); srv != nil {
type listDetacher interface { type listDetacher interface {
ListDetachedRaw() []any ListDetachedRaw() []any
@ -580,7 +579,7 @@ func (a *harness) ListDetached() []DetachedInfo {
return nil return nil
} }
func (a *harness) SignalDetached(pid, signal int) error { func (a *Session) SignalDetached(pid, signal int) error {
if srv := a.execServer(); srv != nil { if srv := a.execServer(); srv != nil {
type signaler interface { type signaler interface {
SignalDetached(int, syscall.Signal) error SignalDetached(int, syscall.Signal) error
@ -592,7 +591,7 @@ func (a *harness) SignalDetached(pid, signal int) error {
return fmt.Errorf("no execute server available") return fmt.Errorf("no execute server available")
} }
func (a *harness) GetDetachedOutput(pid int) (string, error) { func (a *Session) GetDetachedOutput(pid int) (string, error) {
if srv := a.execServer(); srv != nil { if srv := a.execServer(); srv != nil {
type outputGetter interface { type outputGetter interface {
GetDetachedOutput(int) (string, error) GetDetachedOutput(int) (string, error)
@ -604,7 +603,7 @@ func (a *harness) GetDetachedOutput(pid int) (string, error) {
return "", fmt.Errorf("no execute server available") return "", fmt.Errorf("no execute server available")
} }
func (a *harness) DismissDetached(pid int) bool { func (a *Session) DismissDetached(pid int) bool {
if srv := a.execServer(); srv != nil { if srv := a.execServer(); srv != nil {
type dismisser interface { type dismisser interface {
DismissDetached(int) bool DismissDetached(int) bool
@ -616,7 +615,7 @@ func (a *harness) DismissDetached(pid int) bool {
return false return false
} }
func (a *harness) InjectSystemEvent(content string) { func (a *Session) InjectSystemEvent(content string) {
a.Queue("<detached-process-result>\n" + content + "\n</detached-process-result>") a.Queue("<detached-process-result>\n" + content + "\n</detached-process-result>")
} }
@ -638,7 +637,7 @@ func classifyReaction(emoji string) (category, description string, positive bool
} }
} }
func (a *harness) Reactions() map[string]string { func (a *Session) Reactions() map[string]string {
result := make(map[string]string) result := make(map[string]string)
if a.r.history == nil { if a.r.history == nil {
return result return result
@ -649,11 +648,11 @@ func (a *harness) Reactions() map[string]string {
return result return result
} }
func (a *harness) React(emoji string) { func (a *Session) React(emoji string) {
_ = a.ReactTo("", emoji) _ = a.ReactTo("", emoji)
} }
func (a *harness) ReactTo(responseID, emoji string) error { func (a *Session) ReactTo(responseID, emoji string) error {
if a.r.history == nil { if a.r.history == nil {
return fmt.Errorf("no active session") return fmt.Errorf("no active session")
} }
@ -703,33 +702,33 @@ func (a *harness) ReactTo(responseID, emoji string) error {
return nil return nil
} }
func (a *harness) AgentName() string { func (a *Session) AgentName() string {
v := a.r.agentName v := a.r.agentName
a.log.Debug("AgentName() = %q", v) a.log.Debug("AgentName() = %q", v)
return v return v
} }
func (a *harness) BackendName() string { func (a *Session) BackendName() string {
v := a.r.runtime.Backend.Name() v := a.r.runtime.Backend.Name()
a.log.Debug("BackendName() = %q", v) a.log.Debug("BackendName() = %q", v)
return v return v
} }
func (a *harness) ModelName() string { func (a *Session) ModelName() string {
v := a.r.runtime.Backend.Model() v := a.r.runtime.Backend.Model()
a.log.Debug("ModelName() = %q", v) a.log.Debug("ModelName() = %q", v)
return v return v
} }
func (a *harness) State() string { func (a *Session) State() string {
return a.state return a.state
} }
func (a *harness) notifyChange() { func (a *Session) notifyChange() {
a.changeMu.Lock() a.changeMu.Lock()
a.changeCond.Broadcast() a.changeCond.Broadcast()
a.changeMu.Unlock() a.changeMu.Unlock()
} }
func (a *harness) setState(state string) { func (a *Session) setState(state string) {
a.state = state a.state = state
a.log.Debug("state -> %q", state) a.log.Debug("state -> %q", state)
a.notifyChange() a.notifyChange()
@ -737,7 +736,7 @@ func (a *harness) setState(state string) {
// WaitChange blocks until the named field changes from current, then returns // WaitChange blocks until the named field changes from current, then returns
// the new value. Returns ("", false) if ctx is cancelled. // the new value. Returns ("", false) if ctx is cancelled.
func (a *harness) WaitChange(ctx context.Context, field, current string) (string, bool) { func (a *Session) WaitChange(ctx context.Context, field, current string) (string, bool) {
read := func() string { read := func() string {
switch field { switch field {
case WatchState: case WatchState:
@ -775,14 +774,14 @@ func (a *harness) WaitChange(ctx context.Context, field, current string) (string
} }
func (a *harness) Reply() string { func (a *Session) Reply() string {
r := a.reply r := a.reply
a.log.Debug("Reply() len=%d", len(r)) a.log.Debug("Reply() len=%d", len(r))
return r return r
} }
// CWD returns the current working directory for tool execution. // CWD returns the current working directory for tool execution.
func (a *harness) CWD() string { func (a *Session) CWD() string {
if a.cwd != "" { if a.cwd != "" {
a.log.Debug("CWD() = %q", a.cwd) a.log.Debug("CWD() = %q", a.cwd)
return a.cwd return a.cwd
@ -794,7 +793,7 @@ func (a *harness) CWD() string {
// SetCWD changes the working directory for tool execution and updates the // SetCWD changes the working directory for tool execution and updates the
// system prompt. Returns an error if the path does not exist. // system prompt. Returns an error if the path does not exist.
func (a *harness) SetCWD(dir string) error { func (a *Session) SetCWD(dir string) error {
a.log.Debug("SetCWD(%q)", dir) a.log.Debug("SetCWD(%q)", dir)
dir = paths.ExpandHome(dir) dir = paths.ExpandHome(dir)
if dir != "" { if dir != "" {
@ -822,7 +821,7 @@ func (a *harness) SetCWD(dir string) error {
// SetSessionID renames the session. It updates the in-memory ID, renames // SetSessionID renames the session. It updates the in-memory ID, renames
// persisted files on disk, and propagates to the execute server env. // persisted files on disk, and propagates to the execute server env.
func (a *harness) SetSessionID(newID string) error { func (a *Session) SetSessionID(newID string) error {
a.log.Debug("SetSessionID(%q) old=%q", newID, a.id) a.log.Debug("SetSessionID(%q) old=%q", newID, a.id)
oldID := a.id oldID := a.id
if oldID == newID { if oldID == newID {
@ -859,7 +858,7 @@ const defaultContextLength = 128000
const defaultToolResultMaxBytes = 131072 const defaultToolResultMaxBytes = 131072
// autoCompactLimit returns the token threshold for auto-compaction (75%). // autoCompactLimit returns the token threshold for auto-compaction (75%).
func (a *harness) autoCompactLimit(ctx context.Context) int { func (a *Session) autoCompactLimit(ctx context.Context) int {
ctxLen := a.r.runtime.Backend.ContextLength(ctx) ctxLen := a.r.runtime.Backend.ContextLength(ctx)
if ctxLen <= 0 { if ctxLen <= 0 {
ctxLen = defaultContextLength ctxLen = defaultContextLength
@ -868,7 +867,7 @@ func (a *harness) autoCompactLimit(ctx context.Context) int {
} }
// autoWarnLimit returns the token threshold for a context-usage warning (60%). // autoWarnLimit returns the token threshold for a context-usage warning (60%).
func (a *harness) autoWarnLimit(ctx context.Context) int { func (a *Session) autoWarnLimit(ctx context.Context) int {
ctxLen := a.r.runtime.Backend.ContextLength(ctx) ctxLen := a.r.runtime.Backend.ContextLength(ctx)
if ctxLen <= 0 { if ctxLen <= 0 {
ctxLen = defaultContextLength ctxLen = defaultContextLength
@ -879,7 +878,7 @@ func (a *harness) autoWarnLimit(ctx context.Context) int {
// spawnContext assembles the agent context injected at each session refresh // spawnContext assembles the agent context injected at each session refresh
// point (session start, post-clear, post-compaction). It combines the // point (session start, post-clear, post-compaction). It combines the
// agent-specific prompt with any agentSpawn hook output. // agent-specific prompt with any agentSpawn hook output.
func (a *harness) spawnContext(ctx context.Context) string { func (a *Session) spawnContext(ctx context.Context) string {
result := a.r.runtime.Hooks.Run(ctx, HookAgentSpawn, map[string]string{ result := a.r.runtime.Hooks.Run(ctx, HookAgentSpawn, map[string]string{
"session_id": a.id, "session_id": a.id,
"agent": a.r.agentName, "agent": a.r.agentName,
@ -902,7 +901,7 @@ func (a *harness) spawnContext(ctx context.Context) string {
// runCompact executes a full compaction cycle: pre-hook, compact, spawn-context // runCompact executes a full compaction cycle: pre-hook, compact, spawn-context
// re-injection, post-hook. Returns (n compacted, error). Returns (0, nil) if // 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. // the pre-hook blocked or there was nothing to compact. Caller manages setState.
func (a *harness) runCompact(ctx context.Context, trigger string) (int, error) { func (a *Session) runCompact(ctx context.Context, trigger string) (int, error) {
payload := map[string]string{"session_id": a.id, "trigger": trigger, "cwd": a.CWD()} payload := map[string]string{"session_id": a.id, "trigger": trigger, "cwd": a.CWD()}
pre := a.r.runtime.Hooks.Run(ctx, HookPreCompact, payload, a.log) pre := a.r.runtime.Hooks.Run(ctx, HookPreCompact, payload, a.log)
if pre.Warning != "" { if pre.Warning != "" {
@ -949,11 +948,11 @@ func (a *harness) runCompact(ctx context.Context, trigger string) (int, error) {
return n, nil return n, nil
} }
func (a *harness) activeSessionPath(id, suffix string) string { func (a *Session) activeSessionPath(id, suffix string) string {
return filepath.Join(a.sessionsDir, "active", id+suffix) return filepath.Join(a.sessionsDir, "active", id+suffix)
} }
func (a *harness) saveSession() { func (a *Session) saveSession() {
a.saveMu.Lock() a.saveMu.Lock()
a.saveDirty = true a.saveDirty = true
if a.saveTimer == nil { if a.saveTimer == nil {
@ -963,7 +962,7 @@ func (a *harness) saveSession() {
} }
// flushSave immediately persists the session if dirty. // flushSave immediately persists the session if dirty.
func (a *harness) flushSave() { func (a *Session) flushSave() {
a.saveMu.Lock() a.saveMu.Lock()
dirty := a.saveDirty dirty := a.saveDirty
a.saveDirty = false a.saveDirty = false
@ -991,7 +990,7 @@ func (a *harness) flushSave() {
// SaveSession writes the current session state to the given path, including // SaveSession writes the current session state to the given path, including
// backend and model metadata for external restore. // backend and model metadata for external restore.
func (a *harness) SaveSession(path string) error { func (a *Session) SaveSession(path string) error {
a.mu.RLock() a.mu.RLock()
defer a.mu.RUnlock() defer a.mu.RUnlock()
if a.r.history == nil { if a.r.history == nil {
@ -1001,7 +1000,7 @@ func (a *harness) SaveSession(path string) error {
a.r.runtime.Backend.Name(), a.r.runtime.Backend.Model(), a.CWD(), a.remote) a.r.runtime.Backend.Name(), a.r.runtime.Backend.Model(), a.CWD(), a.remote)
} }
func (a *harness) getActionCancel() context.CancelCauseFunc { func (a *Session) getActionCancel() context.CancelCauseFunc {
if a := a.r.currentAction.Load(); a != nil { if a := a.r.currentAction.Load(); a != nil {
return a.cancel return a.cancel
} }
@ -1010,7 +1009,7 @@ func (a *harness) getActionCancel() context.CancelCauseFunc {
// Interrupt cancels the current in-progress agent turn. // Interrupt cancels the current in-progress agent turn.
// Returns true if an action was running and was cancelled. // Returns true if an action was running and was cancelled.
func (a *harness) Interrupt(cause error) bool { func (a *Session) Interrupt(cause error) bool {
a.log.Debug("Interrupt() cause=%v", cause) a.log.Debug("Interrupt() cause=%v", cause)
if cancel := a.getActionCancel(); cancel != nil { if cancel := a.getActionCancel(); cancel != nil {
cancel(cause) cancel(cause)
@ -1019,7 +1018,7 @@ func (a *harness) Interrupt(cause error) bool {
return false return false
} }
func (a *harness) Inject(prompt string) { func (a *Session) Inject(prompt string) {
// If an inject is already pending, fall back to the normal FIFO so nothing // 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. // is lost. Use CompareAndSwap to avoid a race between the nil check and store.
if !a.pendingInject.CompareAndSwap(nil, &prompt) { if !a.pendingInject.CompareAndSwap(nil, &prompt) {
@ -1030,42 +1029,42 @@ func (a *harness) Inject(prompt string) {
a.emit(Event{Role: "user", Content: prompt}) a.emit(Event{Role: "user", Content: prompt})
} }
func (a *harness) injectRewrite(prompt string) { func (a *Session) injectRewrite(prompt string) {
a.pendingInject.Store(&prompt) a.pendingInject.Store(&prompt)
a.emit(Event{Role: "info", Content: "\n"}) a.emit(Event{Role: "info", Content: "\n"})
a.emit(Event{Role: "user", Content: prompt}) a.emit(Event{Role: "user", Content: prompt})
} }
func (a *harness) Queue(prompt string) { func (a *Session) Queue(prompt string) {
a.fifo.Push(prompt) a.fifo.Push(prompt)
a.bus.Publish("queued", prompt) a.bus.Publish("queued", prompt)
} }
func (a *harness) drainQueue() { func (a *Session) drainQueue() {
if prompt, ok := a.fifo.Pop(); ok { if prompt, ok := a.fifo.Pop(); ok {
a.Submit(context.Background(), prompt) a.Submit(context.Background(), prompt)
} }
} }
func (a *harness) Bus() *pubsub.Bus { func (a *Session) Bus() *pubsub.Bus {
return a.bus return a.bus
} }
func (a *harness) emit(ev Event) { func (a *Session) emit(ev Event) {
a.bus.Publish("event", ev) a.bus.Publish("event", ev)
} }
func (a *harness) PopQueue() (string, bool) { func (a *Session) PopQueue() (string, bool) {
return a.fifo.Pop() return a.fifo.Pop()
} }
func (a *harness) IsRunning() bool { func (a *Session) IsRunning() bool {
v := a.r.currentAction.Load() != nil v := a.r.currentAction.Load() != nil
a.log.Debug("IsRunning() = %v", v) a.log.Debug("IsRunning() = %v", v)
return v return v
} }
func (a *harness) CtxSz() string { func (a *Session) CtxSz() string {
if a.r.history == nil { if a.r.history == nil {
a.log.Debug("CtxSz() no session") a.log.Debug("CtxSz() no session")
return "no active session" return "no active session"
@ -1081,7 +1080,7 @@ func (a *harness) CtxSz() string {
return v return v
} }
func (a *harness) Cost() string { func (a *Session) Cost() string {
if a.r.history == nil { if a.r.history == nil {
return "no active session" return "no active session"
} }
@ -1089,7 +1088,7 @@ func (a *harness) Cost() string {
a.r.history.LastTurnCostUSD, a.r.history.SessionCostUSD) a.r.history.LastTurnCostUSD, a.r.history.SessionCostUSD)
} }
func (a *harness) Usage() string { func (a *Session) Usage() string {
if a.r.history == nil { if a.r.history == nil {
a.log.Debug("Usage() no session") a.log.Debug("Usage() no session")
return "no active session" return "no active session"
@ -1107,7 +1106,7 @@ func (a *harness) Usage() string {
return str return str
} }
func (a *harness) Context() []backend.Message { func (a *Session) Context() []backend.Message {
a.mu.RLock() a.mu.RLock()
var msgs []backend.Message var msgs []backend.Message
if a.r.history != nil { if a.r.history != nil {
@ -1120,30 +1119,30 @@ func (a *harness) Context() []backend.Message {
return msgs return msgs
} }
func (a *harness) SystemPrompt() string { func (a *Session) SystemPrompt() string {
a.log.Debug("SystemPrompt() len=%d", len(a.r.runtime.Preamble)) a.log.Debug("SystemPrompt() len=%d", len(a.r.runtime.Preamble))
return a.r.runtime.Preamble return a.r.runtime.Preamble
} }
func (a *harness) GenerationParams() backend.GenerationParams { func (a *Session) GenerationParams() backend.GenerationParams {
a.mu.RLock() a.mu.RLock()
defer a.mu.RUnlock() defer a.mu.RUnlock()
return a.r.runtime.GenParams return a.r.runtime.GenParams
} }
func (a *harness) CompactionModel() string { func (a *Session) CompactionModel() string {
a.mu.RLock() a.mu.RLock()
defer a.mu.RUnlock() defer a.mu.RUnlock()
return a.r.runtime.CompactionModel return a.r.runtime.CompactionModel
} }
func (a *harness) SetCompactionModel(model string) { func (a *Session) SetCompactionModel(model string) {
a.mu.Lock() a.mu.Lock()
defer a.mu.Unlock() defer a.mu.Unlock()
a.r.runtime.CompactionModel = model a.r.runtime.CompactionModel = model
} }
func (a *harness) SetGenerationParams(params backend.GenerationParams) error { func (a *Session) SetGenerationParams(params backend.GenerationParams) error {
if a.IsRunning() { if a.IsRunning() {
return fmt.Errorf("cannot change params while agent is running") return fmt.Errorf("cannot change params while agent is running")
} }
@ -1153,7 +1152,7 @@ func (a *harness) SetGenerationParams(params backend.GenerationParams) error {
return nil return nil
} }
func (a *harness) ListModels() string { func (a *Session) ListModels() string {
a.log.Debug("ListModels()") a.log.Debug("ListModels()")
models := a.r.runtime.Backend.Models(context.Background()) models := a.r.runtime.Backend.Models(context.Background())
slices.Sort(models) slices.Sort(models)
@ -1182,7 +1181,7 @@ func firstSentence(s string) string {
// //
// Continuations (post-turn hook context, unconsumed inject, FIFO drain) are // Continuations (post-turn hook context, unconsumed inject, FIFO drain) are
// handled via an explicit loop rather than recursion to avoid stack growth. // handled via an explicit loop rather than recursion to avoid stack growth.
func (a *harness) Submit(ctx context.Context, input string) { func (a *Session) Submit(ctx context.Context, input string) {
defer func() { defer func() {
if r := recover(); r != nil { if r := recover(); r != nil {
a.log.Error("panic: %v\n%s", r, debug.Stack()) a.log.Error("panic: %v\n%s", r, debug.Stack())
@ -1229,7 +1228,7 @@ func (a *harness) Submit(ctx context.Context, input string) {
// executeTurn runs a single agent turn and returns the next prompt to execute, // executeTurn runs a single agent turn and returns the next prompt to execute,
// or "" if there is nothing more to do. // or "" if there is nothing more to do.
func (a *harness) executeTurn(ctx context.Context, input string) string { func (a *Session) executeTurn(ctx context.Context, input string) string {
a.emit(Event{Role: "user", Content: input}) a.emit(Event{Role: "user", Content: input})
hookResult := a.r.runtime.Hooks.Run(ctx, HookPreTurn, map[string]string{ hookResult := a.r.runtime.Hooks.Run(ctx, HookPreTurn, map[string]string{

View File

@ -1,14 +1,10 @@
package session package session
import ( import (
"context"
"errors" "errors"
"github.com/simonfxr/pubsub"
"ollie/backend"
) )
// WatchField names supported by Core.WaitChange. // WatchField names supported by Session.WaitChange.
const ( const (
WatchState = "state" WatchState = "state"
WatchUsage = "usage" WatchUsage = "usage"
@ -32,151 +28,6 @@ type Event struct {
// EventHandler receives events from the agent. // EventHandler receives events from the agent.
type EventHandler func(Event) type EventHandler func(Event)
// Session is the interface between a frontend (TUI, HTTP handler, etc.) and the
// agent engine. All output from the agent is delivered via the event bus.
type Session interface {
// Submit processes one line of user input. Slash commands and shell
// shortcuts are dispatched synchronously; any other input starts an agent
// turn that publishes events to the bus until the turn is complete.
// After the turn, any queued prompts are drained sequentially.
Submit(ctx context.Context, input string)
// Interrupt cancels the current in-progress agent turn.
// Returns true if an action was running and was cancelled.
Interrupt(cause error) bool
// Inject sends a message that will be appended to the next tool result
// as a user interruption. If no turn is running, it is silently dropped.
Inject(prompt string)
// Queue pushes a prompt onto the FIFO for execution after the current
// turn completes.
Queue(prompt string)
// Bus returns the session event bus.
Bus() *pubsub.Bus
// PopQueue removes and returns the next queued prompt.
// Returns ("", false) if the queue is empty.
PopQueue() (string, bool)
// IsRunning returns true if an agent turn is currently in progress.
IsRunning() bool
// State returns the current agent state: "idle", "thinking", or "calling: <tool>".
State() string
// Reply returns the assistant text from the most recently completed turn.
// Cleared when a new prompt is submitted.
Reply() string
// AgentName returns the name of the active agent.
AgentName() string
// BackendName returns the name of the active backend (e.g. "anthropic", "ollama").
BackendName() string
// ModelName returns the name of the active model.
ModelName() string
// CtxSz returns the estimated context size as a one-line summary.
CtxSz() string
// Usage returns billed token counts as a one-line summary.
Usage() string
// Cost returns per-turn and session cost as a two-line key=value summary.
Cost() string
// ListModels returns available model names, one per line.
ListModels() string
// CWD returns the current working directory used for tool execution.
CWD() string
// SetCWD changes the working directory for tool execution and
// updates the system prompt. Returns an error if the path does not exist.
SetCWD(dir string) error
// SetSessionID renames the session: updates the in-memory ID, renames
// persisted files on disk, and propagates to the execute server env.
SetSessionID(newID string) error
// Context returns the current message history as it would be sent to the
// backend: system prompt prepended, stale reads pruned. Does not include
// tool definitions (see Tools) or generation params (see GenerationParams).
Context() []backend.Message
// SystemPrompt returns the fully rendered system prompt for this session.
SystemPrompt() string
// GenerationParams returns the current sampling parameters.
GenerationParams() backend.GenerationParams
// SetGenerationParams replaces the current sampling parameters.
// Returns an error if the agent is currently running.
SetGenerationParams(params backend.GenerationParams) error
// CompactionModel returns the model override used for context compaction.
CompactionModel() string
// SetCompactionModel changes the model used for context compaction.
SetCompactionModel(model string)
// SetEnv injects a session-scoped environment variable into the shell
// subprocesses. Does not affect the daemon process environment.
SetEnv(key, value string)
// WaitChange blocks until the named field changes from current, then returns
// the new value. Returns ("", false) if ctx is cancelled before a change.
// Supported fields: WatchState, WatchUsage, WatchCtxSz, WatchCWD.
WaitChange(ctx context.Context, field, current string) (string, bool)
// ToolCallCount returns the total number of tool calls executed in this
// session since the agent was created. The counter is monotonically
// increasing and never resets. Blocked calls (pre-tool hook exit 2) are
// not counted. Timed-out or erroring calls are counted because execution
// was attempted. Use modulo arithmetic in hooks to fire every N calls:
// [ $(($(cat tcct) % 10)) -eq 0 ] && ...
ToolCallCount() int64
// SaveSession writes the current session state to the given path.
// The file includes all messages, agent/backend/model metadata, and task state.
SaveSession(path string) error
// Close releases resources associated with the session, including its
// temporary directory under /tmp/ollie/.
Close()
// Detach detaches the currently running process from the agent.
// The process continues running; returns false if nothing is executing.
Detach() bool
// ListDetached returns info about all detached processes.
ListDetached() []DetachedInfo
// SignalDetached sends a signal to a detached process by PID.
SignalDetached(pid, signal int) error
// GetDetachedOutput returns the ring buffer output for a detached process.
GetDetachedOutput(pid int) (string, error)
// DismissDetached removes an exited process from the list.
DismissDetached(pid int) bool
// InjectSystemEvent appends a system-originated message to the session
// context and emits it on the event bus. The agent sees it on its next
// turn; it appears in the chat log immediately.
InjectSystemEvent(content string)
// React records an emoji reaction to the most recent assistant response.
React(emoji string)
// ReactTo records an emoji reaction to a specific assistant response.
ReactTo(responseID, emoji string) error
// Reactions returns the current response ID to emoji mapping.
Reactions() map[string]string
}
// DetachedInfo describes a detached process for external consumers. // DetachedInfo describes a detached process for external consumers.
type DetachedInfo struct { type DetachedInfo struct {
PID int PID int

View File

@ -13,8 +13,8 @@ import (
) )
// newTestCore returns a minimal harness for testing. // newTestCore returns a minimal harness for testing.
func newTestCore(initialState string) *harness { func newTestCore(initialState string) *Session {
a := &harness{ a := &Session{
id: "test", id: "test",
state: initialState, state: initialState,
bus: pubsub.NewBus(), bus: pubsub.NewBus(),