This repository has been archived on 2026-08-16. You can view files and clone it, but cannot push or open issues or pull requests.
ollie-9p/session/session.go

772 lines
20 KiB
Go

package session
import (
"context"
"fmt"
"os"
"sort"
"strings"
"sync"
"sync/atomic"
"ollie/pkg/agent"
"ollie/pkg/backend"
"ollie/pkg/config"
olog "ollie/pkg/log"
"ollie/pkg/paths"
"ollie/pkg/tools"
"ollie/pkg/tools/execute"
"olliesrv/fs"
)
// Session holds all state for one agent session.
type Session struct {
mu sync.RWMutex
id string
uname string // immutable user principal (numeric UID), set at creation
Core agent.Core
Ctx context.Context
cancel context.CancelFunc
log []byte
logVers uint32
ChatOffset int
plan []byte
prevPrompt []byte // last submitted prompt; overwritten on each new submission
AllowTools []string // if non-empty, only these tools are visible/executable
}
func NewSession(id string, core agent.Core, ctx context.Context, cancel context.CancelFunc) *Session {
sess := &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
sess.startEventLog()
return sess
}
func (sess *Session) RunnableID() string { return sess.id }
func (sess *Session) Uname() string { return sess.uname }
func (sess *Session) Cancel() {
sess.cancel()
}
func (sess *Session) Interrupt() {
sess.Core.Interrupt(agent.ErrInterrupted)
}
// AppendLog appends data to the session's log and bumps the version.
func (sess *Session) AppendLog(data []byte) {
if len(data) == 0 {
return
}
sess.mu.Lock()
sess.log = append(sess.log, data...)
sess.logVers++
sess.mu.Unlock()
}
// EnsureTrailingNewline appends a newline if the log doesn't already end with one.
func (sess *Session) EnsureTrailingNewline() {
sess.mu.Lock()
if len(sess.log) > 0 && sess.log[len(sess.log)-1] != '\n' {
sess.log = append(sess.log, '\n')
}
sess.logVers++
sess.mu.Unlock()
}
// startEventLog subscribes to the agent's "event" bus topic and writes
// all events to the session chat log.
func (sess *Session) startEventLog() {
assistantStarted := false
sess.Core.Bus().Subscribe("event", func(ev agent.Event) {
if ev.Role == "user" {
if assistantStarted {
sess.AppendLog([]byte("\n"))
assistantStarted = false
}
sess.AppendLog(FormatEvent(ev))
return
}
switch ev.Role {
case "assistant":
if !assistantStarted {
sess.AppendLog([]byte("assistant: "))
sess.mu.Lock()
sess.ChatOffset = len(sess.log)
sess.mu.Unlock()
assistantStarted = true
}
default:
if assistantStarted {
sess.AppendLog([]byte("\n"))
assistantStarted = false
}
}
sess.AppendLog(FormatEvent(ev))
})
}
// LogInfo returns the current log length and version atomically.
func (sess *Session) LogInfo() (length int, vers uint32) {
sess.mu.RLock()
defer sess.mu.RUnlock()
return len(sess.log), sess.logVers
}
// sessionStoreFiles maps fixed file entries in s/ to their permissions.
var sessionStoreFiles = map[string]os.FileMode{
"new": 0666,
"idx": 0444,
"ls": 0555,
"kill": 0555,
"sh": 0550,
"b": 0555,
"bfg": 0555,
"bbg": 0555,
"cleanup": 0555,
}
var sessionStoreOrder = []string{"new", "idx", "ls", "kill", "sh", "b", "bfg", "bbg", "cleanup"}
// FileMode returns the mode for a fixed session file,
// or 0 and false if the name is not a fixed file.
func FileMode(name string) (os.FileMode, bool) {
m, ok := sessionStoreFiles[name]
return m, ok
}
// ManagerConfig holds the dependencies for a Manager.
type ManagerConfig struct {
AgentsDir string
SessionsDir string
Log *olog.Logger
Sink *olog.Sink
ReadFile func(string) ([]byte, error)
MkdirAll func(string, os.FileMode) error
SaveTranscript func([]byte) error
// ToolTree provides access to the shared tool directory.
ToolTree fs.FileTree
// NewCore, if non-nil, replaces the default backend.New + agent.NewAgentCore
// path. It receives the session ID, agent name, and cwd, and returns a Core.
NewCore func(sessionID, agentName, cwd string) (agent.Core, error)
// Strict rejects inline code steps; only tool steps are allowed.
Strict bool
// Yolo skips the landrun sandbox.
Yolo bool
}
// Manager manages session lifecycle and exposes sessions as a Tree.
type Manager struct {
tree *fs.Tree
cfg ManagerConfig
mu sync.RWMutex
sessions map[string]*Session
nextUID atomic.Uint32 // incrementing principal counter
}
// Tree returns the Tree view of the session namespace.
func (s *Manager) Tree() *fs.Tree { return s.tree }
// nextUname generates the next uname atomically.
func (s *Manager) nextUname() string {
return fmt.Sprintf("%d", s.nextUID.Add(1))
}
func NewManager(cfg ManagerConfig) *Manager {
if cfg.ReadFile == nil {
cfg.ReadFile = os.ReadFile
}
if cfg.MkdirAll == nil {
cfg.MkdirAll = os.MkdirAll
}
ss := &Manager{
cfg: cfg,
sessions: make(map[string]*Session),
}
ss.nextUID.Store(9999)
ss.tree = fs.NewTree(nil, 0,
fs.WithStat(func(_ []string, name string) (os.FileInfo, error) { return ss.stat(name) }),
fs.WithOpener(func(_ []string, name string) (fs.File, error) { return ss.openEntry(name) }),
fs.WithLister(func(_ []string) ([]os.DirEntry, error) { return ss.list() }),
fs.WithReaddir(func(_ []string, name string) ([]os.DirEntry, error) { return ss.Readdir(name) }),
fs.WithDeleter(func(_ []string, name string) error { return ss.del(name) }),
fs.WithCreator(func(_ []string, name string, _ os.FileMode) error { return ss.create(name) }),
fs.WithRenamer(func(_ []string, old, new string) error { return ss.renameSession(old, new) }),
)
return ss
}
// AddSession inserts a pre-built session into the fs.
func (s *Manager) AddSession(sess *Session) {
s.mu.Lock()
s.sessions[sess.RunnableID()] = sess
s.mu.Unlock()
}
func (s *Manager) list() ([]os.DirEntry, error) {
entries := make([]os.DirEntry, 0, len(sessionStoreOrder))
for _, name := range sessionStoreOrder {
entries = append(entries, fs.FileEntry(name, sessionStoreFiles[name]))
}
s.mu.RLock()
for id := range s.sessions {
entries = append(entries, fs.DirEntry(id, 0555))
}
s.mu.RUnlock()
return entries, nil
}
// Readdir lists entries in a subdirectory (e.g. "{id}", "{id}/t", "{id}/t/{subdir}").
func (s *Manager) Readdir(name string) ([]os.DirEntry, error) {
parts := strings.SplitN(name, "/", 3)
sessID := parts[0]
sess := s.Session(sessID)
if sess == nil {
return nil, fmt.Errorf("session not found: %s", sessID)
}
// {id} — list session files + t/
if len(parts) == 1 {
sfs, err := s.openStore(sess)
if err != nil {
return nil, err
}
entries, err := sfs.List()
if err != nil {
return nil, err
}
entries = append(entries, fs.DirEntry("t", 0500))
return entries, nil
}
// {id}/t — list tools (filtered by allowTools)
if parts[1] == "t" && s.cfg.ToolTree != nil {
if len(parts) == 2 {
return s.cfg.ToolTree.List()
}
// {id}/t/{subdir} — delegate to ToolTree.Readdir if available
type listDirer interface {
Readdir(string) ([]os.DirEntry, error)
}
if ld, ok := s.cfg.ToolTree.(listDirer); ok {
return ld.Readdir(parts[2])
}
}
return nil, fmt.Errorf("%s: not a directory", name)
}
func (s *Manager) stat(name string) (os.FileInfo, error) {
// Top-level fixed files (new, idx, sh, etc.)
if mode, ok := sessionStoreFiles[name]; ok {
return &fs.SyntheticFileInfo{Name_: name, Mode_: mode}, nil
}
parts := strings.SplitN(name, "/", 3)
sessID := parts[0]
s.mu.RLock()
sess, ok := s.sessions[sessID]
s.mu.RUnlock()
if !ok {
return nil, fmt.Errorf("%s: not found", name)
}
// Session directory: {id}
if len(parts) == 1 {
return &fs.SyntheticFileInfo{Name_: sessID, Mode_: 0555, IsDir_: true}, nil
}
// Tools directory: {id}/t
if parts[1] == "t" {
if len(parts) == 2 {
return &fs.SyntheticFileInfo{Name_: "t", Mode_: 0500, IsDir_: true}, nil
}
// Tool file: {id}/t/{rel}
if s.cfg.ToolTree != nil {
return s.cfg.ToolTree.Stat(parts[2])
}
return nil, fmt.Errorf("%s: not found", name)
}
// Session file: {id}/{file}
sfs, err := s.openStore(sess)
if err != nil {
return nil, err
}
return sfs.Stat(parts[1])
}
func (s *Manager) openEntry(name string) (fs.File, error) {
notBlocking := func(context.Context, string) ([]byte, string, error) {
return nil, "", fmt.Errorf("blocking read not supported")
}
// Top-level fixed files.
switch name {
case "new":
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) { return &fs.SyntheticFileInfo{Name_: "new", Mode_: 0666}, nil },
ReadFn: func() ([]byte, error) {
return []byte("name=\ncwd=\nbackend=\nmodel=\nagent=\nmaxTokens=\nmaxCompletionTokens=\ntemperature=\ntopP=\ntopK=\nminP=\ntopA=\nfrequencyPenalty=\npresencePenalty=\nrepetitionPenalty=\nreasoning=\nreasoningEffort=\nincludeReasoning=\nresponseFormat=\nstop=\nverbosity=\n"), nil
},
WriteFn: func(data []byte) error {
return s.createSession(strings.Fields(strings.TrimSpace(string(data))))
},
BlockingReadFn: notBlocking,
}, nil
case "idx":
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) { return &fs.SyntheticFileInfo{Name_: "idx", Mode_: 0444}, nil },
ReadFn: func() ([]byte, error) { return s.index(), nil },
WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") },
BlockingReadFn: notBlocking,
}, nil
}
if _, ok := sessionStoreFiles[name]; ok {
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) { return &fs.SyntheticFileInfo{Name_: name, Mode_: sessionStoreFiles[name]}, nil },
ReadFn: func() ([]byte, error) {
return s.cfg.ReadFile(paths.CfgDir() + "/scripts/s/" + name)
},
WriteFn: func([]byte) error { return fmt.Errorf("%s: not writable", name) },
BlockingReadFn: notBlocking,
}, nil
}
// Hierarchical paths: {id}/{file} or {id}/t/{tool}
parts := strings.SplitN(name, "/", 3)
sessID := parts[0]
sess := s.Session(sessID)
if sess == nil {
return nil, fmt.Errorf("%s: not found", name)
}
if len(parts) == 1 {
return nil, fmt.Errorf("%s: is a directory", name)
}
// Tool file: {id}/t/{rel}
if parts[1] == "t" {
if len(parts) < 3 {
return nil, fmt.Errorf("%s: is a directory", name)
}
rel := parts[2]
if s.cfg.ToolTree == nil {
return nil, fmt.Errorf("%s: no tool tree", name)
}
// Apply per-session allowTools filtering on idx.
if rel == "idx" && len(sess.AllowTools) > 0 {
allowed := make(map[string]bool, len(sess.AllowTools))
for _, t := range sess.AllowTools {
allowed[t] = true
}
return s.openFilteredToolIdx(allowed)
}
return s.cfg.ToolTree.Open(rel)
}
// Session file: {id}/{file}
sfs, err := s.openStore(sess)
if err != nil {
return nil, err
}
return sfs.Open(parts[1])
}
func (s *Manager) openFilteredToolIdx(allowed map[string]bool) (fs.File, error) {
base, err := s.cfg.ToolTree.Open("idx")
if err != nil {
return nil, err
}
return &fs.FileConfig{
StatFn: base.Stat,
ReadFn: func() ([]byte, error) {
data, err := base.Read()
if err != nil {
return nil, err
}
return filterToolIdx(data, allowed), nil
},
WriteFn: base.Write,
BlockingReadFn: base.BlockingRead,
}, nil
}
func filterToolIdx(data []byte, allowed map[string]bool) []byte {
var out []byte
for _, section := range strings.Split(string(data), "## ") {
if section == "" {
continue
}
name, _, _ := strings.Cut(section, "\n")
if allowed[name] {
out = append(out, "## "...)
out = append(out, section...)
}
}
return out
}
func (s *Manager) create(name string) error {
// Tool create: {id}/t/{rel}
parts := strings.SplitN(name, "/", 3)
if len(parts) == 3 && parts[1] == "t" && s.cfg.ToolTree != nil {
return s.cfg.ToolTree.Create(parts[2])
}
return fmt.Errorf("create not supported: %s", name)
}
// MkdirAll creates a directory within the session's tool fs.
func (s *Manager) MkdirAll(name string) error {
parts := strings.SplitN(name, "/", 3)
if len(parts) == 3 && parts[1] == "t" {
type mkdirAller interface {
MkdirAll(string) error
}
if md, ok := s.cfg.ToolTree.(mkdirAller); ok {
return md.MkdirAll(parts[2])
}
}
return fmt.Errorf("mkdir not supported: %s", name)
}
func (s *Manager) del(name string) error {
parts := strings.SplitN(name, "/", 3)
sessID := parts[0]
// Delete session directory itself.
if len(parts) == 1 {
s.mu.RLock()
_, ok := s.sessions[sessID]
s.mu.RUnlock()
if ok {
s.KillSession(sessID)
return nil
}
return fmt.Errorf("session not found: %s", sessID)
}
// Delete tool file: {id}/t/{rel}
if parts[1] == "t" && len(parts) == 3 {
if s.cfg.ToolTree != nil {
return s.cfg.ToolTree.Delete(parts[2])
}
return fmt.Errorf("no tool tree")
}
// Session files are synthetic; allow rm -r to continue.
return nil
}
// Session returns the session for the given ID, or nil.
func (s *Manager) Session(id string) *Session {
s.mu.RLock()
defer s.mu.RUnlock()
return s.sessions[id]
}
// SessionByUname returns the session with the given uname (principal), or nil.
func (s *Manager) SessionByUname(uname string) *Session {
s.mu.RLock()
defer s.mu.RUnlock()
for _, sess := range s.sessions {
if sess.uname == uname {
return sess
}
}
return nil
}
// OpenStore returns a Tree for the given session ID.
func (s *Manager) OpenStore(id string) (*fs.Tree, error) {
if sess := s.Session(id); sess != nil {
return s.openStore(sess)
}
return nil, fmt.Errorf("session not found: %s", id)
}
func (s *Manager) openStore(sess *Session) (*fs.Tree, error) {
return NewSessionTree(
sess,
s.cfg.Log,
func() { s.KillSession(sess.id) },
func(newID string) error { return s.renameSession(sess.id, newID) },
s.cfg.SaveTranscript,
), nil
}
// InterruptAll interrupts every active session.
func (s *Manager) InterruptAll() {
s.mu.RLock()
defer s.mu.RUnlock()
for _, sess := range s.sessions {
sess.Core.Interrupt(agent.ErrInterrupted)
}
}
// Shutdown kills all active sessions.
func (s *Manager) Shutdown() {
s.mu.Lock()
ids := make([]string, 0, len(s.sessions))
for id := range s.sessions {
ids = append(ids, id)
}
s.mu.Unlock()
for _, id := range ids {
s.KillSession(id)
}
}
func (s *Manager) KillSession(id string) {
s.mu.Lock()
sess := s.sessions[id]
delete(s.sessions, id)
s.mu.Unlock()
if sess != nil {
sess.Cancel()
sess.Core.Close()
s.cfg.Log.Info("killed session %s", id)
}
}
func (s *Manager) index() []byte {
var sb strings.Builder
s.mu.RLock()
ids := make([]string, 0, len(s.sessions))
for id := range s.sessions {
ids = append(ids, id)
}
sort.Strings(ids)
for _, id := range ids {
sess := s.sessions[id]
sess.mu.RLock()
state := sess.Core.State()
cwd := sess.Core.CWD()
be := sess.Core.BackendName()
model := sess.Core.ModelName()
sess.mu.RUnlock()
fmt.Fprintf(&sb, "%s\t%s\t%s\t%s\t%s\n", id, state, cwd, be, model)
}
s.mu.RUnlock()
return []byte(sb.String())
}
func (s *Manager) createSession(args []string) error {
name := ""
backendOverride := ""
modelOverride := ""
agentName := os.Getenv("OLLIE_DEFAULT_AGENT")
if agentName == "" {
agentName = "default"
}
cwd := ""
for _, arg := range args {
k, v, ok := strings.Cut(arg, "=")
if !ok {
return fmt.Errorf("invalid option %q (expected key=value)", arg)
}
if v == "" {
continue
}
switch k {
case "name":
name = v
case "backend":
backendOverride = v
case "model":
modelOverride = v
case "agent":
agentName = v
case "cwd":
cwd = v
default:
return fmt.Errorf("unknown option %q (valid: name, backend, model, agent, cwd)", k)
}
}
cwd = paths.ExpandHome(os.ExpandEnv(cwd))
if cwd == "" {
return fmt.Errorf("cwd is required (e.g. new cwd=/path/to/project)")
}
sessID := name
if sessID == "" {
sessID = agent.NewSessionID()
}
s.mu.RLock()
_, exists := s.sessions[sessID]
s.mu.RUnlock()
if exists {
return fmt.Errorf("session already exists: %s", sessID)
}
var core agent.Core
var allowTools []string
uname := s.nextUname()
if s.cfg.NewCore != nil {
var err error
core, err = s.cfg.NewCore(sessID, agentName, cwd)
if err != nil {
return err
}
} else {
cfg := LoadAgentConfig(s.cfg.AgentsDir, agentName, nil)
if cfg != nil {
allowTools = cfg.AllowTools
if backendOverride == "" && cfg.Backend != "" {
backendOverride = cfg.Backend
}
if modelOverride == "" && cfg.Model != "" {
modelOverride = cfg.Model
}
}
be, err := backend.NewWithName(backendOverride)
if err != nil {
return fmt.Errorf("backend: %w", err)
}
if modelOverride == "" {
modelOverride = os.Getenv("OLLIE_MODEL")
}
if modelOverride != "" {
be.SetModel(modelOverride)
}
if err := s.cfg.MkdirAll(s.cfg.SessionsDir, 0700); err != nil {
return fmt.Errorf("sessions dir: %w", err)
}
var execOpts []execute.Option
execOpts = append(execOpts, execute.WithMount())
if s.cfg.Strict {
execOpts = append(execOpts, execute.WithStrict())
}
if s.cfg.Yolo {
execOpts = append(execOpts, execute.WithYolo())
}
if cfg != nil && len(cfg.AllowExecutors) > 0 {
execOpts = append(execOpts, execute.WithAllowExecutors(cfg.AllowExecutors))
}
if cfg != nil && len(cfg.AllowTools) > 0 {
execOpts = append(execOpts, execute.WithAllowTools(cfg.AllowTools))
}
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
"execute": execute.Decl(cwd, execOpts...),
})
env := agent.BuildAgentEnv(cfg, newDisp(), cwd, []string{"OLLIE_SESSION_ID=" + sessID, "OLLIE_UNAME=" + uname})
core = agent.NewAgentCore(agent.AgentCoreConfig{
Backend: be,
AgentName: agentName,
AgentsDir: s.cfg.AgentsDir,
SessionsDir: s.cfg.SessionsDir,
SessionID: sessID,
Uname: uname,
CWD: cwd,
Env: env,
NewDispatcher: newDisp,
Log: s.cfg.Sink.NewLogger("core"),
})
}
ctx, cancel := context.WithCancel(context.Background())
sess := NewSession(sessID, core, ctx, cancel)
sess.AllowTools = allowTools
s.mu.Lock()
sess.uname = uname
s.sessions[sessID] = sess
s.mu.Unlock()
s.cfg.Log.Info("new session %s (backend=%s model=%s agent=%s)",
sessID, core.BackendName(), core.ModelName(), core.AgentName())
return nil
}
func (s *Manager) renameSession(old, new string) error {
// Tool rename: {id}/t/{oldRel} -> newName
if strings.Contains(old, "/t/") {
parts := strings.SplitN(old, "/", 3)
if len(parts) == 3 && parts[1] == "t" && s.cfg.ToolTree != nil {
oldRel := parts[2]
// new is just the new base name; reconstruct full rel path.
oldBase := oldRel
if i := strings.LastIndex(oldRel, "/"); i >= 0 {
oldBase = oldRel[i+1:]
}
parentRel := oldRel[:len(oldRel)-len(oldBase)]
newRel := parentRel + new
return s.cfg.ToolTree.Rename(oldRel, newRel)
}
return fmt.Errorf("rename not supported: %s", old)
}
// Session rename: {oldID} -> {newID}
oldID := old
newID := new
s.mu.Lock()
sess, ok := s.sessions[oldID]
if !ok {
s.mu.Unlock()
return fmt.Errorf("session not found: %s", oldID)
}
if _, exists := s.sessions[newID]; exists {
s.mu.Unlock()
return fmt.Errorf("session already exists: %s", newID)
}
if sess.Core.IsRunning() {
s.mu.Unlock()
return fmt.Errorf("cannot rename while agent is running")
}
if err := sess.Core.SetSessionID(newID); err != nil {
s.mu.Unlock()
return err
}
sess.id = newID
s.sessions[newID] = sess
delete(s.sessions, oldID)
s.mu.Unlock()
sess.AppendLog([]byte(fmt.Sprintf(":: session renamed: %s -> %s\n", oldID, newID)))
s.cfg.Log.Info("renamed session %s -> %s", oldID, newID)
return nil
}
// LoadAgentConfig resolves and loads the config for a named agent.
// Returns nil if the config file does not exist; BuildAgentEnv handles nil configs.
func LoadAgentConfig(agentsDir, name string, open func(string) (*os.File, error)) *config.Config {
if open == nil {
open = os.Open
}
path := agent.AgentConfigPath(agentsDir, name)
f, err := open(path)
if err != nil {
return nil
}
defer f.Close()
cfg, _ := config.Load(f)
return cfg
}
// FormatEvent converts an agent Event to bytes for appending to a chat log.
func FormatEvent(ev agent.Event) []byte {
switch ev.Role {
case "user":
return []byte("user: " + ev.Content + "\n")
case "assistant":
return []byte(ev.Content)
case "call":
args := squashWhitespace(ev.Content)
return []byte("-> " + ev.Name + "(" + args + ")\n")
case "tool":
return []byte(strings.TrimRight(ev.Content, "\n") + "\n")
case "retry":
return []byte("retrying in " + ev.Content + "s...\n")
case "error":
return []byte("error: " + ev.Content + "\n")
case "stalled":
return []byte("agent stalled\n")
case "reasoning":
return []byte(ev.Content)
case "exec":
return []byte("```\n" + ev.Content + "\n```\n")
case "info":
return []byte(":: " + ev.Content)
default:
return nil
}
}
func squashWhitespace(s string) string {
return strings.Join(strings.Fields(s), " ")
}