agents: use OLLIE_AGENTS_PATH for multi-dir agent resolution
This commit is contained in:
parent
323ae311fe
commit
182682cbb7
19
server.go
19
server.go
|
|
@ -140,7 +140,7 @@ func (s *ToolStore) Open(name string) (StoreEntry, error) {
|
|||
StatFn: func() (os.FileInfo, error) { return &syntheticFileInfo{Name_: "idx", Mode_: 0444}, nil },
|
||||
ReadFn: func() ([]byte, error) { return s.index() },
|
||||
WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") },
|
||||
BlockingReadFn: func(context.Context, string) ([]byte, error) { return nil, fmt.Errorf("blocking read not supported") },
|
||||
BlockingReadFn: func(context.Context, string) ([]byte, string, error) { return nil, "", fmt.Errorf("blocking read not supported") },
|
||||
}, nil
|
||||
}
|
||||
return s.FlatDirStore.Open(name)
|
||||
|
|
@ -244,13 +244,14 @@ func New(sink *olog.Sink, opts ...ServerOption) *Server {
|
|||
os.MkdirAll(transcriptDir, 0755) //nolint:errcheck
|
||||
tmpDir := defaultTmpDir()
|
||||
os.MkdirAll(tmpDir, 0755) //nolint:errcheck
|
||||
agentsDir := paths.CfgDir() + "/agents"
|
||||
agentsDirs := agent.AgentsDirs()
|
||||
agentsDir := agentsDirs[0]
|
||||
sessionsDir := paths.CfgDir() + "/sessions"
|
||||
s := &Server{
|
||||
log: sink.Logger("9p", olog.LevelDebug),
|
||||
sink: sink,
|
||||
agentsDir: agentsDir,
|
||||
agentStore: NewFlatDirStore(agentsDir, 0644),
|
||||
agentStore: store.NewFlatDirWritableUnion(agentsDirs, 0644),
|
||||
promptStore: store.NewFlatDirUnion(agent.PromptsDirs(), 0444),
|
||||
memStore: NewFlatDirStore(memDir, 0644),
|
||||
toolStore: NewToolStore(),
|
||||
|
|
@ -372,10 +373,10 @@ func storeReadCtx(ctx context.Context, s store.Store, name string) ([]byte, erro
|
|||
}
|
||||
|
||||
// storeBlockingRead opens an entry and performs a blocking read.
|
||||
func storeBlockingRead(s store.Store, name string, ctx context.Context, base string) ([]byte, error) {
|
||||
func storeBlockingRead(s store.Store, name string, ctx context.Context, base string) (content []byte, nextBase string, err error) {
|
||||
e, err := s.Open(name)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, "", err
|
||||
}
|
||||
return e.BlockingRead(ctx, base)
|
||||
}
|
||||
|
|
@ -962,15 +963,15 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
// be delivered to the blocked client process.
|
||||
waitCtx, waitCancel := context.WithTimeout(ctx, 5*time.Second)
|
||||
defer waitCancel()
|
||||
content, err := storeBlockingRead(sfs, parts[2], waitCtx, base)
|
||||
content, nextBase, err := storeBlockingRead(sfs, parts[2], waitCtx, base)
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
// Update the fid's baseline to the returned value for subsequent reads.
|
||||
if content != nil {
|
||||
// Update the fid's baseline for subsequent reads.
|
||||
if nextBase != "" {
|
||||
cs.mu.Lock()
|
||||
if f, ok := cs.fids[fc.Fid]; ok {
|
||||
f.waitBase = strings.TrimSuffix(string(content), "\n")
|
||||
f.waitBase = nextBase
|
||||
}
|
||||
cs.mu.Unlock()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -58,8 +58,8 @@ func (m *MemoryStore) List() ([]os.DirEntry, error) {
|
|||
}
|
||||
|
||||
func (m *MemoryStore) Open(name string) (StoreEntry, error) {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
notBlocking := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return m.Stat(name) },
|
||||
|
|
|
|||
|
|
@ -12,7 +12,7 @@ type FileSpec struct {
|
|||
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
|
||||
Wait func(ctx context.Context, base string) (content []byte, nextBase string, err error) // nil = not waitable
|
||||
Size func() int64 // optional; if nil, len(Read())
|
||||
}
|
||||
|
||||
|
|
@ -90,8 +90,8 @@ func (rs *RunFileStore) open(name string) (StoreEntry, error) {
|
|||
if spec.Write != nil {
|
||||
writeFn = spec.Write
|
||||
}
|
||||
waitFn := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("%s: not a wait file", name)
|
||||
waitFn := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("%s: not a wait file", name)
|
||||
}
|
||||
if spec.Wait != nil {
|
||||
waitFn = spec.Wait
|
||||
|
|
|
|||
|
|
@ -227,8 +227,8 @@ func (s *SessionStore) stat(name string) (os.FileInfo, error) {
|
|||
}
|
||||
|
||||
func (s *SessionStore) openEntry(name string) (StoreEntry, error) {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
notBlocking := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
switch name {
|
||||
case "new":
|
||||
|
|
@ -501,7 +501,8 @@ func LoadAgentConfig(agentsDir, name string, open func(string) (*os.File, error)
|
|||
if open == nil {
|
||||
open = os.Open
|
||||
}
|
||||
f, err := open(agentsDir + "/" + name + ".json")
|
||||
path := agent.AgentConfigPath(agentsDir, name)
|
||||
f, err := open(path)
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ var SessionFileList = []struct {
|
|||
{"chat", 0666},
|
||||
{"offset", 0444},
|
||||
{"cfg", 0666},
|
||||
{"state", 0444},
|
||||
{"statewait", 0444},
|
||||
{"usage", 0444},
|
||||
{"cost", 0444},
|
||||
|
|
@ -178,7 +179,7 @@ func (h *sessionHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
|
||||
// Wait
|
||||
if name == "statewait" {
|
||||
fs.Wait = func(connCtx context.Context, base string) ([]byte, error) {
|
||||
fs.Wait = func(connCtx context.Context, base string) ([]byte, string, error) {
|
||||
ctx, cancel := context.WithCancel(connCtx)
|
||||
defer cancel()
|
||||
context.AfterFunc(h.sess.Ctx, cancel)
|
||||
|
|
@ -187,19 +188,19 @@ func (h *sessionHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
}
|
||||
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
|
||||
s := h.sess.Core.State()
|
||||
return []byte(s + "\n"), s, nil
|
||||
}
|
||||
return []byte(v + "\n"), nil
|
||||
return []byte(v + "\n"), v, nil
|
||||
}
|
||||
}
|
||||
if name == "chatwait" {
|
||||
fs.Wait = func(connCtx context.Context, base string) ([]byte, error) {
|
||||
fs.Wait = func(connCtx context.Context, base string) ([]byte, string, error) {
|
||||
ctx, cancel := context.WithCancel(connCtx)
|
||||
defer cancel()
|
||||
context.AfterFunc(h.sess.Ctx, cancel)
|
||||
if h.sess.Core.State() == "idle" {
|
||||
return nil, nil
|
||||
return nil, base, nil
|
||||
}
|
||||
offset := 0
|
||||
if base == "" {
|
||||
|
|
@ -210,9 +211,9 @@ func (h *sessionHelper) fileSpec(name string, mode os.FileMode) FileSpec {
|
|||
}
|
||||
data := h.sess.WaitLog(ctx, offset)
|
||||
if data == nil {
|
||||
return nil, nil
|
||||
return nil, strconv.Itoa(offset), nil
|
||||
}
|
||||
return data, nil
|
||||
return data, strconv.Itoa(offset + len(data)), nil
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -233,6 +234,8 @@ func (h *sessionHelper) content(name string) string {
|
|||
return h.sess.Core.ListModels() + "\n"
|
||||
case "systemprompt":
|
||||
return h.sess.Core.SystemPrompt()
|
||||
case "state":
|
||||
return h.sess.Core.State() + "\n"
|
||||
case "offset":
|
||||
return fmt.Sprintf("%d\n", h.sess.ChatOffset)
|
||||
case "env":
|
||||
|
|
|
|||
|
|
@ -153,8 +153,8 @@ func (ss *skillState) list() ([]os.DirEntry, error) {
|
|||
}
|
||||
|
||||
func (ss *skillState) open(name string) (StoreEntry, error) {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
notBlocking := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
if name == "idx" {
|
||||
return &EntryConfig{
|
||||
|
|
|
|||
101
store/store.go
101
store/store.go
|
|
@ -15,7 +15,7 @@ type StoreEntry interface {
|
|||
Stat() (os.FileInfo, error)
|
||||
Read() ([]byte, error)
|
||||
Write(data []byte) error
|
||||
BlockingRead(ctx context.Context, base string) ([]byte, error)
|
||||
BlockingRead(ctx context.Context, base string) (content []byte, nextBase string, err error)
|
||||
}
|
||||
|
||||
// Store is a named collection of entries.
|
||||
|
|
@ -48,13 +48,13 @@ type EntryConfig struct {
|
|||
StatFn func() (os.FileInfo, error)
|
||||
ReadFn func() ([]byte, error)
|
||||
WriteFn func([]byte) error
|
||||
BlockingReadFn func(context.Context, string) ([]byte, error)
|
||||
BlockingReadFn func(context.Context, string) ([]byte, string, 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) {
|
||||
func (e *EntryConfig) BlockingRead(ctx context.Context, base string) ([]byte, string, error) {
|
||||
return e.BlockingReadFn(ctx, base)
|
||||
}
|
||||
|
||||
|
|
@ -78,8 +78,8 @@ func (s *storeConfig) Rename(old, new string) error { return s.RenameF
|
|||
// 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) }
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
notBlocking := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
|
||||
return &storeConfig{
|
||||
|
|
@ -111,11 +111,98 @@ func NewFlatDir(dir string, perm os.FileMode) Store {
|
|||
}
|
||||
}
|
||||
|
||||
// NewFlatDirWritableUnion returns a Store that reads from all dirs but writes to the first.
|
||||
// Files are deduplicated by name; the first directory containing a name wins.
|
||||
func NewFlatDirWritableUnion(dirs []string, perm os.FileMode) Store {
|
||||
writeDir := dirs[0]
|
||||
notBlocking := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
|
||||
resolve := func(name string) (string, error) {
|
||||
for _, dir := range dirs {
|
||||
p := filepath.Join(dir, name)
|
||||
if _, err := os.Stat(p); err == nil {
|
||||
return p, nil
|
||||
}
|
||||
}
|
||||
return "", fmt.Errorf("%s: not found", name)
|
||||
}
|
||||
|
||||
return &storeConfig{
|
||||
StatFn: func(name string) (os.FileInfo, error) {
|
||||
p, err := resolve(name)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return os.Stat(p)
|
||||
},
|
||||
ListFn: func() ([]os.DirEntry, error) {
|
||||
seen := make(map[string]bool)
|
||||
var result []os.DirEntry
|
||||
for _, dir := range dirs {
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
for _, e := range entries {
|
||||
if !seen[e.Name()] {
|
||||
seen[e.Name()] = true
|
||||
result = append(result, e)
|
||||
}
|
||||
}
|
||||
}
|
||||
return result, nil
|
||||
},
|
||||
OpenFn: func(name string) (StoreEntry, error) {
|
||||
p, err := resolve(name)
|
||||
if err != nil {
|
||||
// Not found anywhere; write path is in writeDir.
|
||||
p = filepath.Join(writeDir, name)
|
||||
}
|
||||
wp := filepath.Join(writeDir, name)
|
||||
return &EntryConfig{
|
||||
StatFn: func() (os.FileInfo, error) { return os.Stat(p) },
|
||||
ReadFn: func() ([]byte, error) { return os.ReadFile(p) },
|
||||
WriteFn: func(data []byte) error {
|
||||
if err := os.MkdirAll(filepath.Dir(wp), 0755); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.WriteFile(wp, data, perm)
|
||||
},
|
||||
BlockingReadFn: notBlocking,
|
||||
}, nil
|
||||
},
|
||||
CreateFn: func(name string) error {
|
||||
p := filepath.Join(writeDir, name)
|
||||
if err := os.MkdirAll(filepath.Dir(p), 0755); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.WriteFile(p, nil, perm)
|
||||
},
|
||||
DeleteFn: func(name string) error {
|
||||
// Delete from wherever it exists.
|
||||
p, err := resolve(name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return os.Remove(p)
|
||||
},
|
||||
RenameFn: func(old, new string) error {
|
||||
p, err := resolve(old)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return os.Rename(p, filepath.Join(filepath.Dir(p), new))
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// NewFlatDirUnion returns a read-only Store that merges multiple directories.
|
||||
// Files are deduplicated by name; the first directory containing a name wins.
|
||||
func NewFlatDirUnion(dirs []string, perm os.FileMode) Store {
|
||||
notBlocking := func(context.Context, string) ([]byte, error) {
|
||||
return nil, fmt.Errorf("blocking read not supported")
|
||||
notBlocking := func(context.Context, string) ([]byte, string, error) {
|
||||
return nil, "", fmt.Errorf("blocking read not supported")
|
||||
}
|
||||
|
||||
resolve := func(name string) (string, error) {
|
||||
|
|
|
|||
|
|
@ -1061,7 +1061,7 @@ func TestSessionFileStoreBlockingRead(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("Open(statewait): %v", err)
|
||||
}
|
||||
data, err := e.BlockingRead(context.Background(), "idle")
|
||||
data, _, err := e.BlockingRead(context.Background(), "idle")
|
||||
if err != nil {
|
||||
t.Fatalf("BlockingRead: %v", err)
|
||||
}
|
||||
|
|
@ -1079,7 +1079,7 @@ func TestSessionFileStoreBlockingReadCancel(t *testing.T) {
|
|||
cancel()
|
||||
|
||||
e, _ := sf.Open("statewait")
|
||||
data, err := e.BlockingRead(ctx, "idle")
|
||||
data, _, err := e.BlockingRead(ctx, "idle")
|
||||
if err != nil {
|
||||
t.Fatalf("BlockingRead error: %v", err)
|
||||
}
|
||||
|
|
@ -1095,7 +1095,7 @@ func TestSessionFileStoreBlockingReadNotWaitFile(t *testing.T) {
|
|||
sf, _ := newSessionFileStore(t, sess)
|
||||
|
||||
e, _ := sf.Open("chat")
|
||||
if _, err := e.BlockingRead(context.Background(), ""); err == nil {
|
||||
if _, _, err := e.BlockingRead(context.Background(), ""); err == nil {
|
||||
t.Error("BlockingRead(chat) should error")
|
||||
}
|
||||
}
|
||||
|
|
@ -1116,7 +1116,7 @@ func TestSessionFileStoreBlockingReadAllWaitFiles(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("Open(%q): %v", name, err)
|
||||
}
|
||||
data, err := e.BlockingRead(context.Background(), "")
|
||||
data, _, err := e.BlockingRead(context.Background(), "")
|
||||
if err != nil {
|
||||
t.Fatalf("BlockingRead(%q): %v", name, err)
|
||||
}
|
||||
|
|
@ -1138,7 +1138,7 @@ func TestSessionFileStoreBlockingReadIdleTimeout(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Fatalf("Open(statewait): %v", err)
|
||||
}
|
||||
data, err := e.BlockingRead(ctx, "")
|
||||
data, _, err := e.BlockingRead(ctx, "")
|
||||
if err != nil {
|
||||
t.Fatalf("BlockingRead: %v", err)
|
||||
}
|
||||
|
|
|
|||
Reference in New Issue