session: support multiple agents in data structures
- session.Session: r *agent.Agent → agents []*agent.Agent - Agent() returns agents[0] for backward compat - Added Agents(), AgentAt(), AgentCount(), AddAgent() - Close/SetEnv/SetSessionID iterate all agents - Other methods use Agent() (primary agent) - fs/session: AgentLog *AgentLog → AgentLogs map[string]*AgentLog - Added AgentLog() (returns first) and SetAgentLog(id, al) methods - Updated create.go, persist.go, root.go, olliesrv/server.go
This commit is contained in:
parent
107a157a18
commit
951fdc8133
|
|
@ -1373,8 +1373,8 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
if len(parts) == 5 && parts[0] == "session" && parts[2] == "agent" {
|
||||
if sess := fs.Lookup(s.sessionTree, parts[1]); sess != nil {
|
||||
if base == "log" || base == "chat" {
|
||||
if sess.AgentLog != nil {
|
||||
length, vers := sess.AgentLog.LogInfo()
|
||||
if al := sess.AgentLog(); al != nil {
|
||||
length, vers := al.LogInfo()
|
||||
if length > 64*1024 {
|
||||
length = 64 * 1024
|
||||
}
|
||||
|
|
|
|||
|
|
@ -269,7 +269,7 @@ func CreateAgent(rs *rootState, sessName string, args []string) error {
|
|||
sess.mu.Lock()
|
||||
sess.Core = core
|
||||
sess.uname = uname
|
||||
sess.AgentLog = al
|
||||
sess.SetAgentLog(uname, al)
|
||||
sess.proc = proc
|
||||
sess.mu.Unlock()
|
||||
|
||||
|
|
|
|||
|
|
@ -332,26 +332,27 @@ func restoreSession(rs *rootState, ps *agent.PersistedAgent) error {
|
|||
BaseLayers: []string{sysPrompt, opModel, envBlock},
|
||||
Log: rs.cfg.Sink.NewLogger("core"),
|
||||
Output: NewEventHandler(al),
|
||||
ReadPlanStep: func() string {
|
||||
if sessPtr == nil || sessPtr.AgentLog == nil {
|
||||
return ""
|
||||
}
|
||||
sessPtr.AgentLog.mu.RLock()
|
||||
data := make([]byte, len(sessPtr.AgentLog.plan))
|
||||
copy(data, sessPtr.AgentLog.plan)
|
||||
sessPtr.AgentLog.mu.RUnlock()
|
||||
return coresession.NextUncheckedStep(data)
|
||||
},
|
||||
})
|
||||
ReadPlanStep: func() string {
|
||||
if sessPtr == nil || sessPtr.AgentLog() == nil {
|
||||
return ""
|
||||
}
|
||||
al := sessPtr.AgentLog()
|
||||
al.mu.RLock()
|
||||
data := make([]byte, len(al.plan))
|
||||
copy(data, al.plan)
|
||||
al.mu.RUnlock()
|
||||
return coresession.NextUncheckedStep(data)
|
||||
},
|
||||
})
|
||||
|
||||
sessionCtx, sessionCancel := context.WithCancel(rs.cfg.Ctx)
|
||||
sess := NewSession(sessID, core, sessionCtx, sessionCancel)
|
||||
sessPtr = sess
|
||||
sess.uname = uname
|
||||
sess.AgentLog = al
|
||||
sess.proc = proc
|
||||
sessionCtx, sessionCancel := context.WithCancel(rs.cfg.Ctx)
|
||||
sess := NewSession(sessID, core, sessionCtx, sessionCancel)
|
||||
sessPtr = sess
|
||||
sess.uname = uname
|
||||
sess.SetAgentLog(uname, al)
|
||||
sess.proc = proc
|
||||
|
||||
replayMessagesToLog(sess.AgentLog, ps.Messages)
|
||||
replayMessagesToLog(sess.AgentLog(), ps.Messages)
|
||||
|
||||
name := sess.Name()
|
||||
rs.mu.Lock()
|
||||
|
|
|
|||
|
|
@ -231,7 +231,7 @@ func openSessionTree(rs *rootState, sess *Session) *Tree {
|
|||
func openAgentTree(rs *rootState, sess *Session) *Tree {
|
||||
return NewAgentTree(
|
||||
sess,
|
||||
sess.AgentLog,
|
||||
sess.AgentLog(),
|
||||
rs.cfg.Log,
|
||||
nil,
|
||||
rs.cfg.InvalidateModels,
|
||||
|
|
|
|||
|
|
@ -23,8 +23,9 @@ type Session struct {
|
|||
cancel context.CancelFunc
|
||||
proc *toolsrv.Process // toolsrv subprocess (session owns lifecycle)
|
||||
|
||||
// AgentLog for the session's agent (will become a map when multi-agent lands).
|
||||
AgentLog *AgentLog
|
||||
// AgentLogs for the session's agents, keyed by agent ID.
|
||||
// Backward compat: AgentLog() returns the first agent's log (by map iteration order).
|
||||
AgentLogs map[string]*AgentLog
|
||||
|
||||
// Cached model list (expensive API call; refreshed every 24h).
|
||||
modelsMu sync.Mutex
|
||||
|
|
@ -62,6 +63,28 @@ func NewSession(id string, core *session.Session, ctx context.Context, cancel co
|
|||
return sess
|
||||
}
|
||||
|
||||
// AgentLog returns the first agent's log for backward compatibility.
|
||||
// The iteration order over the map is non-deterministic; when multi-agent
|
||||
// support is fully wired, callers should use AgentLogs directly by key.
|
||||
func (sess *Session) AgentLog() *AgentLog {
|
||||
sess.mu.RLock()
|
||||
defer sess.mu.RUnlock()
|
||||
for _, al := range sess.AgentLogs {
|
||||
return al
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetAgentLog stores the agent log for the given agent ID.
|
||||
func (sess *Session) SetAgentLog(agentID string, al *AgentLog) {
|
||||
sess.mu.Lock()
|
||||
defer sess.mu.Unlock()
|
||||
if sess.AgentLogs == nil {
|
||||
sess.AgentLogs = make(map[string]*AgentLog)
|
||||
}
|
||||
sess.AgentLogs[agentID] = al
|
||||
}
|
||||
|
||||
const modelsCacheTTL = 24 * time.Hour
|
||||
|
||||
// CachedListModels returns the model list, using a 24h cache to avoid
|
||||
|
|
|
|||
|
|
@ -43,7 +43,7 @@ type Config struct {
|
|||
}
|
||||
|
||||
// Session is the concrete session type. It owns session-level state and
|
||||
// delegates agent operations to its owned Agent.
|
||||
// delegates agent operations to its owned agents.
|
||||
type Session struct {
|
||||
id string
|
||||
envMu sync.RWMutex
|
||||
|
|
@ -51,7 +51,7 @@ type Session struct {
|
|||
plan []byte
|
||||
prevPrompt string
|
||||
|
||||
r *agent.Agent
|
||||
agents []*agent.Agent
|
||||
log *olog.Logger
|
||||
sessionsDir string
|
||||
listHandlers map[string]func() []string
|
||||
|
|
@ -133,6 +133,7 @@ func New(cfg Config) *Session {
|
|||
a := &Session{
|
||||
id: cfg.SessionID,
|
||||
env: make(map[string]string),
|
||||
agents: make([]*agent.Agent, 0, 1),
|
||||
log: log,
|
||||
auditLog: auditLog,
|
||||
sessionsDir: cfg.SessionsDir,
|
||||
|
|
@ -140,7 +141,7 @@ func New(cfg Config) *Session {
|
|||
listHandlers: cfg.ListHandlers,
|
||||
}
|
||||
|
||||
a.r = agent.NewAgent(agent.AgentCfg{
|
||||
ag := agent.NewAgent(agent.AgentCfg{
|
||||
History: cfg.History,
|
||||
Runtime: rt,
|
||||
AgentName: cfg.AgentName,
|
||||
|
|
@ -159,42 +160,76 @@ func New(cfg Config) *Session {
|
|||
ReadPlanStep: cfg.ReadPlanStep,
|
||||
SaveSession: a.saveSession,
|
||||
FlushSave: a.flushSave,
|
||||
})
|
||||
})
|
||||
|
||||
a.r.SetSessionEnv(a.id)
|
||||
a.agents = append(a.agents, ag)
|
||||
ag.SetSessionEnv(a.id)
|
||||
return a
|
||||
}
|
||||
|
||||
// Agent returns the first (primary) agent. Backward compatible.
|
||||
func (a *Session) Agent() *agent.Agent {
|
||||
if len(a.agents) == 0 {
|
||||
return nil
|
||||
}
|
||||
return a.agents[0]
|
||||
}
|
||||
|
||||
// Agents returns the full agent slice.
|
||||
func (a *Session) Agents() []*agent.Agent { return a.agents }
|
||||
|
||||
// AgentAt returns the agent at the given index, or nil.
|
||||
func (a *Session) AgentAt(idx int) *agent.Agent {
|
||||
if idx < 0 || idx >= len(a.agents) {
|
||||
return nil
|
||||
}
|
||||
return a.agents[idx]
|
||||
}
|
||||
|
||||
// AgentCount returns the number of agents.
|
||||
func (a *Session) AgentCount() int { return len(a.agents) }
|
||||
|
||||
// AddAgent appends an agent to the session.
|
||||
func (a *Session) AddAgent(ag *agent.Agent) {
|
||||
a.agents = append(a.agents, ag)
|
||||
ag.SetSessionEnv(a.id)
|
||||
}
|
||||
|
||||
// Close releases resources for this session.
|
||||
func (a *Session) Close() {
|
||||
a.log.Debug("Close() session=%q", a.id)
|
||||
a.flushSave()
|
||||
a.r.Close()
|
||||
for _, ag := range a.agents {
|
||||
ag.Close()
|
||||
}
|
||||
if a.id != "" {
|
||||
os.RemoveAll(filepath.Join(ollieTmpDir(), a.id)) //nolint:errcheck
|
||||
}
|
||||
}
|
||||
|
||||
// SetEnv stores a session-scoped variable and propagates it to the agent.
|
||||
// SetEnv stores a session-scoped variable and propagates to all agents.
|
||||
func (a *Session) SetEnv(key, value string) {
|
||||
a.envMu.Lock()
|
||||
a.env[key] = value
|
||||
a.envMu.Unlock()
|
||||
a.r.SetEnv(key, value)
|
||||
for _, ag := range a.agents {
|
||||
ag.SetEnv(key, value)
|
||||
}
|
||||
}
|
||||
|
||||
func (a *Session) Agent() *agent.Agent { return a.r }
|
||||
|
||||
// CWD returns the current working directory for tool execution.
|
||||
// CWD returns the current working directory for tool execution
|
||||
// (from the primary agent).
|
||||
func (a *Session) CWD() string {
|
||||
if c := a.r.Cwd(); c != "" {
|
||||
return c
|
||||
if ag := a.Agent(); ag != nil {
|
||||
if c := ag.Cwd(); c != "" {
|
||||
return c
|
||||
}
|
||||
}
|
||||
wd, _ := os.Getwd()
|
||||
return wd
|
||||
}
|
||||
|
||||
// SetCWD validates and sets the working directory.
|
||||
// SetCWD validates and sets the working directory on the primary agent.
|
||||
func (a *Session) SetCWD(dir string) error {
|
||||
dir = paths.ExpandHome(dir)
|
||||
if dir != "" {
|
||||
|
|
@ -202,7 +237,9 @@ func (a *Session) SetCWD(dir string) error {
|
|||
return fmt.Errorf("cwd: %w", err)
|
||||
}
|
||||
}
|
||||
a.r.SetCWD(dir)
|
||||
if ag := a.Agent(); ag != nil {
|
||||
ag.SetCWD(dir)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -223,42 +260,50 @@ func (a *Session) SetSessionID(newID string) error {
|
|||
}
|
||||
}
|
||||
a.id = newID
|
||||
a.r.RenamePreamble(oldID, newID)
|
||||
for _, ag := range a.agents {
|
||||
ag.RenamePreamble(oldID, newID)
|
||||
}
|
||||
oldTemp := filepath.Join(ollieTmpDir(), oldID)
|
||||
newTemp := filepath.Join(ollieTmpDir(), newID)
|
||||
if _, err := os.Stat(oldTemp); err == nil {
|
||||
os.Rename(oldTemp, newTemp) //nolint:errcheck
|
||||
}
|
||||
a.r.SetSessionEnv(newID)
|
||||
if ag := a.Agent(); ag != nil {
|
||||
ag.SetSessionEnv(newID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// WaitChange blocks until the named field changes from current.
|
||||
func (a *Session) WaitChange(ctx context.Context, field, current string) (string, bool) {
|
||||
ag := a.Agent()
|
||||
if ag == nil {
|
||||
return "", false
|
||||
}
|
||||
if field == agent.WatchState {
|
||||
return a.r.WaitChange(ctx, field, current)
|
||||
return ag.WaitChange(ctx, field, current)
|
||||
}
|
||||
// Other fields — read via agent methods, use agent's change signal.
|
||||
read := func() string {
|
||||
switch field {
|
||||
case "usage":
|
||||
return a.r.UsageStr()
|
||||
return ag.UsageStr()
|
||||
case "ctxsz":
|
||||
return a.r.CtxSz()
|
||||
return ag.CtxSz()
|
||||
case "cwd":
|
||||
return a.CWD()
|
||||
case "agent":
|
||||
return a.r.Name()
|
||||
return ag.Name()
|
||||
}
|
||||
return ""
|
||||
}
|
||||
stop := context.AfterFunc(ctx, func() { a.r.BroadcastChange() })
|
||||
stop := context.AfterFunc(ctx, func() { ag.BroadcastChange() })
|
||||
defer stop()
|
||||
for ctx.Err() == nil {
|
||||
if v := read(); v != current {
|
||||
return v, true
|
||||
}
|
||||
a.r.WaitForChange(ctx)
|
||||
ag.WaitForChange(ctx)
|
||||
}
|
||||
return "", false
|
||||
}
|
||||
|
|
@ -293,12 +338,17 @@ func (a *Session) flushSave() {
|
|||
a.log.Error("session save: %v", err)
|
||||
return
|
||||
}
|
||||
if err := a.r.SaveFull(path, a.id, a.CWD(), a.remote); err != nil {
|
||||
a.log.Error("session save: %v", err)
|
||||
if ag := a.Agent(); ag != nil {
|
||||
if err := ag.SaveFull(path, a.id, a.CWD(), a.remote); err != nil {
|
||||
a.log.Error("session save: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SaveSession writes the current session state to the given path.
|
||||
func (a *Session) SaveSession(path string) error {
|
||||
return a.r.SaveFull(path, a.id, a.CWD(), a.remote)
|
||||
if ag := a.Agent(); ag != nil {
|
||||
return ag.SaveFull(path, a.id, a.CWD(), a.remote)
|
||||
}
|
||||
return fmt.Errorf("no agent in session")
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue