restore active sessions from data directory
This commit is contained in:
parent
cbd249e5f9
commit
1a545be58f
36
main.go
36
main.go
|
|
@ -135,12 +135,6 @@ func cmdMount() {
|
|||
fmt.Printf("mounted %s at %s (pid %d)\n", addr, mnt, cmd.Process.Pid)
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
func runServer(sockPath string) {
|
||||
env.EnsureDefaults()
|
||||
|
||||
|
|
@ -162,21 +156,21 @@ func runServer(sockPath string) {
|
|||
sink := olog.NewSink(os.Stdout, os.Stderr, olog.ParseLevel(os.Getenv("OLLIE_LOG"), olog.LevelWarn))
|
||||
|
||||
agentsDirs := agent.AgentsDirs()
|
||||
sessionsDir := paths.CfgDir() + "/sessions"
|
||||
sessionsDir := paths.DataDir() + "/sessions"
|
||||
|
||||
// D-Bus adapter (initialized after manager so callbacks can reference it).
|
||||
var dbusAdapter *DBusAdapter
|
||||
|
||||
mgr := session.NewManager(session.ManagerConfig{
|
||||
AgentsDir: agentsDirs[0],
|
||||
SessionsDir: sessionsDir,
|
||||
Log: sink.NewLogger("9p"),
|
||||
Sink: sink,
|
||||
Strict: *strict,
|
||||
Yolo: *yolo,
|
||||
NoMount: *no9p || *tcpAddr != "",
|
||||
Enable9P: !*no9p,
|
||||
EnableDBus: !*nodbus && (*no9p || *tcpAddr == ""),
|
||||
AgentsDir: agentsDirs[0],
|
||||
SessionsDir: sessionsDir,
|
||||
Log: sink.NewLogger("9p"),
|
||||
Sink: sink,
|
||||
Strict: *strict,
|
||||
Yolo: *yolo,
|
||||
NoMount: *no9p || *tcpAddr != "",
|
||||
Enable9P: !*no9p,
|
||||
EnableDBus: !*nodbus && (*no9p || *tcpAddr == ""),
|
||||
InvalidateModels: modelCache.Invalidate,
|
||||
OnSessionCreated: func(id string, sess *session.Session) {
|
||||
if dbusAdapter != nil {
|
||||
|
|
@ -198,9 +192,9 @@ func runServer(sockPath string) {
|
|||
var srv *server.Server
|
||||
if !*no9p {
|
||||
srv = server.New(server.Config{
|
||||
Sink: sink,
|
||||
SessionMgr: mgr,
|
||||
RootStore: NewRootStore(),
|
||||
Sink: sink,
|
||||
SessionMgr: mgr,
|
||||
RootStore: NewRootStore(),
|
||||
InvalidateModels: modelCache.Invalidate,
|
||||
ToolPrompt: execute.ToolPrompt,
|
||||
})
|
||||
|
|
@ -273,9 +267,6 @@ func runServer(sockPath string) {
|
|||
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
|
||||
<-sigChan
|
||||
|
||||
if srv != nil {
|
||||
srv.InterruptAll()
|
||||
}
|
||||
fmt.Println("shutting down")
|
||||
if dbusAdapter != nil {
|
||||
dbusAdapter.Close()
|
||||
|
|
@ -295,7 +286,6 @@ func runServer(sockPath string) {
|
|||
sink.Flush()
|
||||
}
|
||||
|
||||
|
||||
// ToolIndex generates a tool index from a tree's file listing.
|
||||
// NewRootStore returns a read-only FileTree for synthetic root-level files.
|
||||
func NewRootStore() *fs.Tree {
|
||||
|
|
|
|||
|
|
@ -37,10 +37,10 @@ type Session struct {
|
|||
logVers uint32
|
||||
ChatOffset int
|
||||
plan []byte
|
||||
prevPrompt []byte // last submitted prompt; overwritten on each new submission
|
||||
prevPrompt []byte // last submitted prompt; overwritten on each new submission
|
||||
peers map[string]bool // peer session IDs (bidirectional links)
|
||||
remote string // SSH target for remote execution (empty = local)
|
||||
mount *mountState // non-nil for remote sessions that need a local mount
|
||||
remote string // SSH target for remote execution (empty = local)
|
||||
mount *mountState // non-nil for remote sessions that need a local mount
|
||||
|
||||
// Cached model list (expensive API call; refreshed every 24h).
|
||||
modelsMu sync.Mutex
|
||||
|
|
@ -94,7 +94,7 @@ func (sess *Session) HasPeer(peerID string) bool {
|
|||
}
|
||||
|
||||
func (sess *Session) RunnableID() string { return sess.id }
|
||||
func (sess *Session) Uname() string { return sess.uname }
|
||||
func (sess *Session) Uname() string { return sess.uname }
|
||||
|
||||
func (sess *Session) Cancel() {
|
||||
sess.cancel()
|
||||
|
|
@ -212,8 +212,6 @@ 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 }
|
||||
|
||||
|
||||
|
||||
var sessionStoreOrder = []string{"new", "idx", "ls", "kill", "sh", "b", "bfg", "bbg", "cleanup"}
|
||||
|
||||
// FileMode returns the mode for a fixed session file,
|
||||
|
|
@ -258,14 +256,15 @@ type ManagerConfig struct {
|
|||
// Manager manages session lifecycle and exposes sessions as a Tree.
|
||||
type Manager struct {
|
||||
tree *fs.Tree
|
||||
cfg ManagerConfig
|
||||
mu sync.RWMutex
|
||||
sessions map[string]*Session
|
||||
nextUID atomic.Uint32 // incrementing principal counter
|
||||
cfg ManagerConfig
|
||||
mu sync.RWMutex
|
||||
sessions map[string]*Session
|
||||
nextUID atomic.Uint32 // incrementing principal counter
|
||||
}
|
||||
|
||||
// Tree returns the Tree view of the session namespace.
|
||||
func (s *Manager) Tree() *fs.Tree { return s.tree }
|
||||
|
||||
// nextUname generates the next uname atomically.
|
||||
func (s *Manager) nextUname() string {
|
||||
return fmt.Sprintf("%d", s.nextUID.Add(1))
|
||||
|
|
@ -569,7 +568,9 @@ func (s *Manager) openEntry(name string) (fs.File, error) {
|
|||
switch name {
|
||||
case "new":
|
||||
return &fs.FileConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &fs.SyntheticFileInfo{Name_: "new", Mode_: fs.Perms[fs.PathSessions].Files["new"]}, nil },
|
||||
StatFn: func() (os.FileInfo, error) {
|
||||
return &fs.SyntheticFileInfo{Name_: "new", Mode_: fs.Perms[fs.PathSessions].Files["new"]}, nil
|
||||
},
|
||||
ReadFn: func() ([]byte, error) {
|
||||
return []byte("name=\ncwd=\nremote=\nbackend=\nmodel=\nagent=\nmaxTokens=\nmaxCompletionTokens=\ntemperature=\ntopP=\ntopK=\nminP=\ntopA=\nfrequencyPenalty=\npresencePenalty=\nrepetitionPenalty=\nreasoning=\nreasoningEffort=\nincludeReasoning=\nresponseFormat=\nstop=\nverbosity=\n"), nil
|
||||
},
|
||||
|
|
@ -581,7 +582,9 @@ func (s *Manager) openEntry(name string) (fs.File, error) {
|
|||
}, nil
|
||||
case "idx":
|
||||
return &fs.FileConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &fs.SyntheticFileInfo{Name_: "idx", Mode_: fs.Perms[fs.PathSessions].Files["idx"]}, nil },
|
||||
StatFn: func() (os.FileInfo, error) {
|
||||
return &fs.SyntheticFileInfo{Name_: "idx", Mode_: fs.Perms[fs.PathSessions].Files["idx"]}, nil
|
||||
},
|
||||
ReadFn: func() ([]byte, error) { return s.index(), nil },
|
||||
WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") },
|
||||
BlockingReadFn: notBlocking,
|
||||
|
|
@ -589,7 +592,9 @@ func (s *Manager) openEntry(name string) (fs.File, error) {
|
|||
}
|
||||
if _, ok := fs.Perms[fs.PathSessions].Files[name]; ok {
|
||||
return &fs.FileConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &fs.SyntheticFileInfo{Name_: name, Mode_: fs.Perms[fs.PathSessions].Files[name]}, nil },
|
||||
StatFn: func() (os.FileInfo, error) {
|
||||
return &fs.SyntheticFileInfo{Name_: name, Mode_: fs.Perms[fs.PathSessions].Files[name]}, nil
|
||||
},
|
||||
ReadFn: func() ([]byte, error) {
|
||||
return s.cfg.ReadFile(paths.CfgDir() + "/scripts/s/" + name)
|
||||
},
|
||||
|
|
@ -713,7 +718,7 @@ func (s *Manager) InterruptAll() {
|
|||
// --- Session Persistence ---
|
||||
|
||||
func (s *Manager) activeSessionsDir() string {
|
||||
return filepath.Join(paths.DataDir(), "sessions", "active")
|
||||
return filepath.Join(s.cfg.SessionsDir, "active")
|
||||
}
|
||||
|
||||
func (s *Manager) persistSession(id string) {
|
||||
|
|
@ -766,11 +771,10 @@ func (s *Manager) restoreAllSessions() {
|
|||
}
|
||||
if err := s.restoreSession(ps); err != nil {
|
||||
s.cfg.Log.Error("restore session %s: %v", ps.ID, err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
// Clean up persisted files after successful restore
|
||||
for _, e := range entries {
|
||||
os.Remove(filepath.Join(dir, e.Name()))
|
||||
// Keep active snapshots until the session is explicitly killed.
|
||||
// Shutdown refreshes these files with the latest state.
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -938,6 +942,10 @@ func (s *Manager) restoreSession(ps *agent.PersistedSession) error {
|
|||
}
|
||||
|
||||
func (s *Manager) Shutdown() {
|
||||
// Interrupt all in-progress turns and wait for them to finish
|
||||
// before persisting state, so we capture the latest messages.
|
||||
s.InterruptAll()
|
||||
s.waitIdle(100*time.Millisecond, 5*time.Second)
|
||||
s.saveAllSessions()
|
||||
s.mu.Lock()
|
||||
ids := make([]string, 0, len(s.sessions))
|
||||
|
|
@ -961,6 +969,27 @@ func (s *Manager) Shutdown() {
|
|||
}
|
||||
}
|
||||
|
||||
// waitIdle polls all sessions until they are idle or the timeout expires.
|
||||
func (s *Manager) waitIdle(poll, timeout time.Duration) {
|
||||
deadline := time.Now().Add(timeout)
|
||||
for time.Now().Before(deadline) {
|
||||
allIdle := true
|
||||
s.mu.RLock()
|
||||
for _, sess := range s.sessions {
|
||||
if sess.Core.State() != "idle" {
|
||||
allIdle = false
|
||||
break
|
||||
}
|
||||
}
|
||||
s.mu.RUnlock()
|
||||
if allIdle {
|
||||
return
|
||||
}
|
||||
time.Sleep(poll)
|
||||
}
|
||||
s.cfg.Log.Warn("shutdown: timed out waiting for sessions to become idle, saving anyway")
|
||||
}
|
||||
|
||||
func (s *Manager) KillSession(id string) {
|
||||
s.mu.Lock()
|
||||
sess := s.sessions[id]
|
||||
|
|
@ -1100,9 +1129,9 @@ func (s *Manager) CreateSession(args []string) (string, error) {
|
|||
}
|
||||
|
||||
var execOpts []execute.Option
|
||||
if !s.cfg.NoMount {
|
||||
execOpts = append(execOpts, WithMount9P())
|
||||
}
|
||||
if !s.cfg.NoMount {
|
||||
execOpts = append(execOpts, WithMount9P())
|
||||
}
|
||||
if s.cfg.Strict {
|
||||
execOpts = append(execOpts, execute.WithStrict())
|
||||
}
|
||||
|
|
@ -1424,7 +1453,6 @@ func (s *Manager) Complete(cwd, filePath, prefix, suffix, extraContext string) (
|
|||
return result, nil
|
||||
}
|
||||
|
||||
|
||||
func stripCompletionNoise(s string) string {
|
||||
var lines []string
|
||||
for _, line := range strings.Split(s, "\n") {
|
||||
|
|
|
|||
Reference in New Issue