store: use core Session/SessionStore/SessionFileStore, remove local implementations
Server no longer owns sessions map or session lifecycle methods. SessionStore and SessionFileStore are type aliases to core. Fid path rewriting on rename handled via OnRename callback.
This commit is contained in:
parent
b6d579d74c
commit
73b634b6d3
|
|
@ -49,8 +49,6 @@ import (
|
|||
olog "ollie/pkg/log"
|
||||
"ollie/pkg/paths"
|
||||
"ollie/pkg/store"
|
||||
"ollie/pkg/tools"
|
||||
"ollie/pkg/tools/execute"
|
||||
|
||||
"9fans.net/go/plan9"
|
||||
)
|
||||
|
|
@ -60,31 +58,6 @@ const (
|
|||
QTFile = plan9.QTFILE
|
||||
)
|
||||
|
||||
// session holds all state for one ollie agent session.
|
||||
type session struct {
|
||||
mu sync.RWMutex
|
||||
id string
|
||||
core agent.Core
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
// chatLog is an append-only record of the conversation. Reads are served
|
||||
// directly from this buffer by offset, so tail -f works via polling.
|
||||
chatLog []byte
|
||||
chatVers uint32 // incremented on each append; used as Qid.Vers
|
||||
chatOffset int // byte position of chat EOF immediately before last prompt submit
|
||||
}
|
||||
|
||||
// appendChat appends data to the session's chat log and bumps the Qid version.
|
||||
func (sess *session) appendChat(data []byte) {
|
||||
if len(data) == 0 {
|
||||
return
|
||||
}
|
||||
sess.mu.Lock()
|
||||
sess.chatLog = append(sess.chatLog, data...)
|
||||
sess.chatVers++
|
||||
sess.mu.Unlock()
|
||||
}
|
||||
|
||||
// fid tracks per-descriptor state for a single 9P connection.
|
||||
type fid struct {
|
||||
path string
|
||||
|
|
@ -106,12 +79,10 @@ type connState struct {
|
|||
// Server is the 9P server for ollie sessions.
|
||||
type Server struct {
|
||||
mu sync.RWMutex
|
||||
sessions map[string]*session
|
||||
conns []*connState
|
||||
log *olog.Logger
|
||||
sink *olog.Sink
|
||||
agentsDir string // kept for session creation wiring
|
||||
sessionsDir string
|
||||
agentStore Store
|
||||
promptStore ReadableStore
|
||||
memStore Store
|
||||
|
|
@ -119,7 +90,7 @@ type Server struct {
|
|||
utilStore Store
|
||||
pluginStore Store
|
||||
skillStore Store
|
||||
sessionStore Store
|
||||
sessionStore *SessionStore
|
||||
batchStore *BatchStore
|
||||
transcriptStore Store
|
||||
tmpStore Store
|
||||
|
|
@ -134,12 +105,11 @@ func New(sink *olog.Sink) *Server {
|
|||
tmpDir := defaultTmpDir()
|
||||
os.MkdirAll(tmpDir, 0755) //nolint:errcheck
|
||||
agentsDir := paths.CfgDir() + "/agents"
|
||||
sessionsDir := paths.CfgDir() + "/sessions"
|
||||
s := &Server{
|
||||
sessions: make(map[string]*session),
|
||||
log: sink.Logger("9p", olog.LevelDebug),
|
||||
sink: sink,
|
||||
agentsDir: agentsDir,
|
||||
sessionsDir: paths.CfgDir() + "/sessions",
|
||||
agentStore: NewFlatDirStore(agentsDir, 0644),
|
||||
promptStore: NewFlatDirStore(agent.DefaultPromptsDir(), 0444),
|
||||
memStore: NewFlatDirStore(memDir, 0644),
|
||||
|
|
@ -150,7 +120,26 @@ func New(sink *olog.Sink) *Server {
|
|||
transcriptStore: NewFlatDirStore(transcriptDir, 0444),
|
||||
tmpStore: NewFlatDirStore(tmpDir, 0600),
|
||||
}
|
||||
s.sessionStore = &SessionStore{srv: s}
|
||||
s.sessionStore = store.NewSessionStore(store.SessionStoreConfig{
|
||||
AgentsDir: agentsDir,
|
||||
SessionsDir: sessionsDir,
|
||||
Log: s.log,
|
||||
Sink: s.sink,
|
||||
OnRename: func(oldID, newID string) {
|
||||
oldPrefix := "/s/" + oldID
|
||||
newPrefix := "/s/" + newID
|
||||
for _, c := range s.conns {
|
||||
c.mu.Lock()
|
||||
for _, f := range c.fids {
|
||||
if f.path == oldPrefix || strings.HasPrefix(f.path, oldPrefix+"/") {
|
||||
f.path = newPrefix + f.path[len(oldPrefix):]
|
||||
f.qid.Path = qidPath(f.path)
|
||||
}
|
||||
}
|
||||
c.mu.Unlock()
|
||||
}
|
||||
},
|
||||
})
|
||||
s.batchStore = store.NewBatchStore(store.BatchStoreConfig{
|
||||
AgentsDir: agentsDir,
|
||||
Log: s.log,
|
||||
|
|
@ -186,17 +175,15 @@ func defaultTmpDir() string {
|
|||
|
||||
// sessionFileStore returns a SessionFileStore for the given session ID.
|
||||
func (s *Server) sessionFileStore(sessID string) (*SessionFileStore, bool) {
|
||||
s.mu.RLock()
|
||||
sess, ok := s.sessions[sessID]
|
||||
s.mu.RUnlock()
|
||||
if !ok {
|
||||
sess := s.sessionStore.Session(sessID)
|
||||
if sess == nil {
|
||||
return nil, false
|
||||
}
|
||||
return NewSessionFileStore(
|
||||
return store.NewSessionFileStore(
|
||||
sess,
|
||||
s.log,
|
||||
func() { s.killSession(sessID) },
|
||||
func(newID string) error { return s.renameSession(sessID, newID) },
|
||||
func() { s.sessionStore.KillSession(sessID) },
|
||||
func(newID string) error { return s.sessionStore.Rename(sessID, newID) },
|
||||
func(data []byte) error {
|
||||
name := time.Now().Format("20060102T150405") + "-chat.md"
|
||||
return s.transcriptStore.Put(name, data)
|
||||
|
|
@ -321,7 +308,7 @@ func isSessionStoreFile(path string) bool {
|
|||
if !ok || strings.Contains(name, "/") {
|
||||
return false
|
||||
}
|
||||
_, ok = sessionStoreFiles[name]
|
||||
_, ok = store.SessionStoreFileMode(name)
|
||||
return ok
|
||||
}
|
||||
|
||||
|
|
@ -975,51 +962,6 @@ func (s *Server) wstat(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
return &plan9.Fcall{Type: plan9.Rwstat, Tag: fc.Tag}
|
||||
}
|
||||
|
||||
// renameSession re-keys a session in the sessions map, updates the core's
|
||||
// session ID, and rewrites fid paths across all connections.
|
||||
// Caller must hold NO locks on s.mu (this method acquires it).
|
||||
func (s *Server) renameSession(oldID, newID string) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
sess, ok := s.sessions[oldID]
|
||||
if !ok {
|
||||
return fmt.Errorf("session not found: %s", oldID)
|
||||
}
|
||||
if _, exists := s.sessions[newID]; exists {
|
||||
return fmt.Errorf("session already exists: %s", newID)
|
||||
}
|
||||
if sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot rename while agent is running")
|
||||
}
|
||||
|
||||
if err := sess.core.SetSessionID(newID); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
sess.id = newID
|
||||
s.sessions[newID] = sess
|
||||
delete(s.sessions, oldID)
|
||||
|
||||
// Rewrite fid paths across all connections.
|
||||
oldPrefix := "/s/" + oldID
|
||||
newPrefix := "/s/" + newID
|
||||
for _, c := range s.conns {
|
||||
c.mu.Lock()
|
||||
for _, f := range c.fids {
|
||||
if f.path == oldPrefix || strings.HasPrefix(f.path, oldPrefix+"/") {
|
||||
f.path = newPrefix + f.path[len(oldPrefix):]
|
||||
f.qid.Path = qidPath(f.path)
|
||||
}
|
||||
}
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
sess.appendChat([]byte(fmt.Sprintf("(session renamed: %s -> %s)\n", oldID, newID)))
|
||||
s.log.Info("renamed session %s -> %s", oldID, newID)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Server) clunk(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
||||
cs.mu.Lock()
|
||||
f, ok := cs.fids[fc.Fid]
|
||||
|
|
@ -1175,158 +1117,15 @@ func (s *Server) handleWrite(path, input string) error {
|
|||
}
|
||||
|
||||
// handleNewSession parses KV pairs (one per line or space-separated) and creates a session.
|
||||
func (s *Server) handleNewSession(input string) error {
|
||||
s.log.Debug("handleNewSession input=%q", input)
|
||||
// Accept both newline-separated and space-separated KV pairs.
|
||||
args := strings.Fields(input)
|
||||
return s.createSession(args)
|
||||
}
|
||||
|
||||
// createSession parses key=value options and starts a new agent session.
|
||||
// Recognised keys: name, backend, model, agent, cwd. cwd is required.
|
||||
// Example: new cwd=/home/lkn/src/myproject backend=ollama model=qwen3:8b agent=myagent
|
||||
func (s *Server) createSession(args []string) error {
|
||||
name := ""
|
||||
backendOverride := ""
|
||||
modelOverride := ""
|
||||
agentName := "default"
|
||||
cwd := ""
|
||||
for _, arg := range args {
|
||||
k, v, ok := strings.Cut(arg, "=")
|
||||
if !ok {
|
||||
return fmt.Errorf("invalid option %q (expected key=value)", arg)
|
||||
}
|
||||
if v == "" {
|
||||
continue
|
||||
}
|
||||
switch k {
|
||||
case "name":
|
||||
name = v
|
||||
case "backend":
|
||||
backendOverride = v
|
||||
case "model":
|
||||
modelOverride = v
|
||||
case "agent":
|
||||
agentName = v
|
||||
case "cwd":
|
||||
cwd = v
|
||||
default:
|
||||
return fmt.Errorf("unknown option %q (valid: name, backend, model, agent, cwd)", k)
|
||||
}
|
||||
}
|
||||
|
||||
if cwd == "" {
|
||||
return fmt.Errorf("cwd is required (e.g. new cwd=/path/to/project)")
|
||||
}
|
||||
|
||||
var (
|
||||
be backend.Backend
|
||||
err error
|
||||
)
|
||||
if backendOverride != "" {
|
||||
old := os.Getenv("OLLIE_BACKEND")
|
||||
os.Setenv("OLLIE_BACKEND", backendOverride) //nolint:errcheck
|
||||
be, err = backend.New()
|
||||
os.Setenv("OLLIE_BACKEND", old) //nolint:errcheck
|
||||
} else {
|
||||
be, err = backend.New()
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("backend: %w", err)
|
||||
}
|
||||
|
||||
if modelOverride != "" {
|
||||
be.SetModel(modelOverride)
|
||||
}
|
||||
|
||||
if err := os.MkdirAll(s.sessionsDir, 0700); err != nil {
|
||||
return fmt.Errorf("sessions dir: %w", err)
|
||||
}
|
||||
|
||||
cfg := store.LoadAgentConfig(s.agentsDir, agentName)
|
||||
|
||||
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
|
||||
"execute": execute.Decl(cwd),
|
||||
})
|
||||
|
||||
sessID := name
|
||||
if sessID == "" {
|
||||
sessID = agent.NewSessionID()
|
||||
}
|
||||
|
||||
s.mu.RLock()
|
||||
_, exists := s.sessions[sessID]
|
||||
s.mu.RUnlock()
|
||||
if exists {
|
||||
return fmt.Errorf("session already exists: %s", sessID)
|
||||
}
|
||||
|
||||
env := agent.BuildAgentEnv(cfg, newDisp(), cwd)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
core := agent.NewAgentCore(agent.AgentCoreConfig{
|
||||
Backend: be,
|
||||
AgentName: agentName,
|
||||
AgentsDir: s.agentsDir,
|
||||
SessionsDir: s.sessionsDir,
|
||||
SessionID: sessID,
|
||||
CWD: cwd,
|
||||
Env: env,
|
||||
NewDispatcher: newDisp,
|
||||
Log: s.sink.NewLogger("core"),
|
||||
})
|
||||
sess := &session{
|
||||
id: sessID,
|
||||
core: core,
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
s.sessions[sessID] = sess
|
||||
s.mu.Unlock()
|
||||
|
||||
s.log.Info("new session %s (backend=%s model=%s agent=%s)",
|
||||
sessID, core.BackendName(), core.ModelName(), core.AgentName())
|
||||
return nil
|
||||
}
|
||||
|
||||
// Shutdown kills all active sessions and batch jobs, triggering Close() on each core.
|
||||
// InterruptAll cancels any in-progress agent turn on every active session.
|
||||
// Call this on SIGINT or SIGTERM before Shutdown to allow turns to unwind
|
||||
// cleanly rather than having their context cancelled mid-tool-call.
|
||||
func (s *Server) InterruptAll() {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
for _, sess := range s.sessions {
|
||||
sess.core.Interrupt(agent.ErrInterrupted)
|
||||
}
|
||||
}
|
||||
|
||||
// Shutdown kills all active sessions and batch jobs.
|
||||
func (s *Server) Shutdown() {
|
||||
s.mu.Lock()
|
||||
ids := make([]string, 0, len(s.sessions))
|
||||
for id := range s.sessions {
|
||||
ids = append(ids, id)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
for _, id := range ids {
|
||||
s.killSession(id)
|
||||
}
|
||||
|
||||
s.sessionStore.Shutdown()
|
||||
s.batchStore.Shutdown()
|
||||
}
|
||||
|
||||
func (s *Server) killSession(id string) {
|
||||
s.mu.Lock()
|
||||
sess := s.sessions[id]
|
||||
delete(s.sessions, id)
|
||||
s.mu.Unlock()
|
||||
if sess != nil {
|
||||
sess.cancel()
|
||||
sess.core.Close()
|
||||
s.log.Info("killed session %s", id)
|
||||
}
|
||||
// InterruptAll cancels any in-progress agent turn on every active session.
|
||||
func (s *Server) InterruptAll() {
|
||||
s.sessionStore.InterruptAll()
|
||||
}
|
||||
|
||||
// readDir serializes directory entries for the given path, respecting offset and count.
|
||||
|
|
@ -1545,7 +1344,8 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
if path == "/backends" || path == "/help" {
|
||||
mode = 0444
|
||||
} else if isSessionStoreFile(path) {
|
||||
mode = plan9.Perm(sessionStoreFiles[base])
|
||||
mode_, _ := store.SessionStoreFileMode(base)
|
||||
mode = plan9.Perm(mode_)
|
||||
} else if strings.HasPrefix(path, "/a/") {
|
||||
mode = 0666
|
||||
} else if strings.HasPrefix(path, "/p/") {
|
||||
|
|
@ -1587,15 +1387,11 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
if strings.HasPrefix(path, "/s/") {
|
||||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
|
||||
if len(parts) == 3 && parts[0] == "s" {
|
||||
s.mu.RLock()
|
||||
sess := s.sessions[parts[1]]
|
||||
s.mu.RUnlock()
|
||||
if sess != nil {
|
||||
if sess := s.sessionStore.Session(parts[1]); sess != nil {
|
||||
if base == "chat" {
|
||||
sess.mu.RLock()
|
||||
dir.Length = uint64(len(sess.chatLog))
|
||||
dir.Qid.Vers = sess.chatVers
|
||||
sess.mu.RUnlock()
|
||||
length, vers := sess.ChatInfo()
|
||||
dir.Length = uint64(length)
|
||||
dir.Qid.Vers = vers
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,376 +0,0 @@
|
|||
package p9
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"ollie/pkg/agent"
|
||||
"ollie/pkg/backend"
|
||||
olog "ollie/pkg/log"
|
||||
"ollie/pkg/store"
|
||||
)
|
||||
|
||||
// sessionFileList defines the fixed set of files in a session directory,
|
||||
// with their 9P permission modes.
|
||||
var sessionFileList = []struct {
|
||||
name string
|
||||
mode os.FileMode
|
||||
}{
|
||||
{"ctl", 0200},
|
||||
{"prompt", 0200},
|
||||
{"fifo.in", 0200},
|
||||
{"fifo.out", 0444},
|
||||
{"chat", 0666},
|
||||
{"offset", 0444},
|
||||
{"state", 0444},
|
||||
{"statewait", 0444},
|
||||
{"backend", 0666},
|
||||
{"agent", 0666},
|
||||
{"model", 0666},
|
||||
{"cwd", 0666},
|
||||
{"cwdwait", 0444},
|
||||
{"usage", 0444},
|
||||
{"usagewait", 0444},
|
||||
{"ctxsz", 0444},
|
||||
{"ctxszwait", 0444},
|
||||
{"models", 0444},
|
||||
{"systemprompt", 0444},
|
||||
{"params", 0666},
|
||||
}
|
||||
|
||||
// SessionFileStore implements ReadWriteStore for the files within a single
|
||||
// session directory (/s/{id}/*). The file set is fixed; Create, Delete, and
|
||||
// Rename are not meaningful and are not part of the interface.
|
||||
type SessionFileStore struct {
|
||||
sess *session
|
||||
log *olog.Logger
|
||||
kill func()
|
||||
rename func(newID string) error
|
||||
saveTranscript func([]byte) error
|
||||
}
|
||||
|
||||
func NewSessionFileStore(sess *session, log *olog.Logger, kill func(), rename func(newID string) error, saveTranscript func([]byte) error) *SessionFileStore {
|
||||
return &SessionFileStore{sess: sess, log: log, kill: kill, rename: rename, saveTranscript: saveTranscript}
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) List() ([]os.DirEntry, error) {
|
||||
entries := make([]os.DirEntry, len(sessionFileList))
|
||||
for i, f := range sessionFileList {
|
||||
entries[i] = syntheticEntry(f.name, f.mode)
|
||||
}
|
||||
return entries, nil
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) Stat(name string) (os.FileInfo, error) {
|
||||
for _, f := range sessionFileList {
|
||||
if f.name == name {
|
||||
var size int64
|
||||
switch name {
|
||||
case "chat":
|
||||
s.sess.mu.RLock()
|
||||
size = int64(len(s.sess.chatLog))
|
||||
s.sess.mu.RUnlock()
|
||||
case "statewait", "usagewait", "ctxszwait", "cwdwait":
|
||||
// Blocking reads; size is unknown until resolved.
|
||||
default:
|
||||
size = int64(len(s.content(name)))
|
||||
}
|
||||
return &syntheticFileInfo{Name_: name, Mode_: f.mode, Size_: size}, nil
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) Get(name string) ([]byte, error) {
|
||||
switch name {
|
||||
case "chat":
|
||||
s.sess.mu.RLock()
|
||||
data := make([]byte, len(s.sess.chatLog))
|
||||
copy(data, s.sess.chatLog)
|
||||
s.sess.mu.RUnlock()
|
||||
return data, nil
|
||||
case "offset":
|
||||
return []byte(s.content("offset")), nil
|
||||
case "fifo.out":
|
||||
item, ok := s.sess.core.PopQueue()
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(item), nil
|
||||
default:
|
||||
for _, f := range sessionFileList {
|
||||
if f.name == name {
|
||||
return []byte(s.content(name)), nil
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) Put(name string, data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
switch name {
|
||||
case "chat":
|
||||
return s.saveTranscript([]byte(input))
|
||||
|
||||
case "prompt":
|
||||
s.sess.core.Submit(s.sess.ctx, input, s.makePublish())
|
||||
|
||||
case "fifo.in":
|
||||
s.sess.core.Queue(input)
|
||||
|
||||
case "ctl":
|
||||
return s.handleCtl(input)
|
||||
|
||||
case "backend":
|
||||
if s.sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch backend while agent is running")
|
||||
}
|
||||
s.sess.core.Submit(s.sess.ctx, "/backend "+input, s.makePublish())
|
||||
|
||||
case "agent":
|
||||
if s.sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch agent while agent is running")
|
||||
}
|
||||
s.sess.core.Submit(s.sess.ctx, "/agent "+input, s.makePublish())
|
||||
|
||||
case "model":
|
||||
if s.sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch model while agent is running")
|
||||
}
|
||||
s.sess.core.Submit(s.sess.ctx, "/model "+input, s.makePublish())
|
||||
|
||||
case "cwd":
|
||||
if err := s.sess.core.SetCWD(input); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
case "params":
|
||||
if s.sess.core.IsRunning() {
|
||||
return fmt.Errorf("cannot change params while agent is running")
|
||||
}
|
||||
params, err := parseParams(input, s.sess.core.GenerationParams())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return s.sess.core.SetGenerationParams(params)
|
||||
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func formatParams(p backend.GenerationParams) string {
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "maxTokens=%d\n", p.MaxTokens)
|
||||
if p.Temperature != nil {
|
||||
fmt.Fprintf(&sb, "temperature=%g\n", *p.Temperature)
|
||||
} else {
|
||||
fmt.Fprintf(&sb, "temperature=\n")
|
||||
}
|
||||
if p.FrequencyPenalty != nil {
|
||||
fmt.Fprintf(&sb, "frequencyPenalty=%g\n", *p.FrequencyPenalty)
|
||||
} else {
|
||||
fmt.Fprintf(&sb, "frequencyPenalty=\n")
|
||||
}
|
||||
if p.PresencePenalty != nil {
|
||||
fmt.Fprintf(&sb, "presencePenalty=%g\n", *p.PresencePenalty)
|
||||
} else {
|
||||
fmt.Fprintf(&sb, "presencePenalty=\n")
|
||||
}
|
||||
return sb.String()
|
||||
}
|
||||
|
||||
func parseParams(input string, current backend.GenerationParams) (backend.GenerationParams, error) {
|
||||
p := current
|
||||
for _, line := range strings.Split(input, "\n") {
|
||||
k, v, ok := strings.Cut(line, "=")
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
k = strings.TrimSpace(k)
|
||||
v = strings.TrimSpace(v)
|
||||
switch k {
|
||||
case "maxTokens":
|
||||
if v == "" {
|
||||
p.MaxTokens = 0
|
||||
} else {
|
||||
n, err := strconv.Atoi(v)
|
||||
if err != nil {
|
||||
return p, fmt.Errorf("invalid maxTokens: %s", v)
|
||||
}
|
||||
p.MaxTokens = n
|
||||
}
|
||||
case "temperature":
|
||||
if v == "" {
|
||||
p.Temperature = nil
|
||||
} else {
|
||||
f, err := strconv.ParseFloat(v, 64)
|
||||
if err != nil {
|
||||
return p, fmt.Errorf("invalid temperature: %s", v)
|
||||
}
|
||||
p.Temperature = &f
|
||||
}
|
||||
case "frequencyPenalty":
|
||||
if v == "" {
|
||||
p.FrequencyPenalty = nil
|
||||
} else {
|
||||
f, err := strconv.ParseFloat(v, 64)
|
||||
if err != nil {
|
||||
return p, fmt.Errorf("invalid frequencyPenalty: %s", v)
|
||||
}
|
||||
p.FrequencyPenalty = &f
|
||||
}
|
||||
case "presencePenalty":
|
||||
if v == "" {
|
||||
p.PresencePenalty = nil
|
||||
} else {
|
||||
f, err := strconv.ParseFloat(v, 64)
|
||||
if err != nil {
|
||||
return p, fmt.Errorf("invalid presencePenalty: %s", v)
|
||||
}
|
||||
p.PresencePenalty = &f
|
||||
}
|
||||
}
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
// content returns the string content of a simple readable session file.
|
||||
func (s *SessionFileStore) content(name string) string {
|
||||
s.sess.mu.RLock()
|
||||
defer s.sess.mu.RUnlock()
|
||||
switch name {
|
||||
case "backend":
|
||||
return s.sess.core.BackendName() + "\n"
|
||||
case "agent":
|
||||
return s.sess.core.AgentName() + "\n"
|
||||
case "model":
|
||||
return s.sess.core.ModelName() + "\n"
|
||||
case "state":
|
||||
return s.sess.core.State() + "\n"
|
||||
case "cwd":
|
||||
return s.sess.core.CWD() + "\n"
|
||||
case "usage":
|
||||
return s.sess.core.Usage() + "\n"
|
||||
case "ctxsz":
|
||||
return s.sess.core.CtxSz() + "\n"
|
||||
case "models":
|
||||
return s.sess.core.ListModels() + "\n"
|
||||
case "systemprompt":
|
||||
return s.sess.core.SystemPrompt()
|
||||
case "offset":
|
||||
return fmt.Sprintf("%d\n", s.sess.chatOffset)
|
||||
case "params":
|
||||
return formatParams(s.sess.core.GenerationParams())
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// CurrentWaitValue returns the current value for the named *wait file.
|
||||
func (s *SessionFileStore) CurrentWaitValue(name string) string {
|
||||
switch name {
|
||||
case "statewait":
|
||||
return s.sess.core.State()
|
||||
case "usagewait":
|
||||
return s.sess.core.Usage()
|
||||
case "ctxszwait":
|
||||
return s.sess.core.CtxSz()
|
||||
case "cwdwait":
|
||||
return s.sess.core.CWD()
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Wait blocks until the named *wait file's underlying value changes from base,
|
||||
// then returns the new value. If base is empty, the current value is used.
|
||||
// Unblocks when connCtx or the session context is cancelled.
|
||||
func (s *SessionFileStore) Wait(connCtx context.Context, name, base string) ([]byte, error) {
|
||||
ctx, cancel := context.WithCancel(connCtx)
|
||||
defer cancel()
|
||||
// also unblock when the session itself is killed.
|
||||
context.AfterFunc(s.sess.ctx, cancel)
|
||||
|
||||
var field string
|
||||
switch name {
|
||||
case "statewait":
|
||||
field = agent.WatchState
|
||||
case "usagewait":
|
||||
field = agent.WatchUsage
|
||||
case "ctxszwait":
|
||||
field = agent.WatchCtxSz
|
||||
case "cwdwait":
|
||||
field = agent.WatchCWD
|
||||
default:
|
||||
return nil, fmt.Errorf("%s: not a wait file", name)
|
||||
}
|
||||
|
||||
if base == "" {
|
||||
base = s.CurrentWaitValue(name)
|
||||
}
|
||||
|
||||
v, ok := s.sess.core.WaitChange(ctx, field, base)
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) makePublish() func(agent.Event) {
|
||||
assistantStarted := false
|
||||
return func(ev agent.Event) {
|
||||
switch ev.Role {
|
||||
case "user", "call", "tool":
|
||||
assistantStarted = false
|
||||
case "assistant":
|
||||
if !assistantStarted {
|
||||
s.sess.appendChat([]byte("assistant: "))
|
||||
assistantStarted = true
|
||||
}
|
||||
}
|
||||
s.sess.appendChat(store.FormatEvent(ev))
|
||||
if ev.Role == "user" {
|
||||
s.sess.mu.Lock()
|
||||
s.sess.chatOffset = len(s.sess.chatLog)
|
||||
s.sess.mu.Unlock()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *SessionFileStore) handleCtl(input string) error {
|
||||
cmd := strings.Fields(input)
|
||||
if len(cmd) == 0 {
|
||||
return fmt.Errorf("empty ctl command")
|
||||
}
|
||||
switch cmd[0] {
|
||||
case "stop":
|
||||
s.sess.core.Interrupt(agent.ErrInterrupted)
|
||||
case "kill":
|
||||
s.kill()
|
||||
case "rn":
|
||||
if name := strings.TrimSpace(input[3:]); name != "" {
|
||||
if err := s.rename(name); err != nil {
|
||||
s.log.Error("rename: %v", err)
|
||||
}
|
||||
}
|
||||
case "save":
|
||||
s.sess.mu.RLock()
|
||||
data := make([]byte, len(s.sess.chatLog))
|
||||
copy(data, s.sess.chatLog)
|
||||
s.sess.mu.RUnlock()
|
||||
return s.saveTranscript(data)
|
||||
case "compact", "clear", "backend", "model", "models",
|
||||
"agents", "agent", "sessions", "cwd", "skills",
|
||||
"tools", "context", "usage", "history",
|
||||
"irw", "help":
|
||||
s.sess.core.Submit(s.sess.ctx, "/"+input, s.makePublish())
|
||||
default:
|
||||
return fmt.Errorf("unknown ctl command: %s", cmd[0])
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
@ -1,112 +0,0 @@
|
|||
package p9
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
"ollie/pkg/paths"
|
||||
)
|
||||
|
||||
// sessionStoreFiles maps fixed file entries in s/ to their permissions.
|
||||
// Add an entry here to expose a new file; the server uses this map for
|
||||
// routing, mode, and stat length.
|
||||
var sessionStoreFiles = map[string]os.FileMode{
|
||||
"new": 0666,
|
||||
"idx": 0444,
|
||||
"ls": 0555,
|
||||
"kill": 0555,
|
||||
"sh": 0555,
|
||||
}
|
||||
|
||||
// sessionStoreOrder defines the listing order for fixed s/ entries.
|
||||
var sessionStoreOrder = []string{"new", "idx", "ls", "kill", "sh"}
|
||||
|
||||
// SessionStore implements Store for the /s/ directory.
|
||||
// Entries are session IDs (directories) plus the synthetic files "new" and "idx".
|
||||
// Put("new", data) creates a session; Delete(id) kills one; Rename renames one.
|
||||
type SessionStore struct {
|
||||
srv *Server
|
||||
}
|
||||
|
||||
func (s *SessionStore) List() ([]os.DirEntry, error) {
|
||||
entries := make([]os.DirEntry, 0, len(sessionStoreOrder))
|
||||
for _, name := range sessionStoreOrder {
|
||||
entries = append(entries, syntheticEntry(name, sessionStoreFiles[name]))
|
||||
}
|
||||
s.srv.mu.RLock()
|
||||
for id := range s.srv.sessions {
|
||||
entries = append(entries, syntheticDirEntry(id, 0555))
|
||||
}
|
||||
s.srv.mu.RUnlock()
|
||||
return entries, nil
|
||||
}
|
||||
|
||||
func (s *SessionStore) Stat(name string) (os.FileInfo, error) {
|
||||
if mode, ok := sessionStoreFiles[name]; ok {
|
||||
return &syntheticFileInfo{Name_: name, Mode_: mode}, nil
|
||||
}
|
||||
s.srv.mu.RLock()
|
||||
_, ok := s.srv.sessions[name]
|
||||
s.srv.mu.RUnlock()
|
||||
if ok {
|
||||
return &syntheticFileInfo{Name_: name, Mode_: 0555, IsDir_: true}, nil
|
||||
}
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
|
||||
func (s *SessionStore) Get(name string) ([]byte, error) {
|
||||
switch name {
|
||||
case "new":
|
||||
return []byte("name=\ncwd=\nbackend=\nmodel=\nagent=\n"), nil
|
||||
case "idx":
|
||||
return s.index(), nil
|
||||
default:
|
||||
if _, ok := sessionStoreFiles[name]; ok {
|
||||
return os.ReadFile(paths.CfgDir() + "/scripts/s/" + name)
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("%s: not a readable file", name)
|
||||
}
|
||||
|
||||
func (s *SessionStore) Put(name string, data []byte) error {
|
||||
if name != "new" {
|
||||
return fmt.Errorf("%s: not writable", name)
|
||||
}
|
||||
return s.srv.handleNewSession(strings.TrimSpace(string(data)))
|
||||
}
|
||||
|
||||
func (s *SessionStore) Delete(name string) error {
|
||||
s.srv.mu.RLock()
|
||||
_, ok := s.srv.sessions[name]
|
||||
s.srv.mu.RUnlock()
|
||||
if !ok {
|
||||
return fmt.Errorf("session not found: %s", name)
|
||||
}
|
||||
s.srv.killSession(name)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *SessionStore) Create(name string) error {
|
||||
return fmt.Errorf("create not supported for sessions")
|
||||
}
|
||||
|
||||
func (s *SessionStore) Rename(oldName, newName string) error {
|
||||
return s.srv.renameSession(oldName, newName)
|
||||
}
|
||||
|
||||
func (s *SessionStore) index() []byte {
|
||||
var sb strings.Builder
|
||||
s.srv.mu.RLock()
|
||||
for id, sess := range s.srv.sessions {
|
||||
sess.mu.RLock()
|
||||
state := sess.core.State()
|
||||
cwd := sess.core.CWD()
|
||||
be := sess.core.BackendName()
|
||||
model := sess.core.ModelName()
|
||||
sess.mu.RUnlock()
|
||||
fmt.Fprintf(&sb, "%s\t%s\t%s\t%s\t%s\n", id, state, cwd, be, model)
|
||||
}
|
||||
s.srv.mu.RUnlock()
|
||||
return []byte(sb.String())
|
||||
}
|
||||
|
|
@ -21,6 +21,9 @@ type (
|
|||
SkillStore = store.SkillStore
|
||||
BatchStore = store.BatchStore
|
||||
BatchJobStore = store.BatchJobStore
|
||||
Session = store.Session
|
||||
SessionStore = store.SessionStore
|
||||
SessionFileStore = store.SessionFileStore
|
||||
|
||||
syntheticFileInfo = store.SyntheticFileInfo
|
||||
)
|
||||
|
|
|
|||
Reference in New Issue