agent: delegate bus, fifo, state, reply to session.Session
Move session-scoped concerns out of the agent struct: - Bus() → sess.Bus() - Queue()/PopQueue() → sess.Queue()/PopQueue() - State()/setState() → sess.State()/SetState() - Reply() → sess.Reply()/SetReply() Remove corresponding fields from agent struct. Agent-level WaitChange still uses its own changeMu/changeCond since it observes multiple fields (state, usage, ctxsz, cwd, agent).
This commit is contained in:
parent
457ec2312f
commit
64f3f18573
|
|
@ -22,6 +22,7 @@ import (
|
|||
"ollie/backend"
|
||||
olog "ollie/log"
|
||||
"ollie/paths"
|
||||
"ollie/session"
|
||||
"ollie/tools"
|
||||
)
|
||||
|
||||
|
|
@ -309,6 +310,7 @@ type AgentCoreConfig struct {
|
|||
// agent is the Core implementation. It owns all agent and session state
|
||||
// but has no knowledge of how output is rendered.
|
||||
type agent struct {
|
||||
sess *session.Session
|
||||
history *History
|
||||
runtime *Runtime
|
||||
cfg agentConfig // per-turn config built from runtime; set in executeTurn
|
||||
|
|
@ -330,12 +332,8 @@ type agent struct {
|
|||
startupMessages []string
|
||||
currentAction atomic.Pointer[actionHandle]
|
||||
toolCallCount atomic.Int64
|
||||
fifo PromptFIFO
|
||||
bus *pubsub.Bus
|
||||
pendingInject atomic.Pointer[string]
|
||||
mu sync.RWMutex
|
||||
state string // "idle", "thinking", "calling: <tool>"
|
||||
reply string // assistant text from the most recently completed turn
|
||||
envMu sync.RWMutex
|
||||
env map[string]string // session-scoped env vars
|
||||
changeMu sync.Mutex
|
||||
|
|
@ -478,6 +476,11 @@ func NewAgentCore(cfg AgentCoreConfig) Core {
|
|||
}
|
||||
|
||||
a := &agent{
|
||||
sess: session.New(session.Config{
|
||||
ID: cfg.SessionID,
|
||||
Uname: cfg.Uname,
|
||||
CWD: paths.ExpandHome(cfg.CWD),
|
||||
}),
|
||||
history: cfg.History,
|
||||
runtime: rt,
|
||||
log: log,
|
||||
|
|
@ -496,8 +499,6 @@ func NewAgentCore(cfg AgentCoreConfig) Core {
|
|||
newBackend: cfg.NewBackend,
|
||||
readPlanStep: readPlanStep,
|
||||
listHandlers: cfg.ListHandlers,
|
||||
state: "idle",
|
||||
bus: pubsub.NewBus(),
|
||||
}
|
||||
a.changeCond = sync.NewCond(&a.changeMu)
|
||||
a.turnError = func(_ context.Context, errType, errMsg string) HookResult {
|
||||
|
|
@ -726,9 +727,7 @@ func (s *agent) ModelName() string {
|
|||
}
|
||||
|
||||
func (s *agent) State() string {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
return s.state
|
||||
return s.sess.State()
|
||||
}
|
||||
|
||||
func (s *agent) notifyChange() {
|
||||
|
|
@ -737,6 +736,12 @@ func (s *agent) notifyChange() {
|
|||
s.changeMu.Unlock()
|
||||
}
|
||||
|
||||
func (s *agent) setState(state string) {
|
||||
s.sess.SetState(state)
|
||||
s.log.Debug("state -> %q", state)
|
||||
s.notifyChange()
|
||||
}
|
||||
|
||||
// WaitChange blocks until the named field changes from current, then returns
|
||||
// the new value. Returns ("", false) if ctx is cancelled.
|
||||
func (s *agent) WaitChange(ctx context.Context, field, current string) (string, bool) {
|
||||
|
|
@ -776,19 +781,11 @@ func (s *agent) WaitChange(ctx context.Context, field, current string) (string,
|
|||
return "", false
|
||||
}
|
||||
|
||||
func (s *agent) setState(state string) {
|
||||
s.mu.Lock()
|
||||
s.state = state
|
||||
s.mu.Unlock()
|
||||
s.log.Debug("state -> %q", state)
|
||||
s.notifyChange()
|
||||
}
|
||||
|
||||
func (s *agent) Reply() string {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
s.log.Debug("Reply() len=%d", len(s.reply))
|
||||
return s.reply
|
||||
r := s.sess.Reply()
|
||||
s.log.Debug("Reply() len=%d", len(r))
|
||||
return r
|
||||
}
|
||||
|
||||
// CWD returns the current working directory for tool execution.
|
||||
|
|
@ -1033,7 +1030,7 @@ func (s *agent) Inject(prompt string) {
|
|||
// If an inject is already pending, fall back to the normal FIFO so nothing
|
||||
// is lost. Use CompareAndSwap to avoid a race between the nil check and store.
|
||||
if !s.pendingInject.CompareAndSwap(nil, &prompt) {
|
||||
s.fifo.Push(prompt)
|
||||
s.sess.Queue(prompt)
|
||||
return
|
||||
}
|
||||
s.emit(Event{Role: "info", Content: "\n"})
|
||||
|
|
@ -1047,32 +1044,26 @@ func (s *agent) injectRewrite(prompt string) {
|
|||
}
|
||||
|
||||
func (s *agent) Queue(prompt string) {
|
||||
s.log.Debug("Queue() len=%d", len(prompt))
|
||||
s.fifo.Push(prompt)
|
||||
s.bus.Publish("queued", prompt)
|
||||
if !s.IsRunning() {
|
||||
go s.drainQueue()
|
||||
}
|
||||
s.sess.Queue(prompt)
|
||||
s.sess.Bus().Publish("queued", prompt)
|
||||
}
|
||||
|
||||
func (s *agent) drainQueue() {
|
||||
if prompt, ok := s.fifo.Pop(); ok {
|
||||
if prompt, ok := s.sess.PopQueue(); ok {
|
||||
s.Submit(context.Background(), prompt)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *agent) Bus() *pubsub.Bus {
|
||||
return s.bus
|
||||
return s.sess.Bus()
|
||||
}
|
||||
|
||||
func (s *agent) emit(ev Event) {
|
||||
s.bus.Publish("event", ev)
|
||||
s.sess.Bus().Publish("event", ev)
|
||||
}
|
||||
|
||||
func (s *agent) PopQueue() (string, bool) {
|
||||
v, ok := s.fifo.Pop()
|
||||
s.log.Debug("PopQueue() ok=%v len=%d", ok, len(v))
|
||||
return v, ok
|
||||
return s.sess.PopQueue()
|
||||
}
|
||||
|
||||
func (s *agent) IsRunning() bool {
|
||||
|
|
@ -1221,7 +1212,7 @@ func (s *agent) Submit(ctx context.Context, input string) {
|
|||
if s.handleCommand(ctx, input) {
|
||||
return
|
||||
}
|
||||
s.fifo.Push(input)
|
||||
s.sess.Queue(input)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -1234,7 +1225,7 @@ func (s *agent) Submit(ctx context.Context, input string) {
|
|||
return
|
||||
}
|
||||
if s.IsRunning() {
|
||||
s.fifo.Push(input)
|
||||
s.sess.Queue(input)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -1474,7 +1465,7 @@ func (s *agent) executeTurn(ctx context.Context, input string) string {
|
|||
}
|
||||
|
||||
s.mu.Lock()
|
||||
s.reply = replyBuf.String()
|
||||
s.sess.SetReply(replyBuf.String())
|
||||
s.mu.Unlock()
|
||||
replyBuf.Reset()
|
||||
s.setState("idle")
|
||||
|
|
@ -1493,7 +1484,7 @@ func (s *agent) executeTurn(ctx context.Context, input string) string {
|
|||
}
|
||||
s.emit(Event{Role: "error", Content: err.Error()})
|
||||
// Drain one FIFO item — the turnError hook may have queued a recovery prompt.
|
||||
if next, ok := s.fifo.Pop(); ok {
|
||||
if next, ok := s.sess.PopQueue(); ok {
|
||||
return next
|
||||
}
|
||||
return ""
|
||||
|
|
@ -1520,7 +1511,7 @@ func (s *agent) executeTurn(ctx context.Context, input string) string {
|
|||
s.emit(Event{Role: "info", Content: fmt.Sprintf("costLast=$%.4f\n", s.history.LastTurnCostUSD)})
|
||||
}
|
||||
s.auditLog.Debug("turn: end reply=%s cost=$%.4f session_total=$%.4f session=%s",
|
||||
auditTruncate(s.reply), s.history.LastTurnCostUSD, s.history.SessionCostUSD, s.sessionID)
|
||||
auditTruncate(s.sess.Reply()), s.history.LastTurnCostUSD, s.history.SessionCostUSD, s.sessionID)
|
||||
s.notifyChange()
|
||||
}
|
||||
s.saveSession()
|
||||
|
|
@ -1537,7 +1528,7 @@ func (s *agent) executeTurn(ctx context.Context, input string) string {
|
|||
}
|
||||
|
||||
// Drain one item from the FIFO; the outer loop handles the rest.
|
||||
if next, ok := s.fifo.Pop(); ok {
|
||||
if next, ok := s.sess.PopQueue(); ok {
|
||||
return next
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -8,13 +8,16 @@ import (
|
|||
"time"
|
||||
|
||||
olog "ollie/log"
|
||||
"ollie/session"
|
||||
)
|
||||
|
||||
// newTestCore returns a minimal agent with changeCond wired up.
|
||||
// newTestCore returns a minimal agent with session wired up.
|
||||
func newTestCore(initialState string) *agent {
|
||||
sess := session.New(session.Config{ID: "test"})
|
||||
sess.SetState(initialState)
|
||||
a := &agent{
|
||||
state: initialState,
|
||||
log: olog.NewWriter("test", olog.LevelError+1, io.Discard, io.Discard),
|
||||
sess: sess,
|
||||
log: olog.NewWriter("test", olog.LevelError+1, io.Discard, io.Discard),
|
||||
}
|
||||
a.changeCond = sync.NewCond(&a.changeMu)
|
||||
return a
|
||||
|
|
|
|||
Reference in New Issue