198 lines
5.5 KiB
Go
198 lines
5.5 KiB
Go
package mgr
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"strings"
|
|
"time"
|
|
|
|
"ollie/session"
|
|
"ollie/agent"
|
|
)
|
|
|
|
// 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 *session.Session
|
|
Ctx context.Context
|
|
cancel context.CancelFunc
|
|
log []byte
|
|
logVers uint32
|
|
ChatOffset int
|
|
plan []byte
|
|
prevPrompt []byte // last submitted prompt; overwritten on each new submission
|
|
remote string // SSH target for remote execution (empty = local)
|
|
|
|
// Cached model list (expensive API call; refreshed every 24h).
|
|
modelsMu sync.Mutex
|
|
modelsCache string
|
|
modelsCacheAt time.Time
|
|
}
|
|
|
|
func NewSession(id string, core *session.Session, ctx context.Context, cancel context.CancelFunc) *Session {
|
|
sess := &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
|
|
sess.startEventLog()
|
|
return sess
|
|
}
|
|
|
|
const modelsCacheTTL = 24 * time.Hour
|
|
|
|
// CachedListModels returns the model list, using a 24h cache to avoid
|
|
// repeated expensive API calls.
|
|
// NOTE: Must NOT be called while holding sess.mu (caller content() holds RLock).
|
|
// Uses its own modelsMu to avoid deadlock.
|
|
func (sess *Session) CachedListModels() string {
|
|
sess.modelsMu.Lock()
|
|
if sess.modelsCache != "" && time.Since(sess.modelsCacheAt) < modelsCacheTTL {
|
|
result := sess.modelsCache
|
|
sess.modelsMu.Unlock()
|
|
return result
|
|
}
|
|
sess.modelsMu.Unlock()
|
|
|
|
// Cache miss — fetch and store.
|
|
models := sess.Core.Agent().ListModels()
|
|
result := strings.Join(models, "\n")
|
|
sess.modelsMu.Lock()
|
|
sess.modelsCache = result
|
|
sess.modelsCacheAt = time.Now()
|
|
sess.modelsMu.Unlock()
|
|
return result
|
|
}
|
|
|
|
// InvalidateModelsCache clears the cached model list.
|
|
func (sess *Session) InvalidateModelsCache() {
|
|
sess.modelsMu.Lock()
|
|
sess.modelsCache = ""
|
|
sess.modelsCacheAt = time.Time{}
|
|
sess.modelsMu.Unlock()
|
|
}
|
|
|
|
|
|
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.Agent().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() {
|
|
streamingRole := "" // tracks current streaming role ("assistant", "reasoning", or "tool")
|
|
streamingResponseID := ""
|
|
|
|
sess.Core.Bus().Subscribe("event", func(ev agent.Event) {
|
|
switch ev.Role {
|
|
case "assistant", "reasoning":
|
|
if streamingRole != ev.Role || (ev.Role == "assistant" && ev.ResponseID != streamingResponseID) {
|
|
if streamingRole != "" {
|
|
sess.AppendLog([]byte("\n"))
|
|
}
|
|
if ev.Role == "assistant" && ev.ResponseID != "" {
|
|
sess.AppendLog([]byte("[assistant:" + ev.ResponseID + "]\n"))
|
|
streamingResponseID = ev.ResponseID
|
|
} else {
|
|
sess.AppendLog([]byte("[" + ev.Role + "]\n"))
|
|
}
|
|
if ev.Role == "assistant" {
|
|
sess.mu.Lock()
|
|
sess.ChatOffset = len(sess.log)
|
|
sess.mu.Unlock()
|
|
}
|
|
streamingRole = ev.Role
|
|
}
|
|
sess.AppendLog(FormatEvent(ev))
|
|
|
|
case "tool":
|
|
// Tool events: empty content = header-only (stream start),
|
|
// non-empty while streaming = chunk, non-empty without streaming = full result.
|
|
if streamingRole == "tool" {
|
|
if ev.Content == "" {
|
|
// End of stream.
|
|
sess.AppendLog([]byte("\n"))
|
|
streamingRole = ""
|
|
} else {
|
|
sess.AppendLog([]byte(ev.Content))
|
|
}
|
|
return
|
|
}
|
|
if streamingRole != "" {
|
|
sess.AppendLog([]byte("\n"))
|
|
streamingRole = ""
|
|
}
|
|
if ev.Content == "" {
|
|
// Stream start: write header, enter streaming mode.
|
|
sess.AppendLog([]byte("[tool:" + ev.Name + "]\n"))
|
|
streamingRole = "tool"
|
|
} else {
|
|
// Non-streamed tool result (fast tool, no streaming happened).
|
|
sess.AppendLog(FormatEvent(ev))
|
|
}
|
|
|
|
case "call":
|
|
if streamingRole != "" {
|
|
sess.AppendLog([]byte("\n"))
|
|
streamingRole = ""
|
|
}
|
|
sess.AppendLog(FormatEvent(ev))
|
|
|
|
default:
|
|
if streamingRole != "" {
|
|
sess.AppendLog([]byte("\n"))
|
|
streamingRole = ""
|
|
}
|
|
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
|
|
}
|
|
|
|
// Mu returns the session's mutex for external synchronization (e.g. D-Bus adapter).
|
|
func (sess *Session) Mu() *sync.RWMutex { return &sess.mu }
|
|
|
|
// Log returns the raw chat log bytes. Caller must hold Mu().RLock().
|
|
func (sess *Session) Log() []byte { return sess.log }
|
|
|
|
// Plan returns the plan bytes. Caller must hold Mu().RLock().
|
|
func (sess *Session) Plan() []byte { return sess.plan }
|
|
|
|
// SetPlan sets the plan. Caller must hold Mu().Lock().
|
|
func (sess *Session) SetPlan(p []byte) { sess.plan = p }
|
|
|
|
// PrevPrompt returns the last submitted prompt. Caller must hold Mu().RLock().
|
|
func (sess *Session) PrevPrompt() []byte { return sess.prevPrompt }
|