session: inject bus/log/audit into Agent (prep for Submit move)
Agent now holds all dependencies needed to run turns independently: bus, log, auditLog, sessionID, startupMessages, readPlanStep, saveSession. emit() moved to Agent. Session no longer owns startupMessages/readPlanStep.
This commit is contained in:
parent
8a4f77c118
commit
467e6999b1
|
|
@ -5,7 +5,9 @@ import (
|
|||
"sync"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/simonfxr/pubsub"
|
||||
"ollie/backend"
|
||||
olog "ollie/log"
|
||||
"ollie/tools"
|
||||
)
|
||||
|
||||
|
|
@ -39,6 +41,15 @@ type Agent struct {
|
|||
stateMu sync.RWMutex
|
||||
changeMu sync.Mutex
|
||||
changeCond *sync.Cond
|
||||
|
||||
// Injected session-level dependencies (set at creation, stable for agent lifetime).
|
||||
bus *pubsub.Bus
|
||||
log *olog.Logger
|
||||
auditLog *olog.Logger
|
||||
sessionID string // the owning session's ID
|
||||
startupMessages []string
|
||||
readPlanStep func() string
|
||||
saveSession func() // trigger debounced persistence
|
||||
}
|
||||
|
||||
// Backend returns the active backend from the runtime.
|
||||
|
|
@ -151,6 +162,11 @@ func (ag *Agent) InitCond() {
|
|||
ag.changeCond = sync.NewCond(&ag.changeMu)
|
||||
}
|
||||
|
||||
// emit publishes an event on the agent's bus.
|
||||
func (ag *Agent) emit(ev Event) {
|
||||
ag.bus.Publish("event", ev)
|
||||
}
|
||||
|
||||
// CWD returns the agent's working directory.
|
||||
func (ag *Agent) Cwd() string {
|
||||
ag.stateMu.RLock()
|
||||
|
|
|
|||
|
|
@ -865,7 +865,7 @@ func TestSubmit_BackendError(t *testing.T) {
|
|||
|
||||
func TestSubmit_StartupMessages(t *testing.T) {
|
||||
c := newCore(t, nil, nil)
|
||||
c.startupMessages = []string{"startup msg 1", "startup msg 2"}
|
||||
c.r.startupMessages = []string{"startup msg 1", "startup msg 2"}
|
||||
evs := collectEvents(context.Background(), c, "hello")
|
||||
found := 0
|
||||
for _, s := range byRole(evs, "info") {
|
||||
|
|
@ -876,7 +876,7 @@ func TestSubmit_StartupMessages(t *testing.T) {
|
|||
if found != 2 {
|
||||
t.Errorf("startup messages: found %d info events; want 2; got: %v", found, byRole(evs, "info"))
|
||||
}
|
||||
if c.startupMessages != nil {
|
||||
if c.r.startupMessages != nil {
|
||||
t.Error("startupMessages not cleared after first turn")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -320,11 +320,9 @@ type Session struct {
|
|||
r *Agent
|
||||
log *olog.Logger
|
||||
sessionsDir string
|
||||
readPlanStep func() string
|
||||
listHandlers map[string]func() []string
|
||||
turnError func(ctx context.Context, errType, errMsg string) HookResult // overridable for tests
|
||||
remote string // SSH target for remote execution
|
||||
startupMessages []string
|
||||
mu sync.RWMutex
|
||||
auditLog *olog.Logger
|
||||
|
||||
|
|
@ -455,30 +453,36 @@ func New(cfg Config) *Session {
|
|||
log = olog.NewWriter("core", olog.LevelError+1, io.Discard, io.Discard)
|
||||
}
|
||||
|
||||
bus := pubsub.NewBus()
|
||||
auditLog := log.Sub("audit")
|
||||
|
||||
a := &Session{
|
||||
id: cfg.SessionID,
|
||||
|
||||
bus: pubsub.NewBus(),
|
||||
env: make(map[string]string),
|
||||
id: cfg.SessionID,
|
||||
bus: bus,
|
||||
env: make(map[string]string),
|
||||
r: &Agent{
|
||||
history: cfg.History,
|
||||
runtime: rt,
|
||||
cwd: paths.ExpandHome(cfg.CWD),
|
||||
id: cfg.AgentID,
|
||||
agentName: cfg.AgentName,
|
||||
agentsDir: cfg.AgentsDir,
|
||||
promptEnvExtra: cfg.PromptEnvExtra,
|
||||
baseLayers: cfg.BaseLayers,
|
||||
newDispatcher: cfg.NewDispatcher,
|
||||
newBackend: cfg.NewBackend,
|
||||
history: cfg.History,
|
||||
runtime: rt,
|
||||
cwd: paths.ExpandHome(cfg.CWD),
|
||||
id: cfg.AgentID,
|
||||
agentName: cfg.AgentName,
|
||||
agentsDir: cfg.AgentsDir,
|
||||
promptEnvExtra: cfg.PromptEnvExtra,
|
||||
baseLayers: cfg.BaseLayers,
|
||||
newDispatcher: cfg.NewDispatcher,
|
||||
newBackend: cfg.NewBackend,
|
||||
bus: bus,
|
||||
log: log,
|
||||
auditLog: auditLog,
|
||||
sessionID: cfg.SessionID,
|
||||
startupMessages: rt.Messages,
|
||||
readPlanStep: readPlanStep,
|
||||
},
|
||||
log: log,
|
||||
auditLog: log.Sub("audit"),
|
||||
sessionsDir: cfg.SessionsDir,
|
||||
remote: cfg.Remote,
|
||||
startupMessages: rt.Messages,
|
||||
readPlanStep: readPlanStep,
|
||||
listHandlers: cfg.ListHandlers,
|
||||
log: log,
|
||||
auditLog: auditLog,
|
||||
sessionsDir: cfg.SessionsDir,
|
||||
remote: cfg.Remote,
|
||||
listHandlers: cfg.ListHandlers,
|
||||
}
|
||||
a.r.InitCond()
|
||||
a.r.state = "idle"
|
||||
|
|
@ -1248,11 +1252,11 @@ func (a *Session) executeTurn(ctx context.Context, input string) string {
|
|||
}
|
||||
|
||||
if a.r.history == nil {
|
||||
for _, msg := range a.startupMessages {
|
||||
for _, msg := range a.r.startupMessages {
|
||||
a.log.Debug("startup: %s", msg)
|
||||
a.emit(infoEvent(msg))
|
||||
}
|
||||
a.startupMessages = nil
|
||||
a.r.startupMessages = nil
|
||||
a.r.history = newHistory(input)
|
||||
if sc := a.spawnContext(ctx); sc != "" {
|
||||
a.r.history.appendUserMessage(sc)
|
||||
|
|
@ -1279,7 +1283,7 @@ func (a *Session) executeTurn(ctx context.Context, input string) string {
|
|||
ClassifyTier: a.r.runtime.ClassifyTier,
|
||||
GenerationParams: a.r.runtime.GenParams,
|
||||
MaxSteps: a.r.runtime.MaxSteps,
|
||||
ReadPlanStep: a.readPlanStep,
|
||||
ReadPlanStep: a.r.readPlanStep,
|
||||
TurnError: a.turnError,
|
||||
}
|
||||
|
||||
|
|
|
|||
Reference in New Issue