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/manager.go

596 lines
16 KiB
Go

package session
import (
olog "ollie/log"
"context"
"fmt"
"ollie/paths"
agent "ollie/session"
"ollie/skills"
"ollie/tools"
"olliesrv/fs"
"os"
"strconv"
"strings"
"sync"
"sync/atomic"
)
var sessionFileOrder = []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 := fs.Perms[fs.PathSessions].Files[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
// NewCore, if non-nil, replaces the default backend.New + agent.New
// path. It receives the session ID, agent name, and cwd, and returns a Session.
NewCore func(sessionID, agentName, cwd string) (*agent.Session, error)
// Strict rejects inline code steps; only tool steps are allowed.
Strict bool
// Yolo skips the landrun sandbox.
Yolo bool
// NoMount disables the per-session FUSE mount (e.g. when listening on TCP).
NoMount bool
// Enable9P indicates the 9P listener is active (for operational model injection).
Enable9P bool
// EnableDBus indicates the D-Bus adapter is active (for operational model injection).
EnableDBus bool
// InvalidateModels clears the model cache, forcing a refresh.
InvalidateModels func()
// ResetElevation resets the per-turn elevation rate limiter for a session.
ResetElevation func(sessionID string)
// ToolRegistry is the shared tool registry for lazy tool promotion.
ToolRegistry *tools.Registry
// SkillsRegistry is the shared skills registry for skill loading.
SkillsRegistry *skills.Registry
// OnSessionCreated is called after a new session is added to the manager.
// Receives the session ID and the Session pointer.
OnSessionCreated func(id string, sess *Session)
// OnSessionKilled is called after a session is removed from the manager.
OnSessionKilled func(id string)
// OnSessionRenamed is called after a session is renamed.
OnSessionRenamed func(oldID, newID string)
}
// 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.restoreAllSessions()
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()
}
// List returns all root-level entries (fixed files + session directories).
func (s *Manager) List() ([]os.DirEntry, error) {
return s.list()
}
func (s *Manager) list() ([]os.DirEntry, error) {
entries := make([]os.DirEntry, 0, len(sessionFileOrder))
for _, name := range sessionFileOrder {
entries = append(entries, fs.FileEntry(name, fs.Perms[fs.PathSessions].Files[name]))
}
s.mu.RLock()
for id := range s.sessions {
entries = append(entries, fs.DirEntry(id, fs.Perms[fs.PathSessionDir].DirMode))
}
s.mu.RUnlock()
return entries, nil
}
// Readdir lists entries in a subdirectory.
func (s *Manager) Readdir(name string) ([]os.DirEntry, error) {
parts := strings.SplitN(name, "/", 4)
sessID := parts[0]
sess := s.Session(sessID)
if sess == nil {
return nil, fmt.Errorf("session not found: %s", sessID)
}
// {id} — list session-level files + agent/
if len(parts) == 1 {
entries := []os.DirEntry{
fs.FileEntry("plan", 0666),
fs.FileEntry("env", 0444),
fs.DirEntry("agent", 0755),
}
return entries, nil
}
// {id}/agent — list agent IDs
if parts[1] == "agent" {
if len(parts) == 2 {
// Currently single-agent: use the agent's ID
aid := sess.Core.Agent().ID()
if aid == "" {
aid = "0" // fallback for sessions without explicit agent ID
}
entries := []os.DirEntry{fs.DirEntry(aid, 0755)}
return entries, nil
}
// {id}/agent/{aid} — list agent files + proc/
if len(parts) == 3 {
afs, err := s.openAgentTree(sess)
if err != nil {
return nil, err
}
entries, err := afs.List()
if err != nil {
return nil, err
}
entries = append(entries, fs.DirEntry("proc", 0755))
return entries, nil
}
// {id}/agent/{aid}/proc — list detached process PIDs
if len(parts) == 4 && parts[3] == "proc" {
procs := sess.Core.ListDetached()
entries := make([]os.DirEntry, len(procs))
for i, p := range procs {
entries[i] = fs.FileEntry(fmt.Sprintf("%d", p.PID), 0666)
}
return entries, nil
}
}
// {id}/peer — list peer session IDs
if len(parts) == 2 && parts[1] == "peer" {
peers, _ := s.PeerList(sessID)
entries := make([]os.DirEntry, len(peers))
for i, p := range peers {
entries[i] = fs.FileEntry(p, 0666)
}
return entries, nil
}
return nil, fmt.Errorf("%s: not a directory", name)
}
// --- Peer Management ---
// PeerAdd creates a bidirectional peer link between two sessions.
func (s *Manager) PeerAdd(sessID, peerID string) error {
s.mu.RLock()
sess, ok := s.sessions[sessID]
peer, peerOk := s.sessions[peerID]
s.mu.RUnlock()
if !ok {
return fmt.Errorf("session not found: %s", sessID)
}
if !peerOk {
return fmt.Errorf("peer session not found: %s", peerID)
}
if sessID == peerID {
return fmt.Errorf("cannot peer a session with itself")
}
sess.mu.Lock()
if sess.peers == nil {
sess.peers = make(map[string]bool)
}
sess.peers[peerID] = true
sess.mu.Unlock()
peer.mu.Lock()
if peer.peers == nil {
peer.peers = make(map[string]bool)
}
peer.peers[sessID] = true
peer.mu.Unlock()
return nil
}
// PeerRemove removes a bidirectional peer link between two sessions.
func (s *Manager) PeerRemove(sessID, peerID string) error {
s.mu.RLock()
sess, ok := s.sessions[sessID]
peer, peerOk := s.sessions[peerID]
s.mu.RUnlock()
if !ok {
return fmt.Errorf("session not found: %s", sessID)
}
sess.mu.Lock()
delete(sess.peers, peerID)
sess.mu.Unlock()
if peerOk {
peer.mu.Lock()
delete(peer.peers, sessID)
peer.mu.Unlock()
}
return nil
}
// PeerList returns the peer session IDs for a session.
func (s *Manager) PeerList(sessID string) ([]string, error) {
s.mu.RLock()
sess, ok := s.sessions[sessID]
s.mu.RUnlock()
if !ok {
return nil, fmt.Errorf("session not found: %s", sessID)
}
sess.mu.RLock()
result := make([]string, 0, len(sess.peers))
for id := range sess.peers {
result = append(result, id)
}
sess.mu.RUnlock()
return result, nil
}
// PeerSubmit sends a prompt to a peer session.
func (s *Manager) PeerSubmit(sessID, peerID, prompt string) error {
s.mu.RLock()
sess, ok := s.sessions[sessID]
peer, peerOk := s.sessions[peerID]
s.mu.RUnlock()
if !ok {
return fmt.Errorf("session not found: %s", sessID)
}
if !peerOk {
return fmt.Errorf("peer session not found: %s", peerID)
}
sess.mu.RLock()
isPeer := sess.peers[peerID]
sess.mu.RUnlock()
if !isPeer {
return fmt.Errorf("%s is not a peer of %s", peerID, sessID)
}
go peer.Core.Submit(peer.Ctx, prompt)
return nil
}
func (s *Manager) stat(name string) (os.FileInfo, error) {
// Top-level fixed files (new, idx, sh, etc.)
if mode, ok := fs.Perms[fs.PathSessions].Files[name]; ok {
return &fs.SyntheticFileInfo{Name_: name, Mode_: mode}, nil
}
parts := strings.SplitN(name, "/", 4)
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_: fs.Perms[fs.PathSessionDir].DirMode, IsDir_: true}, nil
}
// Agent directory: {id}/agent
if parts[1] == "agent" {
if len(parts) == 2 {
return &fs.SyntheticFileInfo{Name_: "agent", Mode_: 0755, IsDir_: true}, nil
}
// Agent instance: {id}/agent/{aid}
if len(parts) == 3 {
return &fs.SyntheticFileInfo{Name_: parts[2], Mode_: 0755, IsDir_: true}, nil
}
// Agent file: {id}/agent/{aid}/{file}
if parts[3] == "proc" {
return &fs.SyntheticFileInfo{Name_: "proc", Mode_: 0755, IsDir_: true}, nil
}
afs, err := s.openAgentTree(sess)
if err != nil {
return nil, err
}
return afs.Stat(parts[3])
}
// Peer directory: {id}/peer
if parts[1] == "peer" {
if len(parts) == 2 {
return &fs.SyntheticFileInfo{Name_: "peer", Mode_: 0755, IsDir_: true}, nil
}
// Peer file: {id}/peer/{peer-id}
peerID := parts[2]
sess.mu.RLock()
isPeer := sess.peers[peerID]
sess.mu.RUnlock()
if !isPeer {
return nil, fmt.Errorf("%s: not found", name)
}
return &fs.SyntheticFileInfo{Name_: peerID, Mode_: 0666}, nil
}
// Session file: {id}/{file} (plan, env)
sfs, err := s.openSessionTree(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")
}
// Peer file: {id}/peer/{peer-id} — write submits prompt to peer
parts := strings.SplitN(name, "/", 3)
// Proc file: {id}/proc/{pid} — read returns ring buffer output
if len(parts) == 3 && parts[1] == "proc" {
sessID := parts[0]
sess := s.Session(sessID)
if sess == nil {
return nil, fmt.Errorf("%s: not found", name)
}
pidStr := parts[2]
pid, err := strconv.Atoi(pidStr)
if err != nil {
return nil, fmt.Errorf("invalid pid: %s", pidStr)
}
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) {
return &fs.SyntheticFileInfo{Name_: pidStr, Mode_: 0444}, nil
},
ReadFn: func() ([]byte, error) {
output, err := sess.Core.GetDetachedOutput(pid)
if err != nil {
return nil, err
}
return []byte(output), nil
},
WriteFn: func([]byte) error { return fmt.Errorf("read-only") },
BlockingReadFn: notBlocking,
}, nil
}
if len(parts) == 3 && parts[1] == "peer" {
sessID := parts[0]
peerID := parts[2]
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) {
return &fs.SyntheticFileInfo{Name_: peerID, Mode_: 0666}, nil
},
ReadFn: func() ([]byte, error) {
return nil, nil // empty read
},
WriteFn: func(data []byte) error {
prompt := strings.TrimSpace(string(data))
if prompt == "" {
return nil
}
return s.PeerSubmit(sessID, peerID, prompt)
},
BlockingReadFn: notBlocking,
}, nil
}
// Top-level fixed files.
switch name {
case "new":
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) {
return &fs.SyntheticFileInfo{Name_: "new", Mode_: fs.Perms[fs.PathSessions].Files["new"]}, nil
},
ReadFn: func() ([]byte, error) {
return []byte("name=\ncwd=\nremote=\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 {
_, err := s.CreateSession(strings.Fields(strings.TrimSpace(string(data))))
return err
},
BlockingReadFn: notBlocking,
}, nil
case "idx":
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) {
return &fs.SyntheticFileInfo{Name_: "idx", Mode_: fs.Perms[fs.PathSessions].Files["idx"]}, nil
},
ReadFn: func() ([]byte, error) { return s.index(), nil },
WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") },
BlockingReadFn: notBlocking,
}, nil
}
if _, ok := fs.Perms[fs.PathSessions].Files[name]; ok {
return &fs.FileConfig{
StatFn: func() (os.FileInfo, error) {
return &fs.SyntheticFileInfo{Name_: name, Mode_: fs.Perms[fs.PathSessions].Files[name]}, nil
},
ReadFn: func() ([]byte, error) {
return s.cfg.ReadFile(paths.CfgDir() + "/scripts/session/" + name)
},
WriteFn: func([]byte) error { return fmt.Errorf("%s: not writable", name) },
BlockingReadFn: notBlocking,
}, nil
}
// Hierarchical paths: {id}/{file} or {id}/agent/{aid}/{file}
parts = strings.SplitN(name, "/", 4)
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)
}
// Agent file: {id}/agent/{aid}/{file}
if parts[1] == "agent" {
if len(parts) < 4 {
return nil, fmt.Errorf("%s: is a directory", name)
}
// parts[2] = aid, parts[3] = file
afs, err := s.openAgentTree(sess)
if err != nil {
return nil, err
}
return afs.Open(parts[3])
}
// Session file: {id}/{file}
sfs, err := s.openSessionTree(sess)
if err != nil {
return nil, err
}
return sfs.Open(parts[1])
}
func (s *Manager) create(name string) error {
// Peer creation: {id}/peer/{peer-id}
parts := strings.SplitN(name, "/", 3)
if len(parts) == 3 && parts[1] == "peer" {
return s.PeerAdd(parts[0], parts[2])
}
return fmt.Errorf("create 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)
}
// Peer removal: {id}/peer/{peer-id}
if len(parts) == 3 && parts[1] == "peer" {
return s.PeerRemove(sessID, parts[2])
}
// Proc dismiss: {id}/proc/{pid}
if len(parts) == 3 && parts[1] == "proc" {
sess := s.Session(sessID)
if sess == nil {
return fmt.Errorf("session not found: %s", sessID)
}
pid, err := strconv.Atoi(parts[2])
if err != nil {
return fmt.Errorf("invalid pid: %s", parts[2])
}
if !sess.Core.DismissDetached(pid) {
return fmt.Errorf("process %d not found or still running", pid)
}
return nil
}
// 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
}
// OpenSessionTree returns a Tree for the given session ID.
func (s *Manager) OpenSessionTree(id string) (*fs.Tree, error) {
if sess := s.Session(id); sess != nil {
return s.openSessionTree(sess)
}
return nil, fmt.Errorf("session not found: %s", id)
}
func (s *Manager) openSessionTree(sess *Session) (*fs.Tree, error) {
var resetElev func()
if s.cfg.ResetElevation != nil {
id := sess.id
resetElev = func() { s.cfg.ResetElevation(id) }
}
return NewSessionTree(
sess,
s.cfg.Log,
func() { s.KillSession(sess.id) },
func(newID string) error { return s.renameSession(sess.id, newID) },
nil, // no transcript saving
s.cfg.InvalidateModels,
resetElev,
s.cfg.ToolRegistry,
), nil
}
func (s *Manager) openAgentTree(sess *Session) (*fs.Tree, error) {
var resetElev func()
if s.cfg.ResetElevation != nil {
id := sess.id
resetElev = func() { s.cfg.ResetElevation(id) }
}
return NewAgentTree(
sess,
s.cfg.Log,
nil, // no transcript saving
s.cfg.InvalidateModels,
resetElev,
s.cfg.ToolRegistry,
), 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.