move store package from core into 9p/store
This commit is contained in:
parent
4aa815b50c
commit
61c169af3b
|
|
@ -17,13 +17,13 @@ import (
|
|||
"ollie/pkg/backend"
|
||||
olog "ollie/pkg/log"
|
||||
"ollie/pkg/paths"
|
||||
"ollie/pkg/store"
|
||||
"ollie/pkg/tools/execute"
|
||||
"olliesrv/store"
|
||||
|
||||
"9fans.net/go/plan9"
|
||||
)
|
||||
|
||||
// Re-export core store types so existing 9p code compiles unchanged.
|
||||
// Re-export store types for use in server.go.
|
||||
type (
|
||||
StoreEntry = store.StoreEntry
|
||||
Store = store.Store
|
||||
|
|
|
|||
|
|
@ -0,0 +1,106 @@
|
|||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// MemoryStore is an in-memory key/value store implementing Store.
|
||||
// Values are []byte blobs held in RAM; nothing is persisted to disk.
|
||||
// All operations are safe for concurrent use.
|
||||
type MemoryStore struct {
|
||||
mu sync.RWMutex
|
||||
entries map[string][]byte
|
||||
}
|
||||
|
||||
func NewMemoryStore() *MemoryStore {
|
||||
return &MemoryStore{entries: make(map[string][]byte)}
|
||||
}
|
||||
|
||||
// Get returns the current value for key, or nil if absent.
|
||||
func (m *MemoryStore) Get(key string) []byte {
|
||||
m.mu.RLock()
|
||||
v := m.entries[key]
|
||||
m.mu.RUnlock()
|
||||
return v
|
||||
}
|
||||
|
||||
// Set overwrites the value for key.
|
||||
func (m *MemoryStore) Set(key string, value []byte) {
|
||||
m.mu.Lock()
|
||||
m.entries[key] = value
|
||||
m.mu.Unlock()
|
||||
}
|
||||
|
||||
// --- Store interface ---
|
||||
|
||||
func (m *MemoryStore) Stat(name string) (os.FileInfo, error) {
|
||||
m.mu.RLock()
|
||||
_, ok := m.entries[name]
|
||||
sz := int64(len(m.entries[name]))
|
||||
m.mu.RUnlock()
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
return &SyntheticFileInfo{Name_: name, Mode_: 0444, Size_: sz}, nil
|
||||
}
|
||||
|
||||
func (m *MemoryStore) List() ([]os.DirEntry, error) {
|
||||
m.mu.RLock()
|
||||
entries := make([]os.DirEntry, 0, len(m.entries))
|
||||
for name := range m.entries {
|
||||
entries = append(entries, FileEntry(name, 0444))
|
||||
}
|
||||
m.mu.RUnlock()
|
||||
return entries, nil
|
||||
}
|
||||
|
||||
func (m *MemoryStore) Open(name string) (StoreEntry, error) {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return m.Stat(name) },
|
||||
ReadFn: func() ([]byte, error) {
|
||||
m.mu.RLock()
|
||||
v := make([]byte, len(m.entries[name]))
|
||||
copy(v, m.entries[name])
|
||||
m.mu.RUnlock()
|
||||
return v, nil
|
||||
},
|
||||
WriteFn: func([]byte) error { return fmt.Errorf("%s: read-only", name) },
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (m *MemoryStore) Create(name string) error {
|
||||
m.mu.Lock()
|
||||
m.entries[name] = nil
|
||||
m.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *MemoryStore) Delete(name string) error {
|
||||
m.mu.Lock()
|
||||
_, ok := m.entries[name]
|
||||
delete(m.entries, name)
|
||||
m.mu.Unlock()
|
||||
if !ok {
|
||||
return fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *MemoryStore) Rename(oldName, newName string) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
v, ok := m.entries[oldName]
|
||||
if !ok {
|
||||
return fmt.Errorf("%s: not found", oldName)
|
||||
}
|
||||
m.entries[newName] = v
|
||||
delete(m.entries, oldName)
|
||||
return nil
|
||||
}
|
||||
|
|
@ -0,0 +1,105 @@
|
|||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
)
|
||||
|
||||
// FileSpec describes a single synthetic file exposed by a RunFileStore.
|
||||
type FileSpec struct {
|
||||
Name string
|
||||
Mode os.FileMode
|
||||
Read func() ([]byte, error)
|
||||
Write func([]byte) error // nil = read-only
|
||||
Wait func(ctx context.Context, base string) ([]byte, error) // nil = not waitable
|
||||
Size func() int64 // optional; if nil, len(Read())
|
||||
}
|
||||
|
||||
// RunFileStore implements RunnableStore for any Runnable using a table of FileSpecs.
|
||||
type RunFileStore struct {
|
||||
*storeConfig
|
||||
Runnable
|
||||
specs []FileSpec
|
||||
index map[string]int // name -> index into specs
|
||||
}
|
||||
|
||||
// NewRunFileStore creates a RunnableStore backed by the given Runnable and file table.
|
||||
func NewRunFileStore(r Runnable, specs []FileSpec) *RunFileStore {
|
||||
rs := &RunFileStore{
|
||||
Runnable: r,
|
||||
specs: specs,
|
||||
index: make(map[string]int, len(specs)),
|
||||
}
|
||||
for i, s := range specs {
|
||||
rs.index[s.Name] = i
|
||||
}
|
||||
notSupported := func(string) error { return fmt.Errorf("not supported") }
|
||||
rs.storeConfig = &storeConfig{
|
||||
StatFn: rs.stat,
|
||||
ListFn: rs.list,
|
||||
OpenFn: rs.open,
|
||||
DeleteFn: notSupported,
|
||||
CreateFn: notSupported,
|
||||
RenameFn: func(string, string) error { return fmt.Errorf("not supported") },
|
||||
}
|
||||
return rs
|
||||
}
|
||||
|
||||
func (rs *RunFileStore) lookup(name string) (*FileSpec, bool) {
|
||||
i, ok := rs.index[name]
|
||||
if !ok {
|
||||
return nil, false
|
||||
}
|
||||
return &rs.specs[i], true
|
||||
}
|
||||
|
||||
func (rs *RunFileStore) stat(name string) (os.FileInfo, error) {
|
||||
spec, ok := rs.lookup(name)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
var size int64
|
||||
if spec.Size != nil {
|
||||
size = spec.Size()
|
||||
} else if spec.Wait != nil {
|
||||
// Wait files block on read; report large size so clients read fully.
|
||||
size = 4096
|
||||
} else {
|
||||
if data, err := spec.Read(); err == nil {
|
||||
size = int64(len(data))
|
||||
}
|
||||
}
|
||||
return &SyntheticFileInfo{Name_: name, Mode_: spec.Mode, Size_: size}, nil
|
||||
}
|
||||
|
||||
func (rs *RunFileStore) list() ([]os.DirEntry, error) {
|
||||
entries := make([]os.DirEntry, len(rs.specs))
|
||||
for i, s := range rs.specs {
|
||||
entries[i] = FileEntry(s.Name, s.Mode)
|
||||
}
|
||||
return entries, nil
|
||||
}
|
||||
|
||||
func (rs *RunFileStore) open(name string) (StoreEntry, error) {
|
||||
spec, ok := rs.lookup(name)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
writeFn := func(data []byte) error { return fmt.Errorf("%s: read-only", name) }
|
||||
if spec.Write != nil {
|
||||
writeFn = spec.Write
|
||||
}
|
||||
waitFn := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("%s: not a wait file", name)
|
||||
}
|
||||
if spec.Wait != nil {
|
||||
waitFn = spec.Wait
|
||||
}
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return rs.stat(name) },
|
||||
ReadFn: spec.Read,
|
||||
WriteFn: writeFn,
|
||||
BlockingReadFn: waitFn,
|
||||
}, nil
|
||||
}
|
||||
|
|
@ -0,0 +1,481 @@
|
|||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"ollie/pkg/agent"
|
||||
"ollie/pkg/backend"
|
||||
"ollie/pkg/config"
|
||||
olog "ollie/pkg/log"
|
||||
"ollie/pkg/paths"
|
||||
"ollie/pkg/tools"
|
||||
"ollie/pkg/tools/execute"
|
||||
)
|
||||
|
||||
// Session holds all state for one agent session.
|
||||
type Session struct {
|
||||
mu sync.RWMutex
|
||||
id string
|
||||
Core agent.Core
|
||||
Ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
log []byte
|
||||
logVers uint32
|
||||
ChatOffset int
|
||||
plan []byte
|
||||
prevPrompt []byte // last submitted prompt; overwritten on each new submission
|
||||
}
|
||||
|
||||
func NewSession(id string, core agent.Core, ctx context.Context, cancel context.CancelFunc) *Session {
|
||||
return &Session{id: id, Core: core, Ctx: ctx, cancel: cancel}
|
||||
}
|
||||
|
||||
func (sess *Session) RunnableID() string { return sess.id }
|
||||
|
||||
func (sess *Session) Cancel() {
|
||||
sess.cancel()
|
||||
}
|
||||
|
||||
func (sess *Session) Interrupt() {
|
||||
sess.Core.Interrupt(agent.ErrInterrupted)
|
||||
}
|
||||
|
||||
// AppendLog appends data to the session's log and bumps the version.
|
||||
func (sess *Session) AppendLog(data []byte) {
|
||||
if len(data) == 0 {
|
||||
return
|
||||
}
|
||||
sess.mu.Lock()
|
||||
sess.log = append(sess.log, data...)
|
||||
sess.logVers++
|
||||
sess.mu.Unlock()
|
||||
}
|
||||
|
||||
// EnsureTrailingNewline appends a newline if the log doesn't already end with one.
|
||||
func (sess *Session) EnsureTrailingNewline() {
|
||||
sess.mu.Lock()
|
||||
if len(sess.log) > 0 && sess.log[len(sess.log)-1] != '\n' {
|
||||
sess.log = append(sess.log, '\n')
|
||||
sess.logVers++
|
||||
}
|
||||
sess.mu.Unlock()
|
||||
}
|
||||
|
||||
// LogInfo returns the current log length and version atomically.
|
||||
func (sess *Session) LogInfo() (length int, vers uint32) {
|
||||
sess.mu.RLock()
|
||||
defer sess.mu.RUnlock()
|
||||
return len(sess.log), sess.logVers
|
||||
}
|
||||
|
||||
// sessionStoreFiles maps fixed file entries in s/ to their permissions.
|
||||
var sessionStoreFiles = map[string]os.FileMode{
|
||||
"new": 0666,
|
||||
"idx": 0444,
|
||||
"ls": 0555,
|
||||
"kill": 0555,
|
||||
"sh": 0555,
|
||||
"b": 0555,
|
||||
"bfg": 0555,
|
||||
"bbg": 0555,
|
||||
"cleanup": 0555,
|
||||
}
|
||||
|
||||
var sessionStoreOrder = []string{"new", "idx", "ls", "kill", "sh", "b", "bfg", "bbg", "cleanup"}
|
||||
|
||||
// SessionStoreFileMode returns the mode for a fixed session store file,
|
||||
// or 0 and false if the name is not a fixed file.
|
||||
func SessionStoreFileMode(name string) (os.FileMode, bool) {
|
||||
m, ok := sessionStoreFiles[name]
|
||||
return m, ok
|
||||
}
|
||||
|
||||
// SessionStoreConfig holds the dependencies for a SessionStore.
|
||||
type SessionStoreConfig struct {
|
||||
AgentsDir string
|
||||
SessionsDir string
|
||||
Log *olog.Logger
|
||||
Sink *olog.Sink
|
||||
ReadFile func(string) ([]byte, error)
|
||||
MkdirAll func(string, os.FileMode) error
|
||||
// OnRename is called after a session is renamed in the map,
|
||||
// for protocol-level fixups (e.g. fid path rewriting).
|
||||
OnRename func(oldID, newID string)
|
||||
SaveTranscript func([]byte) error
|
||||
// NewCore, if non-nil, replaces the default backend.New + agent.NewAgentCore
|
||||
// path. It receives the session ID, agent name, and cwd, and returns a Core.
|
||||
NewCore func(sessionID, agentName, cwd string) (agent.Core, error)
|
||||
}
|
||||
|
||||
// SessionStore implements Store for session management.
|
||||
type SessionStore struct {
|
||||
*storeConfig
|
||||
cfg SessionStoreConfig
|
||||
mu sync.RWMutex
|
||||
sessions map[string]*Session
|
||||
}
|
||||
|
||||
func NewSessionStore(cfg SessionStoreConfig) *SessionStore {
|
||||
if cfg.ReadFile == nil {
|
||||
cfg.ReadFile = os.ReadFile
|
||||
}
|
||||
if cfg.MkdirAll == nil {
|
||||
cfg.MkdirAll = os.MkdirAll
|
||||
}
|
||||
ss := &SessionStore{
|
||||
cfg: cfg,
|
||||
sessions: make(map[string]*Session),
|
||||
}
|
||||
ss.storeConfig = &storeConfig{
|
||||
StatFn: ss.stat,
|
||||
ListFn: ss.list,
|
||||
OpenFn: ss.openEntry,
|
||||
DeleteFn: ss.del,
|
||||
CreateFn: func(string) error { return fmt.Errorf("create not supported for sessions") },
|
||||
RenameFn: ss.renameSession,
|
||||
}
|
||||
return ss
|
||||
}
|
||||
|
||||
// AddSession inserts a pre-built session into the store.
|
||||
func (s *SessionStore) AddSession(sess *Session) {
|
||||
s.mu.Lock()
|
||||
s.sessions[sess.RunnableID()] = sess
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
func (s *SessionStore) list() ([]os.DirEntry, error) {
|
||||
entries := make([]os.DirEntry, 0, len(sessionStoreOrder))
|
||||
for _, name := range sessionStoreOrder {
|
||||
entries = append(entries, FileEntry(name, sessionStoreFiles[name]))
|
||||
}
|
||||
s.mu.RLock()
|
||||
for id := range s.sessions {
|
||||
entries = append(entries, DirEntry(id, 0555))
|
||||
}
|
||||
s.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.mu.RLock()
|
||||
_, ok := s.sessions[name]
|
||||
s.mu.RUnlock()
|
||||
if ok {
|
||||
return &SyntheticFileInfo{Name_: name, Mode_: 0555, IsDir_: true}, nil
|
||||
}
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
|
||||
func (s *SessionStore) openEntry(name string) (StoreEntry, error) {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
switch name {
|
||||
case "new":
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &SyntheticFileInfo{Name_: "new", Mode_: 0666}, nil },
|
||||
ReadFn: func() ([]byte, error) { return []byte("name=\ncwd=\nbackend=\nmodel=\nagent=\n"), nil },
|
||||
WriteFn: func(data []byte) error {
|
||||
return s.createSession(strings.Fields(strings.TrimSpace(string(data))))
|
||||
},
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
case "idx":
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &SyntheticFileInfo{Name_: "idx", Mode_: 0444}, nil },
|
||||
ReadFn: func() ([]byte, error) { return s.index(), nil },
|
||||
WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") },
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
default:
|
||||
if _, ok := sessionStoreFiles[name]; ok {
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &SyntheticFileInfo{Name_: name, Mode_: sessionStoreFiles[name]}, nil },
|
||||
ReadFn: func() ([]byte, error) {
|
||||
return s.cfg.ReadFile(paths.CfgDir() + "/scripts/s/" + name)
|
||||
},
|
||||
WriteFn: func([]byte) error { return fmt.Errorf("%s: not writable", name) },
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
|
||||
func (s *SessionStore) del(name string) error {
|
||||
s.mu.RLock()
|
||||
_, ok := s.sessions[name]
|
||||
s.mu.RUnlock()
|
||||
if ok {
|
||||
s.KillSession(name)
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("session not found: %s", name)
|
||||
}
|
||||
|
||||
// Session returns the session for the given ID, or nil.
|
||||
func (s *SessionStore) Session(id string) *Session {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
return s.sessions[id]
|
||||
}
|
||||
|
||||
// OpenStore returns a RunnableStore for the given session ID.
|
||||
func (s *SessionStore) OpenStore(id string) (RunnableStore, error) {
|
||||
if sess := s.Session(id); sess != nil {
|
||||
return NewSessionFileStore(
|
||||
sess,
|
||||
s.cfg.Log,
|
||||
func() { s.KillSession(id) },
|
||||
func(newID string) error { return s.Rename(id, newID) },
|
||||
s.cfg.SaveTranscript,
|
||||
), nil
|
||||
}
|
||||
return nil, fmt.Errorf("session not found: %s", id)
|
||||
}
|
||||
|
||||
// InterruptAll interrupts every active session.
|
||||
func (s *SessionStore) InterruptAll() {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
for _, sess := range s.sessions {
|
||||
sess.Core.Interrupt(agent.ErrInterrupted)
|
||||
}
|
||||
}
|
||||
|
||||
// Shutdown kills all active sessions.
|
||||
func (s *SessionStore) 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)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *SessionStore) 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.cfg.Log.Info("killed session %s", id)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *SessionStore) index() []byte {
|
||||
var sb strings.Builder
|
||||
s.mu.RLock()
|
||||
ids := make([]string, 0, len(s.sessions))
|
||||
for id := range s.sessions {
|
||||
ids = append(ids, id)
|
||||
}
|
||||
sort.Strings(ids)
|
||||
for _, id := range ids {
|
||||
sess := s.sessions[id]
|
||||
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.mu.RUnlock()
|
||||
return []byte(sb.String())
|
||||
}
|
||||
|
||||
func (s *SessionStore) 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)")
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
|
||||
var core agent.Core
|
||||
if s.cfg.NewCore != nil {
|
||||
var err error
|
||||
core, err = s.cfg.NewCore(sessID, agentName, cwd)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
be, err := backend.NewWithName(backendOverride)
|
||||
if err != nil {
|
||||
return fmt.Errorf("backend: %w", err)
|
||||
}
|
||||
|
||||
if modelOverride == "" {
|
||||
modelOverride = os.Getenv("OLLIE_MODEL")
|
||||
}
|
||||
if modelOverride != "" {
|
||||
be.SetModel(modelOverride)
|
||||
}
|
||||
|
||||
if err := s.cfg.MkdirAll(s.cfg.SessionsDir, 0700); err != nil {
|
||||
return fmt.Errorf("sessions dir: %w", err)
|
||||
}
|
||||
|
||||
cfg := LoadAgentConfig(s.cfg.AgentsDir, agentName, nil)
|
||||
|
||||
newDisp := tools.NewDispatcherFunc(map[string]func() tools.Server{
|
||||
"execute": execute.Decl(cwd),
|
||||
})
|
||||
|
||||
env := agent.BuildAgentEnv(cfg, newDisp(), cwd)
|
||||
|
||||
core = agent.NewAgentCore(agent.AgentCoreConfig{
|
||||
Backend: be,
|
||||
AgentName: agentName,
|
||||
AgentsDir: s.cfg.AgentsDir,
|
||||
SessionsDir: s.cfg.SessionsDir,
|
||||
SessionID: sessID,
|
||||
CWD: cwd,
|
||||
Env: env,
|
||||
NewDispatcher: newDisp,
|
||||
Log: s.cfg.Sink.NewLogger("core"),
|
||||
})
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
sess := NewSession(sessID, core, ctx, cancel)
|
||||
|
||||
s.mu.Lock()
|
||||
s.sessions[sessID] = sess
|
||||
s.mu.Unlock()
|
||||
|
||||
s.cfg.Log.Info("new session %s (backend=%s model=%s agent=%s)",
|
||||
sessID, core.BackendName(), core.ModelName(), core.AgentName())
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *SessionStore) 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)
|
||||
|
||||
if s.cfg.OnRename != nil {
|
||||
s.cfg.OnRename(oldID, newID)
|
||||
}
|
||||
|
||||
sess.AppendLog([]byte(fmt.Sprintf(":: session renamed: %s -> %s\n", oldID, newID)))
|
||||
s.cfg.Log.Info("renamed session %s -> %s", oldID, newID)
|
||||
return nil
|
||||
}
|
||||
|
||||
// LoadAgentConfig resolves and loads the config for a named agent.
|
||||
// Returns nil if the config file does not exist; BuildAgentEnv handles nil configs.
|
||||
func LoadAgentConfig(agentsDir, name string, open func(string) (*os.File, error)) *config.Config {
|
||||
if open == nil {
|
||||
open = os.Open
|
||||
}
|
||||
f, err := open(agentsDir + "/" + name + ".json")
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
defer f.Close()
|
||||
cfg, _ := config.Load(f)
|
||||
return cfg
|
||||
}
|
||||
|
||||
// FormatEvent converts an agent Event to bytes for appending to a chat log.
|
||||
func FormatEvent(ev agent.Event) []byte {
|
||||
switch ev.Role {
|
||||
case "user":
|
||||
return []byte("user: " + ev.Content + "\n")
|
||||
case "assistant":
|
||||
return []byte(ev.Content)
|
||||
case "call":
|
||||
args := squashWhitespace(ev.Content)
|
||||
return []byte("-> " + ev.Name + "(" + args + ")\n")
|
||||
case "tool":
|
||||
return []byte(strings.TrimRight(ev.Content, "\n") + "\n")
|
||||
case "retry":
|
||||
return []byte("retrying in " + ev.Content + "s...\n")
|
||||
case "error":
|
||||
return []byte("error: " + ev.Content + "\n")
|
||||
case "stalled":
|
||||
return []byte("agent stalled\n")
|
||||
case "reasoning":
|
||||
return []byte(ev.Content)
|
||||
case "info":
|
||||
return []byte(":: " + ev.Content)
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
func squashWhitespace(s string) string {
|
||||
return strings.Join(strings.Fields(s), " ")
|
||||
}
|
||||
|
|
@ -0,0 +1,483 @@
|
|||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"ollie/pkg/agent"
|
||||
"ollie/pkg/backend"
|
||||
olog "ollie/pkg/log"
|
||||
)
|
||||
|
||||
// SessionFileList defines the fixed set of files in a session directory.
|
||||
var SessionFileList = []struct {
|
||||
Name string
|
||||
Mode os.FileMode
|
||||
}{
|
||||
{"plan", 0666},
|
||||
{"ctl", 0200},
|
||||
{"prompt", 0200},
|
||||
{"fifo.in", 0200},
|
||||
{"fifo.out", 0444},
|
||||
{"chat", 0666},
|
||||
{"offset", 0444},
|
||||
{"cfg", 0666},
|
||||
{"statewait", 0444},
|
||||
{"usage", 0444},
|
||||
{"cost", 0444},
|
||||
{"ctxsz", 0444},
|
||||
{"models", 0444},
|
||||
{"systemprompt", 0444},
|
||||
{"env", 0444},
|
||||
{"tail", 0555},
|
||||
{"prompt.prev", 0444},
|
||||
{"context", 0444},
|
||||
}
|
||||
|
||||
// SessionFileStore is a RunFileStore for a session directory.
|
||||
type SessionFileStore = RunFileStore
|
||||
|
||||
func NewSessionFileStore(sess *Session, log *olog.Logger, kill func(), rename func(newID string) error, saveTranscript func([]byte) error) *SessionFileStore {
|
||||
h := &sessionHelper{sess: sess, log: log, kill: kill, rename: rename, saveTranscript: saveTranscript}
|
||||
specs := make([]FileSpec, len(SessionFileList))
|
||||
for i, f := range SessionFileList {
|
||||
specs[i] = h.fileSpec(f.Name, f.Mode)
|
||||
}
|
||||
return NewRunFileStore(sess, specs)
|
||||
}
|
||||
|
||||
// sessionHelper holds the dependencies needed to build session FileSpecs.
|
||||
type sessionHelper struct {
|
||||
sess *Session
|
||||
log *olog.Logger
|
||||
kill func()
|
||||
rename func(newID string) error
|
||||
saveTranscript func([]byte) error
|
||||
}
|
||||
|
||||
func (h *sessionHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
||||
fs := FileSpec{Name: name, Mode: mode}
|
||||
|
||||
// Read
|
||||
switch name {
|
||||
case "plan":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.plan))
|
||||
copy(data, h.sess.plan)
|
||||
h.sess.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
fs.Size = func() int64 {
|
||||
h.sess.mu.RLock()
|
||||
n := len(h.sess.plan)
|
||||
h.sess.mu.RUnlock()
|
||||
return int64(n)
|
||||
}
|
||||
case "prompt.prev":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.prevPrompt))
|
||||
copy(data, h.sess.prevPrompt)
|
||||
h.sess.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
case "chat":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.log))
|
||||
copy(data, h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
return data, nil
|
||||
}
|
||||
fs.Size = func() int64 {
|
||||
h.sess.mu.RLock()
|
||||
n := len(h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
return int64(n)
|
||||
}
|
||||
case "fifo.out":
|
||||
fs.Read = func() ([]byte, error) {
|
||||
item, ok := h.sess.Core.PopQueue()
|
||||
if !ok {
|
||||
return nil, nil
|
||||
}
|
||||
return []byte(item), nil
|
||||
}
|
||||
case "cfg":
|
||||
fs.Read = func() ([]byte, error) { return []byte(h.cfgContent()), nil }
|
||||
default:
|
||||
fs.Read = func() ([]byte, error) { return []byte(h.content(name)), nil }
|
||||
}
|
||||
|
||||
// Write
|
||||
switch name {
|
||||
case "plan":
|
||||
fs.Write = func(data []byte) error {
|
||||
h.sess.mu.Lock()
|
||||
h.sess.plan = make([]byte, len(data))
|
||||
copy(h.sess.plan, data)
|
||||
h.sess.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
case "chat":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
return h.saveTranscript([]byte(input))
|
||||
}
|
||||
case "prompt":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
h.sess.mu.Lock()
|
||||
h.sess.prevPrompt = []byte(input)
|
||||
h.sess.mu.Unlock()
|
||||
pub := h.makePublish()
|
||||
go func() {
|
||||
h.sess.Core.Submit(h.sess.Ctx, input, pub)
|
||||
h.sess.EnsureTrailingNewline()
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
case "fifo.in":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
h.sess.Core.Queue(input)
|
||||
return nil
|
||||
}
|
||||
case "ctl":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
return h.handleCtl(input)
|
||||
}
|
||||
case "cfg":
|
||||
fs.Write = func(data []byte) error {
|
||||
input := strings.TrimSpace(string(data))
|
||||
if input == "" {
|
||||
return nil
|
||||
}
|
||||
return h.handleCfg(input)
|
||||
}
|
||||
}
|
||||
|
||||
// Wait
|
||||
if name == "statewait" {
|
||||
fs.Wait = func(connCtx context.Context, base string) ([]byte, error) {
|
||||
ctx, cancel := context.WithCancel(connCtx)
|
||||
defer cancel()
|
||||
context.AfterFunc(h.sess.Ctx, cancel)
|
||||
if base == "" {
|
||||
base = h.sess.Core.State()
|
||||
}
|
||||
v, ok := h.sess.Core.WaitChange(ctx, agent.WatchState, base)
|
||||
if !ok {
|
||||
// Timeout or cancel: return current state so caller can check.
|
||||
return []byte(h.sess.Core.State() + "\n"), nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
}
|
||||
}
|
||||
|
||||
return fs
|
||||
}
|
||||
|
||||
func (h *sessionHelper) content(name string) string {
|
||||
h.sess.mu.RLock()
|
||||
defer h.sess.mu.RUnlock()
|
||||
switch name {
|
||||
case "usage":
|
||||
return h.sess.Core.Usage() + "\n"
|
||||
case "cost":
|
||||
return h.sess.Core.Cost()
|
||||
case "ctxsz":
|
||||
return h.sess.Core.CtxSz() + "\n"
|
||||
case "models":
|
||||
return h.sess.Core.ListModels() + "\n"
|
||||
case "systemprompt":
|
||||
return h.sess.Core.SystemPrompt()
|
||||
case "offset":
|
||||
return fmt.Sprintf("%d\n", h.sess.ChatOffset)
|
||||
case "env":
|
||||
return h.envContent()
|
||||
case "tail":
|
||||
return "#!/bin/sh\nexec tail -f \"$(dirname \"$0\")/chat\"\n"
|
||||
case "context":
|
||||
var sb strings.Builder
|
||||
enc := json.NewEncoder(&sb)
|
||||
enc.SetEscapeHTML(false)
|
||||
for _, m := range h.sess.Core.Context() {
|
||||
enc.Encode(m)
|
||||
}
|
||||
return sb.String()
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (h *sessionHelper) envContent() string {
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "OLLIE_SESSION_ID=%s\n", h.sess.RunnableID())
|
||||
for _, e := range os.Environ() {
|
||||
if strings.HasPrefix(e, "OLLIE_") && !strings.HasPrefix(e, "OLLIE_SESSION_ID=") {
|
||||
sb.WriteString(e)
|
||||
sb.WriteByte('\n')
|
||||
}
|
||||
}
|
||||
return sb.String()
|
||||
}
|
||||
|
||||
|
||||
func (h *sessionHelper) cfgContent() string {
|
||||
h.sess.mu.RLock()
|
||||
defer h.sess.mu.RUnlock()
|
||||
p := h.sess.Core.GenerationParams()
|
||||
var sb strings.Builder
|
||||
fmt.Fprintf(&sb, "state=%s\n", h.sess.Core.State())
|
||||
fmt.Fprintf(&sb, "backend=%s\n", h.sess.Core.BackendName())
|
||||
fmt.Fprintf(&sb, "model=%s\n", h.sess.Core.ModelName())
|
||||
fmt.Fprintf(&sb, "agent=%s\n", h.sess.Core.AgentName())
|
||||
fmt.Fprintf(&sb, "cwd=%s\n", h.sess.Core.CWD())
|
||||
fmt.Fprintf(&sb, "maxTokens=%d\n", p.MaxTokens)
|
||||
if p.Temperature != nil {
|
||||
fmt.Fprintf(&sb, "temperature=%g\n", *p.Temperature)
|
||||
} else {
|
||||
sb.WriteString("temperature=\n")
|
||||
}
|
||||
if p.FrequencyPenalty != nil {
|
||||
fmt.Fprintf(&sb, "frequencyPenalty=%g\n", *p.FrequencyPenalty)
|
||||
} else {
|
||||
sb.WriteString("frequencyPenalty=\n")
|
||||
}
|
||||
if p.PresencePenalty != nil {
|
||||
fmt.Fprintf(&sb, "presencePenalty=%g\n", *p.PresencePenalty)
|
||||
} else {
|
||||
sb.WriteString("presencePenalty=\n")
|
||||
}
|
||||
return sb.String()
|
||||
}
|
||||
|
||||
func (h *sessionHelper) handleCfg(input string) error {
|
||||
p := h.sess.Core.GenerationParams()
|
||||
hasParams := false
|
||||
for _, line := range strings.Split(input, "\n") {
|
||||
line = strings.TrimSpace(line)
|
||||
if line == "" {
|
||||
continue
|
||||
}
|
||||
k, v, ok := strings.Cut(line, "=")
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
k, v = strings.TrimSpace(k), strings.TrimSpace(v)
|
||||
switch k {
|
||||
case "backend", "model", "agent":
|
||||
if v == "" {
|
||||
continue
|
||||
}
|
||||
if h.sess.Core.IsRunning() {
|
||||
return fmt.Errorf("cannot switch %s while agent is running", k)
|
||||
}
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/"+k+" "+v, h.makePublish())
|
||||
case "cwd":
|
||||
if v == "" {
|
||||
continue
|
||||
}
|
||||
if err := h.sess.Core.SetCWD(v); err != nil {
|
||||
return err
|
||||
}
|
||||
case "maxTokens":
|
||||
hasParams = true
|
||||
if v == "" {
|
||||
p.MaxTokens = 0
|
||||
} else if n, err := strconv.Atoi(v); err == nil {
|
||||
p.MaxTokens = n
|
||||
}
|
||||
case "temperature":
|
||||
hasParams = true
|
||||
if v == "" {
|
||||
p.Temperature = nil
|
||||
} else if f, err := strconv.ParseFloat(v, 64); err == nil {
|
||||
p.Temperature = &f
|
||||
}
|
||||
case "frequencyPenalty":
|
||||
hasParams = true
|
||||
if v == "" {
|
||||
p.FrequencyPenalty = nil
|
||||
} else if f, err := strconv.ParseFloat(v, 64); err == nil {
|
||||
p.FrequencyPenalty = &f
|
||||
}
|
||||
case "presencePenalty":
|
||||
hasParams = true
|
||||
if v == "" {
|
||||
p.PresencePenalty = nil
|
||||
} else if f, err := strconv.ParseFloat(v, 64); err == nil {
|
||||
p.PresencePenalty = &f
|
||||
}
|
||||
// state → read-only, silently ignored
|
||||
}
|
||||
}
|
||||
if hasParams {
|
||||
if h.sess.Core.IsRunning() {
|
||||
return fmt.Errorf("cannot change params while agent is running")
|
||||
}
|
||||
return h.sess.Core.SetGenerationParams(p)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
func (h *sessionHelper) makePublish() func(agent.Event) {
|
||||
assistantStarted := false
|
||||
return func(ev agent.Event) {
|
||||
if ev.Role == "user" {
|
||||
if assistantStarted {
|
||||
h.sess.AppendLog([]byte("\n"))
|
||||
assistantStarted = false
|
||||
}
|
||||
h.sess.AppendLog(FormatEvent(ev))
|
||||
return
|
||||
} else {
|
||||
switch ev.Role {
|
||||
case "assistant":
|
||||
if !assistantStarted {
|
||||
h.sess.AppendLog([]byte("assistant: "))
|
||||
h.sess.mu.Lock()
|
||||
h.sess.ChatOffset = len(h.sess.log)
|
||||
h.sess.mu.Unlock()
|
||||
assistantStarted = true
|
||||
}
|
||||
default:
|
||||
if assistantStarted {
|
||||
h.sess.AppendLog([]byte("\n"))
|
||||
assistantStarted = false
|
||||
}
|
||||
}
|
||||
}
|
||||
h.sess.AppendLog(FormatEvent(ev))
|
||||
}
|
||||
}
|
||||
|
||||
func (h *sessionHelper) handleCtl(input string) error {
|
||||
cmd := strings.Fields(input)
|
||||
if len(cmd) == 0 {
|
||||
return fmt.Errorf("empty ctl command")
|
||||
}
|
||||
switch cmd[0] {
|
||||
case "stop":
|
||||
h.sess.Core.Interrupt(agent.ErrInterrupted)
|
||||
case "kill":
|
||||
h.kill()
|
||||
case "rn":
|
||||
if name := strings.TrimSpace(input[3:]); name != "" {
|
||||
if err := h.rename(name); err != nil {
|
||||
h.log.Error("rename: %v", err)
|
||||
}
|
||||
}
|
||||
case "save":
|
||||
h.sess.mu.RLock()
|
||||
data := make([]byte, len(h.sess.log))
|
||||
copy(data, h.sess.log)
|
||||
h.sess.mu.RUnlock()
|
||||
return h.saveTranscript(data)
|
||||
case "compact", "clear", "backend", "model", "models",
|
||||
"agents", "agent", "sessions", "cwd", "skills",
|
||||
"tools", "context", "usage", "cost", "history",
|
||||
"irw", "help":
|
||||
h.sess.Core.Submit(h.sess.Ctx, "/"+input, h.makePublish())
|
||||
default:
|
||||
return fmt.Errorf("unknown ctl command: %s", cmd[0])
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// FormatParams formats generation parameters as key=value lines.
|
||||
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()
|
||||
}
|
||||
|
||||
// ParseParams parses key=value lines into generation parameters.
|
||||
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
|
||||
}
|
||||
|
|
@ -0,0 +1,235 @@
|
|||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"ollie/pkg/paths"
|
||||
"ollie/pkg/skills"
|
||||
)
|
||||
|
||||
// SkillStoreConfig holds optional dependencies for NewSkillStore.
|
||||
// Nil functions default to their os package equivalents.
|
||||
type SkillStoreConfig struct {
|
||||
Dirs []string
|
||||
ReadDir func(string) ([]os.DirEntry, error)
|
||||
Open func(string) (*os.File, error)
|
||||
ReadFile func(string) ([]byte, error)
|
||||
WriteFile func(string, []byte, os.FileMode) error
|
||||
MkdirAll func(string, os.FileMode) error
|
||||
RemoveAll func(string) error
|
||||
Rename func(string, string) error
|
||||
}
|
||||
|
||||
type skillState struct {
|
||||
dirs []string
|
||||
readDir func(string) ([]os.DirEntry, error)
|
||||
openFile func(string) (*os.File, error)
|
||||
readFile func(string) ([]byte, error)
|
||||
writeFile func(string, []byte, os.FileMode) error
|
||||
mkdirAll func(string, os.FileMode) error
|
||||
removeAll func(string) error
|
||||
rename func(string, string) error
|
||||
}
|
||||
|
||||
func NewSkillStore() Store {
|
||||
return NewSkillStoreWith(SkillStoreConfig{})
|
||||
}
|
||||
|
||||
func NewSkillStoreWith(cfg SkillStoreConfig) Store {
|
||||
if cfg.Dirs == nil {
|
||||
cfg.Dirs = skillDirs()
|
||||
}
|
||||
if cfg.ReadDir == nil {
|
||||
cfg.ReadDir = os.ReadDir
|
||||
}
|
||||
if cfg.Open == nil {
|
||||
cfg.Open = os.Open
|
||||
}
|
||||
if cfg.ReadFile == nil {
|
||||
cfg.ReadFile = os.ReadFile
|
||||
}
|
||||
if cfg.WriteFile == nil {
|
||||
cfg.WriteFile = os.WriteFile
|
||||
}
|
||||
if cfg.MkdirAll == nil {
|
||||
cfg.MkdirAll = os.MkdirAll
|
||||
}
|
||||
if cfg.RemoveAll == nil {
|
||||
cfg.RemoveAll = os.RemoveAll
|
||||
}
|
||||
if cfg.Rename == nil {
|
||||
cfg.Rename = os.Rename
|
||||
}
|
||||
|
||||
ss := &skillState{
|
||||
dirs: cfg.Dirs,
|
||||
readDir: cfg.ReadDir,
|
||||
openFile: cfg.Open,
|
||||
readFile: cfg.ReadFile,
|
||||
writeFile: cfg.WriteFile,
|
||||
mkdirAll: cfg.MkdirAll,
|
||||
removeAll: cfg.RemoveAll,
|
||||
rename: cfg.Rename,
|
||||
}
|
||||
|
||||
return &storeConfig{
|
||||
StatFn: ss.stat,
|
||||
ListFn: ss.list,
|
||||
OpenFn: ss.open,
|
||||
DeleteFn: ss.del,
|
||||
CreateFn: ss.create,
|
||||
RenameFn: ss.ren,
|
||||
}
|
||||
}
|
||||
|
||||
func skillDirs() []string {
|
||||
if env := os.Getenv("OLLIE_SKILLS_PATH"); env != "" {
|
||||
return strings.Split(env, ":")
|
||||
}
|
||||
return []string{paths.CfgDir() + "/skills"}
|
||||
}
|
||||
|
||||
func (ss *skillState) listSkills() []skills.Meta {
|
||||
seen := make(map[string]bool)
|
||||
var result []skills.Meta
|
||||
for _, dir := range ss.dirs {
|
||||
entries, err := ss.readDir(dir)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
for _, e := range entries {
|
||||
if !e.IsDir() || seen[e.Name()] {
|
||||
continue
|
||||
}
|
||||
skillDir := filepath.Join(dir, e.Name())
|
||||
f, err := ss.openFile(filepath.Join(skillDir, "SKILL.md"))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
meta, err := skills.ParseFrontMatter(f, filepath.Base(skillDir), skillDir)
|
||||
f.Close()
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
seen[meta.Name] = true
|
||||
result = append(result, *meta)
|
||||
}
|
||||
}
|
||||
sort.Slice(result, func(i, j int) bool { return result[i].Name < result[j].Name })
|
||||
return result
|
||||
}
|
||||
|
||||
func (ss *skillState) readSkill(name string) ([]byte, error) {
|
||||
for _, m := range ss.listSkills() {
|
||||
if m.Name == name {
|
||||
return ss.readFile(filepath.Join(m.Dir, "SKILL.md"))
|
||||
}
|
||||
}
|
||||
return nil, os.ErrNotExist
|
||||
}
|
||||
|
||||
func (ss *skillState) stat(name string) (os.FileInfo, error) {
|
||||
if name == "idx" {
|
||||
return &SyntheticFileInfo{Name_: "idx", Mode_: 0444}, nil
|
||||
}
|
||||
skillName := strings.TrimSuffix(name, ".md")
|
||||
if _, err := ss.readSkill(skillName); err != nil {
|
||||
return nil, fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
return &SyntheticFileInfo{Name_: name, Mode_: 0666}, nil
|
||||
}
|
||||
|
||||
func (ss *skillState) list() ([]os.DirEntry, error) {
|
||||
result := []os.DirEntry{FileEntry("idx", 0444)}
|
||||
for _, m := range ss.listSkills() {
|
||||
result = append(result, FileEntry(m.Name+".md", 0666))
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (ss *skillState) open(name string) (StoreEntry, error) {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
if name == "idx" {
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return &SyntheticFileInfo{Name_: "idx", Mode_: 0444}, nil },
|
||||
ReadFn: func() ([]byte, error) { return ss.index() },
|
||||
WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") },
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
}
|
||||
skillName := strings.TrimSuffix(name, ".md")
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) {
|
||||
if _, err := ss.readSkill(skillName); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &SyntheticFileInfo{Name_: name, Mode_: 0666}, nil
|
||||
},
|
||||
ReadFn: func() ([]byte, error) { return ss.readSkill(skillName) },
|
||||
WriteFn: func(data []byte) error {
|
||||
dir := ""
|
||||
for _, m := range ss.listSkills() {
|
||||
if m.Name == skillName {
|
||||
dir = m.Dir
|
||||
break
|
||||
}
|
||||
}
|
||||
if dir == "" {
|
||||
dir = filepath.Join(ss.dirs[0], skillName)
|
||||
}
|
||||
if err := ss.mkdirAll(dir, 0755); err != nil {
|
||||
return err
|
||||
}
|
||||
return ss.writeFile(filepath.Join(dir, "SKILL.md"), data, 0644)
|
||||
},
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (ss *skillState) del(name string) error {
|
||||
skillName := strings.TrimSuffix(name, ".md")
|
||||
for _, m := range ss.listSkills() {
|
||||
if m.Name == skillName {
|
||||
return ss.removeAll(m.Dir)
|
||||
}
|
||||
}
|
||||
return fmt.Errorf("skill not found: %s", skillName)
|
||||
}
|
||||
|
||||
func (ss *skillState) ren(oldName, newName string) error {
|
||||
oldSkill := strings.TrimSuffix(oldName, ".md")
|
||||
newSkill := strings.TrimSuffix(newName, ".md")
|
||||
for _, m := range ss.listSkills() {
|
||||
if m.Name == oldSkill {
|
||||
newDir := filepath.Join(filepath.Dir(m.Dir), newSkill)
|
||||
return ss.rename(m.Dir, newDir)
|
||||
}
|
||||
}
|
||||
return fmt.Errorf("skill not found: %s", oldSkill)
|
||||
}
|
||||
|
||||
func (ss *skillState) create(name string) error {
|
||||
skillName := strings.TrimSuffix(name, ".md")
|
||||
dir := filepath.Join(ss.dirs[0], skillName)
|
||||
if err := ss.mkdirAll(dir, 0755); err != nil {
|
||||
return err
|
||||
}
|
||||
return ss.writeFile(filepath.Join(dir, "SKILL.md"), nil, 0644)
|
||||
}
|
||||
|
||||
func (ss *skillState) index() ([]byte, error) {
|
||||
var sb strings.Builder
|
||||
for _, m := range ss.listSkills() {
|
||||
fmt.Fprintf(&sb, "## %s\n", m.Name)
|
||||
fmt.Fprintf(&sb, "description: %s\n", m.Description)
|
||||
sb.WriteString("\n")
|
||||
}
|
||||
return []byte(sb.String()), nil
|
||||
}
|
||||
|
|
@ -0,0 +1,156 @@
|
|||
// Package store defines storage interfaces and a filesystem-backed
|
||||
// implementation used by the 9P server.
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"time"
|
||||
)
|
||||
|
||||
// StoreEntry is an open handle to a single named blob.
|
||||
type StoreEntry interface {
|
||||
Stat() (os.FileInfo, error)
|
||||
Read() ([]byte, error)
|
||||
Write(data []byte) error
|
||||
BlockingRead(ctx context.Context, base string) ([]byte, error)
|
||||
}
|
||||
|
||||
// Store is a named collection of entries.
|
||||
type Store interface {
|
||||
Stat(name string) (os.FileInfo, error)
|
||||
List() ([]os.DirEntry, error)
|
||||
Open(name string) (StoreEntry, error)
|
||||
Create(name string) error
|
||||
Delete(name string) error
|
||||
Rename(oldName, newName string) error
|
||||
}
|
||||
|
||||
// Runnable is a running agent that can be observed and controlled.
|
||||
type Runnable interface {
|
||||
RunnableID() string
|
||||
Cancel()
|
||||
Interrupt()
|
||||
AppendLog([]byte)
|
||||
LogInfo() (length int, vers uint32)
|
||||
}
|
||||
|
||||
// RunnableStore is a Store backed by a running agent.
|
||||
type RunnableStore interface {
|
||||
Store
|
||||
Runnable
|
||||
}
|
||||
|
||||
// EntryConfig implements StoreEntry via function pointers.
|
||||
type EntryConfig struct {
|
||||
StatFn func() (os.FileInfo, error)
|
||||
ReadFn func() ([]byte, error)
|
||||
WriteFn func([]byte) error
|
||||
BlockingReadFn func(context.Context, string) ([]byte, error)
|
||||
}
|
||||
|
||||
func (e *EntryConfig) Stat() (os.FileInfo, error) { return e.StatFn() }
|
||||
func (e *EntryConfig) Read() ([]byte, error) { return e.ReadFn() }
|
||||
func (e *EntryConfig) Write(data []byte) error { return e.WriteFn(data) }
|
||||
func (e *EntryConfig) BlockingRead(ctx context.Context, base string) ([]byte, error) {
|
||||
return e.BlockingReadFn(ctx, base)
|
||||
}
|
||||
|
||||
// storeConfig implements Store via function pointers.
|
||||
type storeConfig struct {
|
||||
StatFn func(string) (os.FileInfo, error)
|
||||
ListFn func() ([]os.DirEntry, error)
|
||||
OpenFn func(string) (StoreEntry, error)
|
||||
CreateFn func(string) error
|
||||
DeleteFn func(string) error
|
||||
RenameFn func(string, string) error
|
||||
}
|
||||
|
||||
func (s *storeConfig) Stat(name string) (os.FileInfo, error) { return s.StatFn(name) }
|
||||
func (s *storeConfig) List() ([]os.DirEntry, error) { return s.ListFn() }
|
||||
func (s *storeConfig) Open(name string) (StoreEntry, error) { return s.OpenFn(name) }
|
||||
func (s *storeConfig) Create(name string) error { return s.CreateFn(name) }
|
||||
func (s *storeConfig) Delete(name string) error { return s.DeleteFn(name) }
|
||||
func (s *storeConfig) Rename(old, new string) error { return s.RenameFn(old, new) }
|
||||
|
||||
// NewFlatDir returns a Store backed by a directory on the local filesystem.
|
||||
func NewFlatDir(dir string, perm os.FileMode) Store {
|
||||
join := func(name string) string { return filepath.Join(dir, name) }
|
||||
ensureDir := func() error { return os.MkdirAll(dir, 0755) }
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
|
||||
return &storeConfig{
|
||||
StatFn: func(name string) (os.FileInfo, error) { return os.Stat(join(name)) },
|
||||
ListFn: func() ([]os.DirEntry, error) { return os.ReadDir(dir) },
|
||||
OpenFn: func(name string) (StoreEntry, error) {
|
||||
path := join(name)
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return os.Stat(path) },
|
||||
ReadFn: func() ([]byte, error) { return os.ReadFile(path) },
|
||||
WriteFn: func(data []byte) error {
|
||||
if err := ensureDir(); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.WriteFile(path, data, perm)
|
||||
},
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
},
|
||||
CreateFn: func(name string) error {
|
||||
if err := ensureDir(); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.WriteFile(join(name), nil, perm)
|
||||
},
|
||||
DeleteFn: func(name string) error { return os.Remove(join(name)) },
|
||||
RenameFn: func(old, new string) error { return os.Rename(join(old), join(new)) },
|
||||
}
|
||||
}
|
||||
|
||||
// SyntheticFileInfo implements os.FileInfo for entries with no backing file.
|
||||
type SyntheticFileInfo struct {
|
||||
Name_ string
|
||||
Mode_ os.FileMode
|
||||
Size_ int64
|
||||
IsDir_ bool
|
||||
}
|
||||
|
||||
func (f *SyntheticFileInfo) Name() string { return f.Name_ }
|
||||
func (f *SyntheticFileInfo) Size() int64 { return f.Size_ }
|
||||
func (f *SyntheticFileInfo) Mode() os.FileMode { return f.Mode_ }
|
||||
func (f *SyntheticFileInfo) ModTime() time.Time { return time.Time{} }
|
||||
func (f *SyntheticFileInfo) IsDir() bool { return f.IsDir_ }
|
||||
func (f *SyntheticFileInfo) Sys() any { return nil }
|
||||
|
||||
// SyntheticEntry implements os.DirEntry for entries with no backing file.
|
||||
type SyntheticEntry struct {
|
||||
Name_ string
|
||||
Mode_ os.FileMode
|
||||
IsDir_ bool
|
||||
}
|
||||
|
||||
func (e *SyntheticEntry) Name() string { return e.Name_ }
|
||||
func (e *SyntheticEntry) IsDir() bool { return e.IsDir_ }
|
||||
func (e *SyntheticEntry) Type() os.FileMode {
|
||||
if e.IsDir_ {
|
||||
return os.ModeDir
|
||||
}
|
||||
return 0
|
||||
}
|
||||
func (e *SyntheticEntry) Info() (os.FileInfo, error) {
|
||||
return &SyntheticFileInfo{Name_: e.Name_, Mode_: e.Mode_, IsDir_: e.IsDir_}, nil
|
||||
}
|
||||
|
||||
// FileEntry returns a synthetic file DirEntry.
|
||||
func FileEntry(name string, mode os.FileMode) os.DirEntry {
|
||||
return &SyntheticEntry{Name_: name, Mode_: mode}
|
||||
}
|
||||
|
||||
// DirEntry returns a synthetic directory DirEntry.
|
||||
func DirEntry(name string, mode os.FileMode) os.DirEntry {
|
||||
return &SyntheticEntry{Name_: name, Mode_: mode, IsDir_: true}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue