split Session into Session + AgentState; move all agent-scoped state to AgentState
Session struct now only holds: id, uname, mu, Core, SessionCtx, cancel, proc, modelsMu/Cache/At. New AgentState struct holds: log, logVers, ChatOffset, chatCond, plan, prevPrompt, remote — all per-agent state. agentHelper now takes an *AgentState alongside *Session. All agent file operations (plan, log, chat, prevPrompt, offset, remote, cfg) access AgentState fields instead of Session fields. Removed startEventLog from NewSession (was writing to session log). replayMessagesToLog takes *AgentState. Rename no longer appends to log. ReadPlanStep stubbed out (TODO: look up correct AgentState). server.go makeStat stubbed out AgentState log size tracking (TODO).
This commit is contained in:
parent
405ae3137a
commit
83304e2896
|
|
@ -1267,32 +1267,16 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
|
||||
// For chat and tailable mutable files, report actual size and
|
||||
// Qid version so polling tools (tail -f) can detect changes via stat.
|
||||
// Path format: /session/{sessid}/{file} or /session/{sessid}/agent/{aid}/{file}
|
||||
// Path format: /session/{sessid}/agent/{aid}/{file}
|
||||
if strings.HasPrefix(path, fs.PathSessions) {
|
||||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
|
||||
if len(parts) == 3 && parts[0] == "session" {
|
||||
if sess := fs.Lookup(s.sessionTree, parts[1]); sess != nil {
|
||||
if base == "log" {
|
||||
length, vers := sess.LogInfo()
|
||||
if length > 64*1024 {
|
||||
length = 64 * 1024
|
||||
}
|
||||
dir.Length = uint64(length)
|
||||
dir.Qid.Vers = vers
|
||||
}
|
||||
}
|
||||
}
|
||||
// Also handle agent paths: /session/{sessid}/agent/{aid}/{file}
|
||||
parts = strings.SplitN(strings.TrimPrefix(path, "/"), "/", 5)
|
||||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 5)
|
||||
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" {
|
||||
length, vers := sess.LogInfo()
|
||||
if length > 64*1024 {
|
||||
length = 64 * 1024
|
||||
}
|
||||
dir.Length = uint64(length)
|
||||
dir.Qid.Vers = vers
|
||||
// TODO: look up the correct AgentState for this agent
|
||||
// length, vers := as.LogInfo()
|
||||
dir.Length = 4096
|
||||
// dir.Qid.Vers = vers
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -231,11 +231,8 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
if sessPtr == nil {
|
||||
return ""
|
||||
}
|
||||
sessPtr.mu.RLock()
|
||||
data := make([]byte, len(sessPtr.plan))
|
||||
copy(data, sessPtr.plan)
|
||||
sessPtr.mu.RUnlock()
|
||||
return coresession.NextUncheckedStep(data)
|
||||
// TODO: look up the correct AgentState for the active agent
|
||||
return ""
|
||||
},
|
||||
})
|
||||
}
|
||||
|
|
@ -243,7 +240,7 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
sessionCtx, sessionCancel := context.WithCancel(rs.cfg.Ctx)
|
||||
sess := NewSession(sessID, core, sessionCtx, sessionCancel)
|
||||
sessPtr = sess
|
||||
sess.remote = remoteTarget
|
||||
// remote is now on AgentState; set via NewAgentState
|
||||
sess.proc = proc
|
||||
|
||||
rs.mu.Lock()
|
||||
|
|
|
|||
|
|
@ -75,8 +75,8 @@ func NewSessionTree(sess *Session, log *olog.Logger, kill func(), rename func(ne
|
|||
}
|
||||
|
||||
// NewAgentTree builds the file tree for a single agent within a session.
|
||||
func NewAgentTree(sess *Session, log *olog.Logger, saveTranscript func([]byte) error, invalidateModels func(), toolRegistry *toolsrv.Registry) *Tree {
|
||||
h := &agentHelper{sess: sess, log: log, saveTranscript: saveTranscript, invalidateModels: invalidateModels, toolRegistry: toolRegistry}
|
||||
func NewAgentTree(sess *Session, as *AgentState, log *olog.Logger, saveTranscript func([]byte) error, invalidateModels func(), toolRegistry *toolsrv.Registry) *Tree {
|
||||
h := &agentHelper{as: as, sess: sess, log: log, saveTranscript: saveTranscript, invalidateModels: invalidateModels, toolRegistry: toolRegistry}
|
||||
perms := sessionFilePerms()
|
||||
specs := make([]FileSpec, len(AgentFileList))
|
||||
for i, f := range AgentFileList {
|
||||
|
|
@ -100,6 +100,7 @@ type sessionHelper struct {
|
|||
|
||||
// agentHelper holds the dependencies needed to build agent FileSpecs.
|
||||
type agentHelper struct {
|
||||
as *AgentState
|
||||
sess *Session
|
||||
log *olog.Logger
|
||||
saveTranscript func([]byte) error
|
||||
|
|
@ -153,11 +154,7 @@ func (h *sessionHelper) handleCtl(input string) error {
|
|||
}
|
||||
}
|
||||
case "save":
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.log))
|
||||
copy(data, h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
return h.saveTranscript(data)
|
||||
h.sess.Core.SaveSession("") // triggers session persistence
|
||||
case "invalidate":
|
||||
if h.invalidateModels != nil {
|
||||
h.invalidateModels()
|
||||
|
|
@ -180,45 +177,45 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
}
|
||||
case "plan":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.plan))
|
||||
copy(data, h.sess.plan)
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
data := make([]byte, len(h.as.plan))
|
||||
copy(data, h.as.plan)
|
||||
h.as.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
fs.Size = func() int64 {
|
||||
h.sess.mu.RLock()
|
||||
n := len(h.sess.plan)
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
n := len(h.as.plan)
|
||||
h.as.mu.RUnlock()
|
||||
return int64(n)
|
||||
}
|
||||
case "prompt.prev":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.prevPrompt))
|
||||
copy(data, h.sess.prevPrompt)
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
data := make([]byte, len(h.as.prevPrompt))
|
||||
copy(data, h.as.prevPrompt)
|
||||
h.as.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
case "log":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
const maxWindow = 64 * 1024
|
||||
h.sess.mu.RLock()
|
||||
log := h.sess.log
|
||||
h.as.mu.RLock()
|
||||
log := h.as.log
|
||||
start := 0
|
||||
if len(log) > maxWindow {
|
||||
start = len(log) - maxWindow
|
||||
}
|
||||
data := make([]byte, len(log)-start)
|
||||
copy(data, log[start:])
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
fs.Size = func() int64 {
|
||||
const maxWindow = 64 * 1024
|
||||
h.sess.mu.RLock()
|
||||
n := len(h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
n := len(h.as.log)
|
||||
h.as.mu.RUnlock()
|
||||
if n > maxWindow {
|
||||
return int64(maxWindow)
|
||||
}
|
||||
|
|
@ -250,10 +247,10 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
}
|
||||
case "plan":
|
||||
fs.Write = func(data []byte) error {
|
||||
h.sess.mu.Lock()
|
||||
h.sess.plan = make([]byte, len(data))
|
||||
copy(h.sess.plan, data)
|
||||
h.sess.mu.Unlock()
|
||||
h.as.mu.Lock()
|
||||
h.as.plan = make([]byte, len(data))
|
||||
copy(h.as.plan, data)
|
||||
h.as.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
case "log":
|
||||
|
|
@ -278,12 +275,12 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
h.sess.InvalidateModelsCache()
|
||||
return nil
|
||||
}
|
||||
h.sess.mu.Lock()
|
||||
h.sess.prevPrompt = []byte(input)
|
||||
h.sess.mu.Unlock()
|
||||
h.as.mu.Lock()
|
||||
h.as.prevPrompt = []byte(input)
|
||||
h.as.mu.Unlock()
|
||||
go func() {
|
||||
h.sess.Core.Agent().Submit(h.sess.SessionCtx, input)
|
||||
h.sess.EnsureTrailingNewline()
|
||||
h.as.EnsureTrailingNewline()
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
|
@ -349,26 +346,26 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
if base != "" {
|
||||
fmt.Sscanf(base, "%d", &offset)
|
||||
} else {
|
||||
h.sess.mu.RLock()
|
||||
offset = len(h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
offset = len(h.as.log)
|
||||
h.as.mu.RUnlock()
|
||||
}
|
||||
|
||||
// Block until new data is available or context cancelled
|
||||
h.sess.mu.RLock()
|
||||
for len(h.sess.log) <= offset {
|
||||
h.as.mu.RLock()
|
||||
for len(h.as.log) <= offset {
|
||||
if ctx.Err() != nil {
|
||||
h.sess.mu.RUnlock()
|
||||
h.as.mu.RUnlock()
|
||||
return nil, "", nil
|
||||
}
|
||||
h.sess.chatCond.Wait()
|
||||
h.as.chatCond.Wait()
|
||||
}
|
||||
|
||||
// Copy new data
|
||||
data := make([]byte, len(h.sess.log)-offset)
|
||||
copy(data, h.sess.log[offset:])
|
||||
newOffset := len(h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
data := make([]byte, len(h.as.log)-offset)
|
||||
copy(data, h.as.log[offset:])
|
||||
newOffset := len(h.as.log)
|
||||
h.as.mu.RUnlock()
|
||||
|
||||
return data, fmt.Sprintf("%d", newOffset), nil
|
||||
}
|
||||
|
|
@ -388,8 +385,8 @@ func (h *sessionHelper) content(name string) string {
|
|||
}
|
||||
|
||||
func (h *agentHelper) content(name string) string {
|
||||
h.sess.mu.RLock()
|
||||
defer h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
defer h.as.mu.RUnlock()
|
||||
switch name {
|
||||
case "usage":
|
||||
return h.sess.Core.Agent().UsageStr() + "\n"
|
||||
|
|
@ -405,7 +402,7 @@ func (h *agentHelper) content(name string) string {
|
|||
case "state":
|
||||
return h.sess.Core.Agent().State() + "\n"
|
||||
case "offset":
|
||||
return fmt.Sprintf("%d\n", h.sess.ChatOffset)
|
||||
return fmt.Sprintf("%d\n", h.as.ChatOffset)
|
||||
case "tail":
|
||||
return "#!/bin/sh\nexec tail -f \"$(dirname \"$0\")/chat\"\n"
|
||||
case "context":
|
||||
|
|
@ -456,8 +453,8 @@ func (h *agentHelper) content(name string) string {
|
|||
}
|
||||
|
||||
func (h *agentHelper) cfgContent() string {
|
||||
h.sess.mu.RLock()
|
||||
defer h.sess.mu.RUnlock()
|
||||
h.as.mu.RLock()
|
||||
defer h.as.mu.RUnlock()
|
||||
p := h.sess.Core.Agent().GenParams()
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "name=%s\n", h.sess.id)
|
||||
|
|
@ -465,7 +462,7 @@ func (h *agentHelper) cfgContent() string {
|
|||
fmt.Fprintf(&sb, "model=%s\n", h.sess.Core.Agent().ModelName())
|
||||
fmt.Fprintf(&sb, "agent=%s\n", h.sess.Core.Agent().Name())
|
||||
fmt.Fprintf(&sb, "cwd=%s\n", h.sess.Core.CWD())
|
||||
fmt.Fprintf(&sb, "remote=%s\n", h.sess.remote)
|
||||
fmt.Fprintf(&sb, "remote=%s\n", h.as.remote)
|
||||
fmt.Fprintf(&sb, "maxTokens=%d\n", p.MaxTokens)
|
||||
fmt.Fprintf(&sb, "maxCompletionTokens=%d\n", p.MaxCompletionTokens)
|
||||
if p.Temperature != nil {
|
||||
|
|
|
|||
|
|
@ -58,7 +58,7 @@ func squashWhitespace(s string) string {
|
|||
|
||||
// replayMessagesToLog renders the tail of a persisted message list into the
|
||||
// session's chat log so that restored sessions show recent history.
|
||||
func replayMessagesToLog(sess *Session, messages []backend.Message) {
|
||||
func replayMessagesToLog(as *AgentState, messages []backend.Message) {
|
||||
const maxReplay = 20
|
||||
start := len(messages) - maxReplay
|
||||
if start < 0 {
|
||||
|
|
@ -73,22 +73,22 @@ func replayMessagesToLog(sess *Session, messages []backend.Message) {
|
|||
case "system":
|
||||
// skip
|
||||
case "user":
|
||||
sess.AppendLog([]byte("[user]\n" + m.Content + "\n"))
|
||||
as.AppendLog([]byte("[user]\n" + m.Content + "\n"))
|
||||
case "assistant":
|
||||
header := "[assistant]"
|
||||
if m.ID != "" {
|
||||
header = "[assistant:" + m.ID + "]"
|
||||
}
|
||||
sess.AppendLog([]byte(header + "\n"))
|
||||
as.AppendLog([]byte(header + "\n"))
|
||||
if m.Content != "" {
|
||||
sess.AppendLog([]byte(m.Content + "\n"))
|
||||
as.AppendLog([]byte(m.Content + "\n"))
|
||||
}
|
||||
for _, tc := range m.ToolCalls {
|
||||
args := squashWhitespace(string(tc.Arguments))
|
||||
sess.AppendLog([]byte("[call:" + tc.Name + "]\n" + args + "\n"))
|
||||
as.AppendLog([]byte("[call:" + tc.Name + "]\n" + args + "\n"))
|
||||
}
|
||||
case "tool":
|
||||
sess.AppendLog([]byte("[tool]\n" + strings.TrimRight(m.Content, "\n") + "\n"))
|
||||
as.AppendLog([]byte("[tool]\n" + strings.TrimRight(m.Content, "\n") + "\n"))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -65,7 +65,6 @@ func Rename(rs *rootState, old, new string) error {
|
|||
delete(rs.sessions, old)
|
||||
rs.mu.Unlock()
|
||||
|
||||
sess.AppendLog([]byte(fmt.Sprintf(":: session renamed: %s -> %s\n", old, new)))
|
||||
rs.cfg.Log.Info("renamed session %s -> %s", old, new)
|
||||
return nil
|
||||
}
|
||||
|
|
@ -326,26 +325,24 @@ func restoreSession(rs *rootState, ps *agent.PersistedAgent) error {
|
|||
PromptEnvExtra: promptEnv,
|
||||
BaseLayers: []string{sysPrompt, opModel, envBlock},
|
||||
Log: rs.cfg.Sink.NewLogger("core"),
|
||||
ReadPlanStep: func() string {
|
||||
if sessPtr == nil {
|
||||
ReadPlanStep: func() string {
|
||||
if sessPtr == nil {
|
||||
return ""
|
||||
}
|
||||
// TODO: look up the correct AgentState for the active agent
|
||||
return ""
|
||||
}
|
||||
sessPtr.mu.RLock()
|
||||
data := make([]byte, len(sessPtr.plan))
|
||||
copy(data, sessPtr.plan)
|
||||
sessPtr.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.remote = remoteTarget
|
||||
// sess.remote is now on AgentState; set via NewAgentState above
|
||||
sess.proc = proc
|
||||
|
||||
replayMessagesToLog(sess, ps.Messages)
|
||||
as := NewAgentState(remoteTarget)
|
||||
replayMessagesToLog(as, ps.Messages)
|
||||
|
||||
rs.mu.Lock()
|
||||
rs.sessions[sessID] = sess
|
||||
|
|
|
|||
|
|
@ -170,6 +170,7 @@ func openSessionTree(rs *rootState, sess *Session) *Tree {
|
|||
func openAgentTree(rs *rootState, sess *Session) *Tree {
|
||||
return NewAgentTree(
|
||||
sess,
|
||||
nil, // AgentState — will be wired when multi-agent lands
|
||||
rs.cfg.Log,
|
||||
nil,
|
||||
rs.cfg.InvalidateModels,
|
||||
|
|
|
|||
|
|
@ -11,7 +11,8 @@ import (
|
|||
"ollie/toolsrv"
|
||||
)
|
||||
|
||||
// Session holds all state for one agent session exposed via 9P.
|
||||
// Session holds session-level state exposed via 9P.
|
||||
// Agent-scoped state (log, plan, etc.) lives in AgentState.
|
||||
type Session struct {
|
||||
mu sync.RWMutex
|
||||
id string
|
||||
|
|
@ -19,28 +20,37 @@ type Session struct {
|
|||
Core *session.Session
|
||||
SessionCtx 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)
|
||||
proc *toolsrv.Process // toolsrv subprocess (session owns lifecycle)
|
||||
|
||||
// Chat stream: condition variable broadcast on every log append.
|
||||
// Readers of the `chat` file block on this until new data arrives.
|
||||
chatCond *sync.Cond
|
||||
|
||||
// Cached model list (expensive API call; refreshed every 24h).
|
||||
modelsMu sync.Mutex
|
||||
modelsCache string
|
||||
modelsCacheAt time.Time
|
||||
}
|
||||
|
||||
// AgentState holds all per-agent state exposed via 9P.
|
||||
type AgentState struct {
|
||||
mu sync.RWMutex
|
||||
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)
|
||||
|
||||
// Chat stream: condition variable broadcast on every log append.
|
||||
// Readers of the `chat` file block on this until new data arrives.
|
||||
chatCond *sync.Cond
|
||||
}
|
||||
|
||||
func NewAgentState(remote string) *AgentState {
|
||||
as := &AgentState{remote: remote}
|
||||
as.chatCond = sync.NewCond(as.mu.RLocker())
|
||||
return as
|
||||
}
|
||||
|
||||
func NewSession(id string, core *session.Session, ctx context.Context, cancel context.CancelFunc) *Session {
|
||||
sess := &Session{id: id, Core: core, SessionCtx: ctx, cancel: cancel}
|
||||
sess.chatCond = sync.NewCond(sess.mu.RLocker())
|
||||
sess.startEventLog()
|
||||
return sess
|
||||
}
|
||||
|
||||
|
|
@ -88,125 +98,51 @@ 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) {
|
||||
// --- AgentState methods ---
|
||||
|
||||
// AppendLog appends data to the agent's log and bumps the version.
|
||||
func (as *AgentState) AppendLog(data []byte) {
|
||||
if len(data) == 0 {
|
||||
return
|
||||
}
|
||||
sess.mu.Lock()
|
||||
sess.log = append(sess.log, data...)
|
||||
sess.logVers++
|
||||
sess.mu.Unlock()
|
||||
sess.chatCond.Broadcast()
|
||||
as.mu.Lock()
|
||||
as.log = append(as.log, data...)
|
||||
as.logVers++
|
||||
as.mu.Unlock()
|
||||
as.chatCond.Broadcast()
|
||||
}
|
||||
|
||||
// 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')
|
||||
func (as *AgentState) EnsureTrailingNewline() {
|
||||
as.mu.Lock()
|
||||
if len(as.log) > 0 && as.log[len(as.log)-1] != '\n' {
|
||||
as.log = append(as.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":
|
||||
if streamingRole == "tool" {
|
||||
if ev.Content == "" {
|
||||
sess.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
} else {
|
||||
sess.AppendLog([]byte(ev.Content))
|
||||
}
|
||||
return
|
||||
}
|
||||
if streamingRole != "" {
|
||||
sess.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
if ev.Content == "" {
|
||||
sess.AppendLog([]byte("[tool:" + ev.Name + "]\n"))
|
||||
streamingRole = "tool"
|
||||
} else {
|
||||
sess.AppendLog(FormatEvent(ev))
|
||||
}
|
||||
|
||||
case "retry":
|
||||
if streamingRole == "retry" {
|
||||
sess.AppendLog([]byte("\n" + ev.Content))
|
||||
return
|
||||
}
|
||||
if streamingRole != "" {
|
||||
sess.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
sess.AppendLog([]byte("[retry]\n"))
|
||||
streamingRole = "retry"
|
||||
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))
|
||||
}
|
||||
})
|
||||
as.logVers++
|
||||
as.mu.Unlock()
|
||||
}
|
||||
|
||||
// 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
|
||||
func (as *AgentState) LogInfo() (length int, vers uint32) {
|
||||
as.mu.RLock()
|
||||
defer as.mu.RUnlock()
|
||||
return len(as.log), as.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 }
|
||||
func (as *AgentState) Log() []byte { return as.log }
|
||||
|
||||
// Plan returns the plan bytes. Caller must hold Mu().RLock().
|
||||
func (sess *Session) Plan() []byte { return sess.plan }
|
||||
func (as *AgentState) Plan() []byte { return as.plan }
|
||||
|
||||
// SetPlan sets the plan. Caller must hold Mu().Lock().
|
||||
func (sess *Session) SetPlan(p []byte) { sess.plan = p }
|
||||
func (as *AgentState) SetPlan(p []byte) { as.plan = p }
|
||||
|
||||
// PrevPrompt returns the last submitted prompt. Caller must hold Mu().RLock().
|
||||
func (sess *Session) PrevPrompt() []byte { return sess.prevPrompt }
|
||||
func (as *AgentState) PrevPrompt() []byte { return as.prevPrompt }
|
||||
|
||||
// Mu returns the agent state's mutex.
|
||||
func (as *AgentState) Mu() *sync.RWMutex { return &as.mu }
|
||||
|
||||
// Remote returns the SSH target for remote execution.
|
||||
func (as *AgentState) Remote() string { return as.remote }
|
||||
|
|
|
|||
Loading…
Reference in New Issue