226 lines
4.9 KiB
Go
226 lines
4.9 KiB
Go
// Package session defines the runtime environment for an agent.
|
|
// A Session holds identity, working directory, environment, state machine,
|
|
// event bus, and prompt queue — everything that persists across agent swaps.
|
|
package session
|
|
|
|
import (
|
|
"sync"
|
|
|
|
"github.com/simonfxr/pubsub"
|
|
)
|
|
|
|
// Session is the runtime environment in which an agent operates.
|
|
// It outlives any particular agent configuration and maintains
|
|
// identity, state, and communication channels.
|
|
type Session struct {
|
|
mu sync.RWMutex
|
|
id string
|
|
uname string // immutable user principal
|
|
cwd string
|
|
state string // "idle", "thinking", "calling: <tool>"
|
|
reply string // assistant text from last completed turn
|
|
|
|
bus *pubsub.Bus
|
|
fifo Fifo
|
|
|
|
envMu sync.RWMutex
|
|
env map[string]string
|
|
|
|
plan []byte
|
|
prevPrompt string
|
|
peers map[string]bool // peer session IDs (bidirectional)
|
|
|
|
changeMu sync.Mutex
|
|
changeCond *sync.Cond
|
|
}
|
|
|
|
// Config holds parameters for creating a new Session.
|
|
type Config struct {
|
|
ID string
|
|
Uname string
|
|
CWD string
|
|
}
|
|
|
|
// New creates a Session with the given configuration.
|
|
func New(cfg Config) *Session {
|
|
s := &Session{
|
|
id: cfg.ID,
|
|
uname: cfg.Uname,
|
|
cwd: cfg.CWD,
|
|
state: "idle",
|
|
bus: pubsub.NewBus(),
|
|
env: make(map[string]string),
|
|
}
|
|
s.changeCond = sync.NewCond(&s.changeMu)
|
|
return s
|
|
}
|
|
|
|
// ID returns the session identifier.
|
|
func (s *Session) ID() string { return s.id }
|
|
|
|
// Uname returns the immutable user principal.
|
|
func (s *Session) Uname() string { return s.uname }
|
|
|
|
// CWD returns the current working directory.
|
|
func (s *Session) CWD() string {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.cwd
|
|
}
|
|
|
|
// SetCWD updates the working directory.
|
|
func (s *Session) SetCWD(dir string) {
|
|
s.mu.Lock()
|
|
s.cwd = dir
|
|
s.mu.Unlock()
|
|
s.notifyChange()
|
|
}
|
|
|
|
// State returns the current session state.
|
|
func (s *Session) State() string {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.state
|
|
}
|
|
|
|
// SetState transitions to a new state and notifies waiters.
|
|
func (s *Session) SetState(state string) {
|
|
s.mu.Lock()
|
|
s.state = state
|
|
s.mu.Unlock()
|
|
s.notifyChange()
|
|
}
|
|
|
|
// Reply returns the assistant text from the last completed turn.
|
|
func (s *Session) Reply() string {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.reply
|
|
}
|
|
|
|
// SetReply stores the assistant reply text.
|
|
func (s *Session) SetReply(text string) {
|
|
s.mu.Lock()
|
|
s.reply = text
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// Bus returns the session event bus.
|
|
func (s *Session) Bus() *pubsub.Bus { return s.bus }
|
|
|
|
// SetEnv sets a session-scoped environment variable.
|
|
func (s *Session) SetEnv(key, value string) {
|
|
s.envMu.Lock()
|
|
s.env[key] = value
|
|
s.envMu.Unlock()
|
|
}
|
|
|
|
// Env returns a copy of all session environment variables.
|
|
func (s *Session) Env() map[string]string {
|
|
s.envMu.RLock()
|
|
defer s.envMu.RUnlock()
|
|
out := make(map[string]string, len(s.env))
|
|
for k, v := range s.env {
|
|
out[k] = v
|
|
}
|
|
return out
|
|
}
|
|
|
|
// SetID renames the session.
|
|
func (s *Session) SetID(id string) {
|
|
s.mu.Lock()
|
|
s.id = id
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// Queue pushes a prompt onto the FIFO.
|
|
func (s *Session) Queue(prompt string) { s.fifo.Push(prompt) }
|
|
|
|
// PopQueue removes and returns the next queued prompt.
|
|
func (s *Session) PopQueue() (string, bool) { return s.fifo.Pop() }
|
|
|
|
// WaitChange blocks until the state changes from current, then returns the new state.
|
|
// Returns ("", false) if ctx-based cancellation would be needed (not implemented here).
|
|
func (s *Session) WaitChange(current string) string {
|
|
s.changeMu.Lock()
|
|
defer s.changeMu.Unlock()
|
|
for {
|
|
s.mu.RLock()
|
|
now := s.state
|
|
s.mu.RUnlock()
|
|
if now != current {
|
|
return now
|
|
}
|
|
s.changeCond.Wait()
|
|
}
|
|
}
|
|
|
|
func (s *Session) notifyChange() {
|
|
s.changeMu.Lock()
|
|
s.changeCond.Broadcast()
|
|
s.changeMu.Unlock()
|
|
}
|
|
|
|
// Plan returns the session plan.
|
|
func (s *Session) Plan() []byte {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.plan
|
|
}
|
|
|
|
// SetPlan updates the session plan.
|
|
func (s *Session) SetPlan(p []byte) {
|
|
s.mu.Lock()
|
|
s.plan = p
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// PrevPrompt returns the last submitted prompt.
|
|
func (s *Session) PrevPrompt() string {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.prevPrompt
|
|
}
|
|
|
|
// SetPrevPrompt stores the last submitted prompt.
|
|
func (s *Session) SetPrevPrompt(p string) {
|
|
s.mu.Lock()
|
|
s.prevPrompt = p
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// HasPeer reports whether peerID is linked to this session.
|
|
func (s *Session) HasPeer(id string) bool {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.peers[id]
|
|
}
|
|
|
|
// AddPeer adds a bidirectional peer link.
|
|
func (s *Session) AddPeer(id string) {
|
|
s.mu.Lock()
|
|
if s.peers == nil {
|
|
s.peers = make(map[string]bool)
|
|
}
|
|
s.peers[id] = true
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// RemovePeer removes a peer link.
|
|
func (s *Session) RemovePeer(id string) {
|
|
s.mu.Lock()
|
|
delete(s.peers, id)
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
// Peers returns a copy of all peer IDs.
|
|
func (s *Session) Peers() []string {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
out := make([]string, 0, len(s.peers))
|
|
for id := range s.peers {
|
|
out = append(out, id)
|
|
}
|
|
return out
|
|
}
|