update: agent loop, session persistence, 9P filesystem, server, and remote sandbox
This commit is contained in:
parent
83304e2896
commit
937e6aef3b
|
|
@ -10,7 +10,6 @@ import (
|
|||
"sync/atomic"
|
||||
"syscall"
|
||||
|
||||
"github.com/simonfxr/pubsub"
|
||||
"ollie/backend"
|
||||
olog "ollie/log"
|
||||
"ollie/toolsrv"
|
||||
|
|
@ -47,8 +46,8 @@ type Agent struct {
|
|||
changeMu sync.Mutex
|
||||
changeCond *sync.Cond
|
||||
|
||||
// Injected session-level dependencies (set at creation, stable for agent lifetime).
|
||||
bus *pubsub.Bus
|
||||
// Event handler — receives events from this agent.
|
||||
output EventHandler
|
||||
log *olog.Logger
|
||||
auditLog *olog.Logger
|
||||
sessionID string // the owning session's ID
|
||||
|
|
@ -166,9 +165,11 @@ func (ag *Agent) InitCond() {
|
|||
ag.changeCond = sync.NewCond(&ag.changeMu)
|
||||
}
|
||||
|
||||
// emit publishes an event on the agent's bus.
|
||||
// emit sends an event to the agent's output handler.
|
||||
func (ag *Agent) emit(ev Event) {
|
||||
ag.bus.Publish("event", ev)
|
||||
if ag.output != nil {
|
||||
ag.output(ev)
|
||||
}
|
||||
}
|
||||
|
||||
// CWD returns the agent's working directory.
|
||||
|
|
|
|||
|
|
@ -3,7 +3,6 @@ package agent
|
|||
import (
|
||||
"sync"
|
||||
|
||||
"github.com/simonfxr/pubsub"
|
||||
"ollie/backend"
|
||||
olog "ollie/log"
|
||||
"ollie/toolsrv"
|
||||
|
|
@ -21,7 +20,7 @@ type AgentCfg struct {
|
|||
PromptEnvExtra []string
|
||||
NewToolServer func() toolsrv.Runner
|
||||
NewBackend func(string) (backend.Backend, error)
|
||||
Bus *pubsub.Bus
|
||||
Output EventHandler
|
||||
Log *olog.Logger
|
||||
AuditLog *olog.Logger
|
||||
SessionID string
|
||||
|
|
@ -44,7 +43,7 @@ func NewAgent(cfg AgentCfg) *Agent {
|
|||
promptEnvExtra: cfg.PromptEnvExtra,
|
||||
newToolServer: cfg.NewToolServer,
|
||||
newBackend: cfg.NewBackend,
|
||||
bus: cfg.Bus,
|
||||
output: cfg.Output,
|
||||
log: cfg.Log,
|
||||
auditLog: cfg.AuditLog,
|
||||
sessionID: cfg.SessionID,
|
||||
|
|
|
|||
|
|
@ -133,6 +133,7 @@ func (ag *Agent) executeTurn(ctx context.Context, input string) string {
|
|||
ag.cfg = agentConfig{
|
||||
Backend: ag.runtime.Backend,
|
||||
preamble: ag.runtime.Preamble,
|
||||
Output: ag.output,
|
||||
Tools: tools,
|
||||
Exec: ag.runtime.Exec,
|
||||
ClassifyTool: ag.runtime.ClassifyTool,
|
||||
|
|
|
|||
|
|
@ -1,174 +0,0 @@
|
|||
# Default execute_code sandbox configuration.
|
||||
filesystem:
|
||||
ro:
|
||||
- "{XDG_CONFIG_HOME}/git"
|
||||
- "{HOME}/.netrc"
|
||||
- "/etc/gitconfig"
|
||||
- "/etc/gitattributes"
|
||||
- "/etc/ssh/ssh_config"
|
||||
- "/etc/ssh/ssh_known_hosts"
|
||||
- "/etc/magic"
|
||||
- "/var/log"
|
||||
- "/run/systemd"
|
||||
- "/var/log/journal"
|
||||
- "/etc/systemd"
|
||||
- "/proc"
|
||||
- "/sys"
|
||||
- "{HOME}/.jenkins"
|
||||
- "{HOME}/.psqlrc"
|
||||
- "{HOME}/.pgpass"
|
||||
- "{HOME}/.aws"
|
||||
- "{HOME}/.ansible"
|
||||
- "{HOME}/.ansible.cfg"
|
||||
- "{HOME}/.jira.d"
|
||||
- "{HOME}/.my.cnf"
|
||||
- "{HOME}/.mylogin.cnf"
|
||||
- "{XDG_CONFIG_HOME}/gh"
|
||||
- "{XDG_DATA_HOME}/gh/extensions"
|
||||
- "/dev"
|
||||
- "{HOME}/img"
|
||||
- "{OLLIE_CFG_PATH}"
|
||||
rw:
|
||||
- "{OLLIE_MEMORY_PATH}"
|
||||
- "{OLLIE_PLAN_PATH}"
|
||||
- "{HOME}/.ssh/known_hosts"
|
||||
- "{HOME}/pCloudDrive"
|
||||
- "{HOME}/.pcloud"
|
||||
- "{HOME}/.git-credentials"
|
||||
- "{XDG_CONFIG_HOME}/git/credentials"
|
||||
- "{XDG_RUNTIME_DIR}/git/"
|
||||
- "{HOME}/.mcpt"
|
||||
- "{HOME}/doc"
|
||||
- "{HOME}/.cache/mise"
|
||||
- "{HOME}/.locale/state/mise"
|
||||
- "{HOME}/.gnupg"
|
||||
- "{XDG_DOCUMENTS_DIR}"
|
||||
- "{XDG_CONFIG_HOME}/.jira"
|
||||
- "{HOME}/.m2"
|
||||
- "{HOME}/.gradle"
|
||||
- "/var/run/docker.sock"
|
||||
- "{HOME}/.docker"
|
||||
- "{HOME}/.pnpm-store"
|
||||
- "{XDG_CONFIG_HOME}/yarn"
|
||||
- "{HOME}/.yarnrc.yml"
|
||||
- "{HOME}/.npmrc"
|
||||
- "{HOME}/.ansible/tmp"
|
||||
- "{XDG_RUNTIME_DIR}"
|
||||
- "{XDG_DATA_HOME}/containers"
|
||||
- "{XDG_CONFIG_HOME}/containers"
|
||||
- "{XDG_CONFIG_HOME}/pip"
|
||||
- "{HOME}/.pip"
|
||||
- "{HOME}/.cargo"
|
||||
- "{HOME}/.rustup"
|
||||
- "{XDG_CONFIG_HOME}/nvim"
|
||||
- "{XDG_CONFIG_HOME}/vim"
|
||||
- "{XDG_DATA_HOME}/nvim"
|
||||
- "{XDG_STATE_HOME}/nvim"
|
||||
- "{OLLIE_MEMORY_PATH}"
|
||||
- "{OLLIE_PLAN_PATH}"
|
||||
- "{HOME}/.vimrc"
|
||||
- "{HOME}/.vim"
|
||||
- "{HOME}/.emacs.d"
|
||||
- "{HOME}/.config/emacs"
|
||||
- "/etc"
|
||||
- "/usr/local/etc"
|
||||
- "/dev/urandom"
|
||||
- "/dev/random"
|
||||
- "/dev/null"
|
||||
- "/dev/fuse"
|
||||
- "{OLLIE_DATA_PATH}"
|
||||
- "{HOME}/notes"
|
||||
rox:
|
||||
- "/usr"
|
||||
- "/lib"
|
||||
- "/lib64"
|
||||
- "/bin"
|
||||
- "/sbin"
|
||||
- "{HOME}/go/bin"
|
||||
- "{HOME}/env"
|
||||
- "/proc/self/fd"
|
||||
- "/proc/self/cmdline"
|
||||
- "{OLLIE_TOOLS_PATH}"
|
||||
rwx:
|
||||
- "{CWD}"
|
||||
- "{HOME}/Qt"
|
||||
- "{HOME}/.local/share/pnpm"
|
||||
- "{HOME}/.local/share/pipx"
|
||||
- "{HOME}/.local/share/aqt/"
|
||||
- "{HOME}/prj/aicoder"
|
||||
- "{TMPDIR}"
|
||||
- "{PLAN9}"
|
||||
- "{HOME}/.cache/uv"
|
||||
- "{HOME}/bin"
|
||||
- "{OLLIE}"
|
||||
- "{OLLIE_CFG_PATH}"
|
||||
- "{OLLIE_DATA_PATH}"
|
||||
- "{OLLIE_MEMORY_PATH}"
|
||||
- "{OLLIE_PLAN_PATH}"
|
||||
- "{OLLIE_SKILLS_PATH}"
|
||||
- "{OLLIE_TMP_PATH}"
|
||||
- "{OLLIE_TOOLS_PATH}"
|
||||
- "{OLLIE_TRANSCRIPT_PATH}"
|
||||
- "{OLLIE_AGENTS_PATH}"
|
||||
- "{OLLIE_PROMPTS_PATH}"
|
||||
# XDG defaults for OLLIE_* paths (in case env vars are unset)
|
||||
- "{XDG_CONFIG_HOME}/ollie"
|
||||
- "{XDG_DATA_HOME}/ollie"
|
||||
- "{HOME}/.pyenv"
|
||||
- "{HOME}/.local/bin"
|
||||
- "{HOME}/.sdkman"
|
||||
- "{HOME}/go"
|
||||
- "{HOME}/.cache/go-build"
|
||||
- "{HOME}/opt"
|
||||
- "{HOME}/src"
|
||||
- "{HOME}/prj"
|
||||
- "{HOME}/.nvm"
|
||||
- "{HOME}/.please"
|
||||
- "{HOME}/.npm"
|
||||
- "{HOME}/.mcp-auth"
|
||||
- "{HOME}/.yarn"
|
||||
- "{HOME}/.local/share/mise"
|
||||
- "{PCFS}"
|
||||
- "{XDG_RUNTIME_DIR}/{WAYLAND_DISPLAY}"
|
||||
- "{HOME}/ws"
|
||||
env:
|
||||
- PATH
|
||||
- NAMESPACE
|
||||
- OLLIE
|
||||
- OLLIE_DEFAULT_AGENT
|
||||
- OLLIE_CFG_PATH
|
||||
- OLLIE_DATA_PATH
|
||||
- OLLIE_MEMORY_PATH
|
||||
- OLLIE_PLAN_PATH
|
||||
- OLLIE_SESSION_ID
|
||||
- OLLIE_SKILLS_PATH
|
||||
- OLLIE_TMP_PATH
|
||||
- OLLIE_ELEVATE_SOCKET
|
||||
- OLLIE_TOOLS_PATH
|
||||
- OLLIE_TRANSCRIPT_PATH
|
||||
- OLLIE_AGENTS_PATH
|
||||
- OLLIE_PROMPTS_PATH
|
||||
- OLLIE_UNAME
|
||||
- OLLIE_LOGSEQ_TOKEN
|
||||
- HOME
|
||||
- TMPDIR
|
||||
- PLAN9
|
||||
- JENKINS_USER_ID
|
||||
- JENKINS_API_TOKEN
|
||||
- JENKINS_MCP_AUTH
|
||||
- JIRA_API_TOKEN
|
||||
- BITBUCKET_ACCESS_TOKEN
|
||||
- BITBUCKET_TOKEN
|
||||
- BITBUCKET_USERNAME
|
||||
- PCLOUD_USER
|
||||
- EDITOR
|
||||
- GIT_SEQUENCE_EDITOR
|
||||
- DISPLAY
|
||||
- DBUS_SESSION_BUS_ADDRESS
|
||||
- PCFS
|
||||
- SSH_AUTH_SOCK
|
||||
- XDG_SESSION_TYPE
|
||||
- XDG_RUNTIME_DIR
|
||||
- WAYLAND_DISPLAY
|
||||
network:
|
||||
unrestricted: true
|
||||
|
|
@ -301,16 +301,19 @@ func runServer(sockPath string) {
|
|||
<-sigChan
|
||||
|
||||
fmt.Println("shutting down")
|
||||
daemonCancel() // signal all sessions via context propagation
|
||||
if srv != nil {
|
||||
srv.Shutdown()
|
||||
}
|
||||
|
||||
// Stop accepting new connections first.
|
||||
if listener != nil {
|
||||
listener.Close() //nolint:errcheck
|
||||
}
|
||||
if tcpListener != nil {
|
||||
tcpListener.Close() //nolint:errcheck
|
||||
}
|
||||
|
||||
daemonCancel() // signal all sessions via context propagation
|
||||
if srv != nil {
|
||||
srv.Shutdown()
|
||||
}
|
||||
if true {
|
||||
os.Remove(sockPath)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -75,6 +75,11 @@ type Server struct {
|
|||
invalidateModels func()
|
||||
toolPrompt func(string) string
|
||||
toolRegistry *toolsrv.Registry
|
||||
serveWG sync.WaitGroup
|
||||
shutdownCtx context.Context
|
||||
shutdownCancel context.CancelFunc
|
||||
connMu sync.Mutex
|
||||
activeConns map[net.Conn]struct{}
|
||||
}
|
||||
|
||||
// Config holds the pre-built trees and manager for the server.
|
||||
|
|
@ -94,6 +99,7 @@ type Option func(*Server)
|
|||
|
||||
// New creates a new Server from a pre-built Config.
|
||||
func New(cfg Config) *Server {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
s := &Server{
|
||||
log: cfg.Sink.Logger("9p", olog.LevelDebug),
|
||||
sink: cfg.Sink,
|
||||
|
|
@ -103,6 +109,9 @@ func New(cfg Config) *Server {
|
|||
invalidateModels: cfg.InvalidateModels,
|
||||
toolPrompt: cfg.ToolPrompt,
|
||||
toolRegistry: cfg.ToolRegistry,
|
||||
shutdownCtx: ctx,
|
||||
shutdownCancel: cancel,
|
||||
activeConns: make(map[net.Conn]struct{}),
|
||||
}
|
||||
if cfg.ElevateBroker != nil {
|
||||
s.elevateTree = newElevateTree(cfg.ElevateBroker)
|
||||
|
|
@ -338,7 +347,19 @@ const readTimeout = 10 * time.Second
|
|||
// goroutine so blocking reads (e.g. *wait files) do not stall the serve loop.
|
||||
func (s *Server) Serve(conn net.Conn) {
|
||||
defer conn.Close()
|
||||
connCtx, connCancel := context.WithCancel(context.Background())
|
||||
s.serveWG.Add(1)
|
||||
defer s.serveWG.Done()
|
||||
|
||||
s.connMu.Lock()
|
||||
s.activeConns[conn] = struct{}{}
|
||||
s.connMu.Unlock()
|
||||
defer func() {
|
||||
s.connMu.Lock()
|
||||
delete(s.activeConns, conn)
|
||||
s.connMu.Unlock()
|
||||
}()
|
||||
|
||||
connCtx, connCancel := context.WithCancel(s.shutdownCtx)
|
||||
cs := &connState{
|
||||
fids: make(map[uint32]*fid),
|
||||
ctx: connCtx,
|
||||
|
|
@ -400,9 +421,13 @@ func (s *Server) Serve(conn net.Conn) {
|
|||
close(responses)
|
||||
}
|
||||
|
||||
func (s *Server) handle(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan9.Fcall {
|
||||
func (s *Server) handle(cs *connState, fc *plan9.Fcall, ctx context.Context) (resp *plan9.Fcall) {
|
||||
start := time.Now()
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
s.log.Error("panic in handle: %v", r)
|
||||
resp = errFcall(fc, "internal error")
|
||||
}
|
||||
s.log.Debug(" -> %s %dµs", fcallTypeName(fc.Type), time.Since(start).Microseconds())
|
||||
}()
|
||||
switch fc.Type {
|
||||
|
|
@ -1153,7 +1178,31 @@ func (s *Server) handleWrite(path, input, uname string) error {
|
|||
|
||||
// Shutdown kills all active sessions and batch jobs.
|
||||
func (s *Server) Shutdown() {
|
||||
// Cancel the shutdown context — all connection contexts derive from it,
|
||||
// so in-flight request handlers will see context cancellation.
|
||||
s.shutdownCancel()
|
||||
|
||||
// Close all active connections — unblocks ReadFcall in Serve goroutines.
|
||||
s.connMu.Lock()
|
||||
for conn := range s.activeConns {
|
||||
conn.Close()
|
||||
}
|
||||
s.connMu.Unlock()
|
||||
|
||||
// Kill sessions.
|
||||
fs.Shutdown(s.sessionTree)
|
||||
|
||||
// Wait for all Serve goroutines to finish, with timeout.
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
s.serveWG.Wait()
|
||||
close(done)
|
||||
}()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(3 * time.Second):
|
||||
s.log.Warn("shutdown: timed out waiting for Serve goroutines")
|
||||
}
|
||||
}
|
||||
|
||||
// InterruptAll cancels any in-progress agent turn on every active fs.
|
||||
|
|
@ -1273,10 +1322,16 @@ 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" {
|
||||
// TODO: look up the correct AgentState for this agent
|
||||
// length, vers := as.LogInfo()
|
||||
dir.Length = 4096
|
||||
// dir.Qid.Vers = vers
|
||||
if sess.AgentLog != nil {
|
||||
length, vers := sess.AgentLog.LogInfo()
|
||||
if length > 64*1024 {
|
||||
length = 64 * 1024
|
||||
}
|
||||
dir.Length = uint64(length)
|
||||
dir.Qid.Vers = vers
|
||||
} else {
|
||||
dir.Length = 4096
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -85,6 +85,7 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
var core *coresession.Session
|
||||
var sessPtr *Session
|
||||
var proc *toolsrv.Process
|
||||
var al *AgentLog
|
||||
uname := rs.nextUname()
|
||||
if rs.cfg.NewCore != nil {
|
||||
var err error
|
||||
|
|
@ -92,6 +93,7 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
al = NewAgentLog(remoteTarget)
|
||||
} else {
|
||||
cfg := LoadAgentConfig(rs.cfg.AgentsDir, agentName, nil)
|
||||
if cfg != nil {
|
||||
|
|
@ -213,6 +215,7 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
toolSrv := newToolServer()
|
||||
rt := agent.BuildRuntime(cfg, toolSrv, cwd, env, sysPrompt, opModel, envBlock)
|
||||
|
||||
al = NewAgentLog(remoteTarget)
|
||||
core = coresession.New(coresession.Config{
|
||||
Backend: be,
|
||||
AgentName: agentName,
|
||||
|
|
@ -227,12 +230,16 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
PromptEnvExtra: promptEnv,
|
||||
BaseLayers: []string{sysPrompt, opModel, envBlock},
|
||||
Log: rs.cfg.Sink.NewLogger("core"),
|
||||
Output: NewEventHandler(al),
|
||||
ReadPlanStep: func() string {
|
||||
if sessPtr == nil {
|
||||
if sessPtr == nil || sessPtr.AgentLog == nil {
|
||||
return ""
|
||||
}
|
||||
// TODO: look up the correct AgentState for the active agent
|
||||
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)
|
||||
},
|
||||
})
|
||||
}
|
||||
|
|
@ -240,7 +247,7 @@ func Create(rs *rootState, args []string) (string, error) {
|
|||
sessionCtx, sessionCancel := context.WithCancel(rs.cfg.Ctx)
|
||||
sess := NewSession(sessID, core, sessionCtx, sessionCancel)
|
||||
sessPtr = sess
|
||||
// remote is now on AgentState; set via NewAgentState
|
||||
sess.AgentLog = al
|
||||
sess.proc = proc
|
||||
|
||||
rs.mu.Lock()
|
||||
|
|
|
|||
|
|
@ -32,6 +32,7 @@ var AgentFileList = []struct {
|
|||
}{
|
||||
{"cwd", false, false},
|
||||
{"plan", false, false},
|
||||
{"ctl", false, true},
|
||||
{"prompt", false, true},
|
||||
{"fifo.in", false, true},
|
||||
{"fifo.out", true, false},
|
||||
|
|
@ -75,8 +76,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, 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}
|
||||
func NewAgentTree(sess *Session, al *AgentLog, log *olog.Logger, saveTranscript func([]byte) error, invalidateModels func(), toolRegistry *toolsrv.Registry) *Tree {
|
||||
h := &agentHelper{al: al, sess: sess, log: log, saveTranscript: saveTranscript, invalidateModels: invalidateModels, toolRegistry: toolRegistry}
|
||||
perms := sessionFilePerms()
|
||||
specs := make([]FileSpec, len(AgentFileList))
|
||||
for i, f := range AgentFileList {
|
||||
|
|
@ -100,7 +101,7 @@ type sessionHelper struct {
|
|||
|
||||
// agentHelper holds the dependencies needed to build agent FileSpecs.
|
||||
type agentHelper struct {
|
||||
as *AgentState
|
||||
al *AgentLog
|
||||
sess *Session
|
||||
log *olog.Logger
|
||||
saveTranscript func([]byte) error
|
||||
|
|
@ -115,6 +116,7 @@ func (h *sessionHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
case "env":
|
||||
fs.Read = func() ([]byte, error) { return []byte(h.envContent()), nil }
|
||||
case "ctl":
|
||||
fs.Read = func() ([]byte, error) { return nil, nil }
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
|
|
@ -175,47 +177,49 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
fs.Read = func() ([]byte, error) {
|
||||
return []byte(h.sess.Core.CWD() + "\n"), nil
|
||||
}
|
||||
case "ctl":
|
||||
fs.Read = func() ([]byte, error) { return nil, nil }
|
||||
case "plan":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.as.mu.RLock()
|
||||
data := make([]byte, len(h.as.plan))
|
||||
copy(data, h.as.plan)
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
data := make([]byte, len(h.al.plan))
|
||||
copy(data, h.al.plan)
|
||||
h.al.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
fs.Size = func() int64 {
|
||||
h.as.mu.RLock()
|
||||
n := len(h.as.plan)
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
n := len(h.al.plan)
|
||||
h.al.mu.RUnlock()
|
||||
return int64(n)
|
||||
}
|
||||
case "prompt.prev":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.as.mu.RLock()
|
||||
data := make([]byte, len(h.as.prevPrompt))
|
||||
copy(data, h.as.prevPrompt)
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
data := make([]byte, len(h.al.prevPrompt))
|
||||
copy(data, h.al.prevPrompt)
|
||||
h.al.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
case "log":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
const maxWindow = 64 * 1024
|
||||
h.as.mu.RLock()
|
||||
log := h.as.log
|
||||
h.al.mu.RLock()
|
||||
log := h.al.log
|
||||
start := 0
|
||||
if len(log) > maxWindow {
|
||||
start = len(log) - maxWindow
|
||||
}
|
||||
data := make([]byte, len(log)-start)
|
||||
copy(data, log[start:])
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
fs.Size = func() int64 {
|
||||
const maxWindow = 64 * 1024
|
||||
h.as.mu.RLock()
|
||||
n := len(h.as.log)
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
n := len(h.al.log)
|
||||
h.al.mu.RUnlock()
|
||||
if n > maxWindow {
|
||||
return int64(maxWindow)
|
||||
}
|
||||
|
|
@ -247,10 +251,10 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
}
|
||||
case "plan":
|
||||
fs.Write = func(data []byte) error {
|
||||
h.as.mu.Lock()
|
||||
h.as.plan = make([]byte, len(data))
|
||||
copy(h.as.plan, data)
|
||||
h.as.mu.Unlock()
|
||||
h.al.mu.Lock()
|
||||
h.al.plan = make([]byte, len(data))
|
||||
copy(h.al.plan, data)
|
||||
h.al.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
case "log":
|
||||
|
|
@ -275,12 +279,12 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
h.sess.InvalidateModelsCache()
|
||||
return nil
|
||||
}
|
||||
h.as.mu.Lock()
|
||||
h.as.prevPrompt = []byte(input)
|
||||
h.as.mu.Unlock()
|
||||
h.al.mu.Lock()
|
||||
h.al.prevPrompt = []byte(input)
|
||||
h.al.mu.Unlock()
|
||||
go func() {
|
||||
h.sess.Core.Agent().Submit(h.sess.SessionCtx, input)
|
||||
h.as.EnsureTrailingNewline()
|
||||
h.al.EnsureTrailingNewline()
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
|
@ -301,6 +305,14 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
}
|
||||
return h.handleCfg(input)
|
||||
}
|
||||
case "ctl":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
return h.handleCtl(input)
|
||||
}
|
||||
case "tools":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
|
|
@ -340,32 +352,34 @@ func (h *agentHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
ctx, cancel := context.WithCancel(connCtx)
|
||||
defer cancel()
|
||||
context.AfterFunc(h.sess.SessionCtx, cancel)
|
||||
// Wake up on context cancellation (e.g. shutdown).
|
||||
context.AfterFunc(ctx, func() { h.al.chatCond.Broadcast() })
|
||||
|
||||
// Parse offset from base; empty = start from current end ("now")
|
||||
var offset int
|
||||
if base != "" {
|
||||
fmt.Sscanf(base, "%d", &offset)
|
||||
} else {
|
||||
h.as.mu.RLock()
|
||||
offset = len(h.as.log)
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
offset = len(h.al.log)
|
||||
h.al.mu.RUnlock()
|
||||
}
|
||||
|
||||
// Block until new data is available or context cancelled
|
||||
h.as.mu.RLock()
|
||||
for len(h.as.log) <= offset {
|
||||
h.al.mu.RLock()
|
||||
for len(h.al.log) <= offset {
|
||||
if ctx.Err() != nil {
|
||||
h.as.mu.RUnlock()
|
||||
h.al.mu.RUnlock()
|
||||
return nil, "", nil
|
||||
}
|
||||
h.as.chatCond.Wait()
|
||||
h.al.chatCond.Wait()
|
||||
}
|
||||
|
||||
// Copy new data
|
||||
data := make([]byte, len(h.as.log)-offset)
|
||||
copy(data, h.as.log[offset:])
|
||||
newOffset := len(h.as.log)
|
||||
h.as.mu.RUnlock()
|
||||
data := make([]byte, len(h.al.log)-offset)
|
||||
copy(data, h.al.log[offset:])
|
||||
newOffset := len(h.al.log)
|
||||
h.al.mu.RUnlock()
|
||||
|
||||
return data, fmt.Sprintf("%d", newOffset), nil
|
||||
}
|
||||
|
|
@ -385,8 +399,8 @@ func (h *sessionHelper) content(name string) string {
|
|||
}
|
||||
|
||||
func (h *agentHelper) content(name string) string {
|
||||
h.as.mu.RLock()
|
||||
defer h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
defer h.al.mu.RUnlock()
|
||||
switch name {
|
||||
case "usage":
|
||||
return h.sess.Core.Agent().UsageStr() + "\n"
|
||||
|
|
@ -402,7 +416,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.as.ChatOffset)
|
||||
return fmt.Sprintf("%d\n", h.al.ChatOffset)
|
||||
case "tail":
|
||||
return "#!/bin/sh\nexec tail -f \"$(dirname \"$0\")/chat\"\n"
|
||||
case "context":
|
||||
|
|
@ -453,8 +467,8 @@ func (h *agentHelper) content(name string) string {
|
|||
}
|
||||
|
||||
func (h *agentHelper) cfgContent() string {
|
||||
h.as.mu.RLock()
|
||||
defer h.as.mu.RUnlock()
|
||||
h.al.mu.RLock()
|
||||
defer h.al.mu.RUnlock()
|
||||
p := h.sess.Core.Agent().GenParams()
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "name=%s\n", h.sess.id)
|
||||
|
|
@ -462,7 +476,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.as.remote)
|
||||
fmt.Fprintf(&sb, "remote=%s\n", h.al.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(as *AgentState, messages []backend.Message) {
|
||||
func replayMessagesToLog(as *AgentLog, messages []backend.Message) {
|
||||
const maxReplay = 20
|
||||
start := len(messages) - maxReplay
|
||||
if start < 0 {
|
||||
|
|
|
|||
|
|
@ -310,6 +310,7 @@ func restoreSession(rs *rootState, ps *agent.PersistedAgent) error {
|
|||
restoredSession := agent.RestoreHistory(ps)
|
||||
|
||||
var sessPtr *Session
|
||||
al := NewAgentLog(remoteTarget)
|
||||
core := coresession.New(coresession.Config{
|
||||
Backend: be,
|
||||
AgentName: agentName,
|
||||
|
|
@ -325,24 +326,27 @@ 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 {
|
||||
return ""
|
||||
}
|
||||
// TODO: look up the correct AgentState for the active agent
|
||||
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)
|
||||
},
|
||||
})
|
||||
|
||||
sessionCtx, sessionCancel := context.WithCancel(rs.cfg.Ctx)
|
||||
sess := NewSession(sessID, core, sessionCtx, sessionCancel)
|
||||
sessPtr = sess
|
||||
sess.uname = uname
|
||||
// sess.remote is now on AgentState; set via NewAgentState above
|
||||
sess.AgentLog = al
|
||||
sess.proc = proc
|
||||
|
||||
as := NewAgentState(remoteTarget)
|
||||
replayMessagesToLog(as, ps.Messages)
|
||||
replayMessagesToLog(sess.AgentLog, ps.Messages)
|
||||
|
||||
rs.mu.Lock()
|
||||
rs.sessions[sessID] = sess
|
||||
|
|
|
|||
|
|
@ -170,7 +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
|
||||
sess.AgentLog,
|
||||
rs.cfg.Log,
|
||||
nil,
|
||||
rs.cfg.InvalidateModels,
|
||||
|
|
|
|||
|
|
@ -12,7 +12,7 @@ import (
|
|||
)
|
||||
|
||||
// Session holds session-level state exposed via 9P.
|
||||
// Agent-scoped state (log, plan, etc.) lives in AgentState.
|
||||
// Agent-scoped state (log, plan, etc.) lives in AgentLog.
|
||||
type Session struct {
|
||||
mu sync.RWMutex
|
||||
id string
|
||||
|
|
@ -22,14 +22,17 @@ 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
|
||||
|
||||
// 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 {
|
||||
// AgentLog holds all per-agent state exposed via 9P.
|
||||
type AgentLog struct {
|
||||
mu sync.RWMutex
|
||||
log []byte
|
||||
logVers uint32
|
||||
|
|
@ -43,10 +46,10 @@ type AgentState struct {
|
|||
chatCond *sync.Cond
|
||||
}
|
||||
|
||||
func NewAgentState(remote string) *AgentState {
|
||||
as := &AgentState{remote: remote}
|
||||
as.chatCond = sync.NewCond(as.mu.RLocker())
|
||||
return as
|
||||
func NewAgentLog(remote string) *AgentLog {
|
||||
al := &AgentLog{remote: remote}
|
||||
al.chatCond = sync.NewCond(al.mu.RLocker())
|
||||
return al
|
||||
}
|
||||
|
||||
func NewSession(id string, core *session.Session, ctx context.Context, cancel context.CancelFunc) *Session {
|
||||
|
|
@ -98,10 +101,10 @@ func (sess *Session) Interrupt() {
|
|||
sess.Core.Agent().Interrupt(agent.ErrInterrupted)
|
||||
}
|
||||
|
||||
// --- AgentState methods ---
|
||||
// --- AgentLog methods ---
|
||||
|
||||
// AppendLog appends data to the agent's log and bumps the version.
|
||||
func (as *AgentState) AppendLog(data []byte) {
|
||||
func (as *AgentLog) AppendLog(data []byte) {
|
||||
if len(data) == 0 {
|
||||
return
|
||||
}
|
||||
|
|
@ -113,7 +116,7 @@ func (as *AgentState) AppendLog(data []byte) {
|
|||
}
|
||||
|
||||
// EnsureTrailingNewline appends a newline if the log doesn't already end with one.
|
||||
func (as *AgentState) EnsureTrailingNewline() {
|
||||
func (as *AgentLog) EnsureTrailingNewline() {
|
||||
as.mu.Lock()
|
||||
if len(as.log) > 0 && as.log[len(as.log)-1] != '\n' {
|
||||
as.log = append(as.log, '\n')
|
||||
|
|
@ -123,26 +126,105 @@ func (as *AgentState) EnsureTrailingNewline() {
|
|||
}
|
||||
|
||||
// LogInfo returns the current log length and version atomically.
|
||||
func (as *AgentState) LogInfo() (length int, vers uint32) {
|
||||
func (as *AgentLog) LogInfo() (length int, vers uint32) {
|
||||
as.mu.RLock()
|
||||
defer as.mu.RUnlock()
|
||||
return len(as.log), as.logVers
|
||||
}
|
||||
|
||||
// Log returns the raw chat log bytes. Caller must hold Mu().RLock().
|
||||
func (as *AgentState) Log() []byte { return as.log }
|
||||
func (as *AgentLog) Log() []byte { return as.log }
|
||||
|
||||
// Plan returns the plan bytes. Caller must hold Mu().RLock().
|
||||
func (as *AgentState) Plan() []byte { return as.plan }
|
||||
func (as *AgentLog) Plan() []byte { return as.plan }
|
||||
|
||||
// SetPlan sets the plan. Caller must hold Mu().Lock().
|
||||
func (as *AgentState) SetPlan(p []byte) { as.plan = p }
|
||||
func (as *AgentLog) SetPlan(p []byte) { as.plan = p }
|
||||
|
||||
// PrevPrompt returns the last submitted prompt. Caller must hold Mu().RLock().
|
||||
func (as *AgentState) PrevPrompt() []byte { return as.prevPrompt }
|
||||
func (as *AgentLog) PrevPrompt() []byte { return as.prevPrompt }
|
||||
|
||||
// Mu returns the agent state's mutex.
|
||||
func (as *AgentState) Mu() *sync.RWMutex { return &as.mu }
|
||||
func (as *AgentLog) Mu() *sync.RWMutex { return &as.mu }
|
||||
|
||||
// Remote returns the SSH target for remote execution.
|
||||
func (as *AgentState) Remote() string { return as.remote }
|
||||
func (as *AgentLog) Remote() string { return as.remote }
|
||||
|
||||
// NewEventHandler returns an agent.EventHandler that writes formatted events
|
||||
// to this AgentLog. Use as the Output callback when creating an agent.
|
||||
func NewEventHandler(as *AgentLog) agent.EventHandler {
|
||||
streamingRole := ""
|
||||
streamingResponseID := ""
|
||||
|
||||
return func(ev agent.Event) {
|
||||
switch ev.Role {
|
||||
case "assistant", "reasoning":
|
||||
if streamingRole != ev.Role || (ev.Role == "assistant" && ev.ResponseID != streamingResponseID) {
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n"))
|
||||
}
|
||||
if ev.Role == "assistant" && ev.ResponseID != "" {
|
||||
as.AppendLog([]byte("[assistant:" + ev.ResponseID + "]\n"))
|
||||
streamingResponseID = ev.ResponseID
|
||||
} else {
|
||||
as.AppendLog([]byte("[" + ev.Role + "]\n"))
|
||||
}
|
||||
if ev.Role == "assistant" {
|
||||
as.mu.Lock()
|
||||
as.ChatOffset = len(as.log)
|
||||
as.mu.Unlock()
|
||||
}
|
||||
streamingRole = ev.Role
|
||||
}
|
||||
as.AppendLog(FormatEvent(ev))
|
||||
|
||||
case "tool":
|
||||
if streamingRole == "tool" {
|
||||
if ev.Content == "" {
|
||||
as.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
} else {
|
||||
as.AppendLog([]byte(ev.Content))
|
||||
}
|
||||
return
|
||||
}
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
if ev.Content == "" {
|
||||
as.AppendLog([]byte("[tool:" + ev.Name + "]\n"))
|
||||
streamingRole = "tool"
|
||||
} else {
|
||||
as.AppendLog(FormatEvent(ev))
|
||||
}
|
||||
|
||||
case "retry":
|
||||
if streamingRole == "retry" {
|
||||
as.AppendLog([]byte("\n" + ev.Content))
|
||||
return
|
||||
}
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
as.AppendLog([]byte("[retry]\n"))
|
||||
streamingRole = "retry"
|
||||
as.AppendLog(FormatEvent(ev))
|
||||
|
||||
case "call":
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
as.AppendLog(FormatEvent(ev))
|
||||
|
||||
default:
|
||||
if streamingRole != "" {
|
||||
as.AppendLog([]byte("\n"))
|
||||
streamingRole = ""
|
||||
}
|
||||
as.AppendLog(FormatEvent(ev))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -12,7 +12,6 @@ import (
|
|||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/simonfxr/pubsub"
|
||||
"ollie/agent"
|
||||
"ollie/backend"
|
||||
olog "ollie/log"
|
||||
|
|
@ -41,13 +40,13 @@ type Config struct {
|
|||
PromptEnvExtra []string
|
||||
Remote string
|
||||
BaseLayers []string
|
||||
Output agent.EventHandler
|
||||
}
|
||||
|
||||
// Session is the concrete session type. It owns session-level state and
|
||||
// delegates agent operations to its owned Agent.
|
||||
type Session struct {
|
||||
id string
|
||||
bus *pubsub.Bus
|
||||
envMu sync.RWMutex
|
||||
env map[string]string
|
||||
plan []byte
|
||||
|
|
@ -127,12 +126,10 @@ 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: bus,
|
||||
env: make(map[string]string),
|
||||
log: log,
|
||||
auditLog: auditLog,
|
||||
|
|
@ -152,7 +149,7 @@ func New(cfg Config) *Session {
|
|||
PromptEnvExtra: cfg.PromptEnvExtra,
|
||||
NewToolServer: cfg.NewToolServer,
|
||||
NewBackend: cfg.NewBackend,
|
||||
Bus: bus,
|
||||
Output: cfg.Output,
|
||||
Log: log,
|
||||
AuditLog: auditLog,
|
||||
SessionID: cfg.SessionID,
|
||||
|
|
@ -185,7 +182,6 @@ func (a *Session) SetEnv(key, value string) {
|
|||
}
|
||||
|
||||
func (a *Session) Agent() *agent.Agent { return a.r }
|
||||
func (a *Session) Bus() *pubsub.Bus { return a.bus }
|
||||
|
||||
// CWD returns the current working directory for tool execution.
|
||||
func (a *Session) CWD() string {
|
||||
|
|
|
|||
Loading…
Reference in New Issue