diff --git a/server.go b/server.go index 30c60dd..0ae466b 100644 --- a/server.go +++ b/server.go @@ -99,6 +99,60 @@ func (s *ExecStore) Open(name string) (StoreEntry, error) { return s.FlatDirStore.Open(name) } +// NewRootStore returns a read-only Store for synthetic root-level files. +func NewRootStore() Store { + helpPath := paths.CfgDir() + "/help.md" + notBlocking := func(context.Context, string) ([]byte, string, error) { + return nil, "", fmt.Errorf("blocking read not supported") + } + readOnly := func([]byte) error { return fmt.Errorf("read-only") } + + entries := map[string]func() ([]byte, error){ + "backends": func() ([]byte, error) { + return []byte(strings.Join(backend.Backends(), "\n") + "\n"), nil + }, + "help": func() ([]byte, error) { + return os.ReadFile(helpPath) + }, + } + + return &rootStore{entries: entries, notBlocking: notBlocking, readOnly: readOnly} +} + +type rootStore struct { + entries map[string]func() ([]byte, error) + notBlocking func(context.Context, string) ([]byte, string, error) + readOnly func([]byte) error +} + +func (r *rootStore) Stat(name string) (os.FileInfo, error) { + if _, ok := r.entries[name]; ok { + return &syntheticFileInfo{Name_: name, Mode_: 0444}, nil + } + return nil, fmt.Errorf("%s: not found", name) +} +func (r *rootStore) List() ([]os.DirEntry, error) { + return []os.DirEntry{ + syntheticEntry("backends", 0444), + syntheticEntry("help", 0444), + }, nil +} +func (r *rootStore) Open(name string) (StoreEntry, error) { + readFn, ok := r.entries[name] + if !ok { + return nil, fmt.Errorf("%s: not found", name) + } + return &store.EntryConfig{ + StatFn: func() (os.FileInfo, error) { return &syntheticFileInfo{Name_: name, Mode_: 0444}, nil }, + ReadFn: readFn, + WriteFn: r.readOnly, + BlockingReadFn: r.notBlocking, + }, nil +} +func (r *rootStore) Create(string) error { return fmt.Errorf("read-only store") } +func (r *rootStore) Delete(string) error { return fmt.Errorf("read-only store") } +func (r *rootStore) Rename(string, string) error { return fmt.Errorf("read-only store") } + // --- tools --- // ToolStore is a BlobStore backed by the tools directory. @@ -147,6 +201,25 @@ func (s *ToolStore) Open(name string) (StoreEntry, error) { return s.FlatDirStore.Open(name) } +// OpenFiltered opens an entry, applying allowTools filtering to idx if non-nil. +func (s *ToolStore) OpenFiltered(name string, allowed map[string]bool) (StoreEntry, error) { + if name == "idx" && allowed != nil { + return &store.EntryConfig{ + StatFn: func() (os.FileInfo, error) { return &syntheticFileInfo{Name_: "idx", Mode_: 0444}, nil }, + ReadFn: func() ([]byte, error) { + data, err := s.index() + if err != nil { + return nil, err + } + return filterToolIndex(data, allowed), nil + }, + WriteFn: func([]byte) error { return fmt.Errorf("idx: read-only") }, + BlockingReadFn: func(context.Context, string) ([]byte, string, error) { return nil, "", fmt.Errorf("blocking read not supported") }, + }, nil + } + return s.Open(name) +} + func (s *ToolStore) index() ([]byte, error) { entries, err := s.FlatDirStore.List() if err != nil { @@ -228,6 +301,7 @@ type Server struct { sessionStore *SessionStore transcriptStore Store tmpStore Store + rootStore Store strict bool yolo bool groups map[string]map[string]bool // group → set of members @@ -263,6 +337,7 @@ func New(sink *olog.Sink, opts ...ServerOption) *Server { skillStore: NewSkillStore(), transcriptStore: NewFlatDirStore(transcriptDir, 0444), tmpStore: NewFlatDirStore(tmpDir, 0600), + rootStore: NewRootStore(), groups: make(map[string]map[string]bool), } for _, o := range opts { @@ -302,6 +377,92 @@ func New(sink *olog.Sink, opts ...ServerOption) *Server { return s } +// storeRoute maps a path prefix to its backing store. +type storeRoute struct { + prefix string + store func() Store +} + +// storeRoutes returns the route table for generic store-backed paths. +func (s *Server) storeRoutes() []storeRoute { + return []storeRoute{ + {"/a/", func() Store { return s.agentStore }}, + {"/p/", func() Store { return s.promptStore }}, + {"/m/", func() Store { return s.memStore }}, + {"/sk/", func() Store { return s.skillStore }}, + {"/u/", func() Store { return s.utilStore }}, + {"/x/", func() Store { return s.pluginStore }}, + {"/tmp/", func() Store { return s.tmpStore }}, + {"/tr/", func() Store { return s.transcriptStore }}, + } +} + +// routeStore returns the store and entry name for a path, or nil if no route matches. +func (s *Server) routeStore(path string) (Store, string) { + for _, r := range s.storeRoutes() { + if strings.HasPrefix(path, r.prefix) { + return r.store(), pathBase(path) + } + } + // Root-level files (e.g. /backends, /help). + name := strings.TrimPrefix(path, "/") + if !strings.Contains(name, "/") { + if _, err := s.rootStore.Stat(name); err == nil { + return s.rootStore, name + } + } + // Session store files (/s/sh, /s/bfg, etc.) + if isSessionStoreFile(path) { + return s.sessionStore, pathBase(path) + } + // Session paths: /s/{id}/{file} or /s/{id}/t/{tool} + if strings.HasPrefix(path, "/s/") { + parts := strings.SplitN(strings.TrimPrefix(path, "/s/"), "/", 3) + if len(parts) >= 2 { + sessID := parts[0] + if parts[1] == "t" && len(parts) == 3 { + // Tool file: /s/{id}/t/{rel} + return s.toolStore, parts[2] + } + if parts[1] != "t" { + // Session file: /s/{id}/{file} + if sfs, ok := s.sessionFileStore(sessID); ok { + return sfs, parts[1] + } + } + } + } + return nil, "" +} + +// openEntry opens a StoreEntry for the given path, applying any path-specific +// logic (e.g. tool idx filtering by session allowTools). +func (s *Server) openEntry(path string) (StoreEntry, error) { + // Tool idx with per-session filtering. + if strings.HasPrefix(path, "/s/") { + parts := strings.SplitN(strings.TrimPrefix(path, "/s/"), "/", 3) + if len(parts) == 3 && parts[1] == "t" { + allowed := s.sessionAllowTools(parts[0]) + return s.toolStore.OpenFiltered(parts[2], allowed) + } + } + st, name := s.routeStore(path) + if st == nil { + return nil, fmt.Errorf("%s: not found", path) + } + return st.Open(name) +} + +// routeDir returns the store for a directory path (e.g. "/a" → agentStore). +func (s *Server) routeDir(path string) Store { + for _, r := range s.storeRoutes() { + if path+"/" == r.prefix { + return r.store() + } + } + return nil +} + // AddGroup adds a user to a group. func (s *Server) AddGroup(group, user string) { s.mu.Lock() @@ -478,39 +639,6 @@ func storeWrite(s store.Store, name string, data []byte) error { // readTimeout is the server-side deadline for all non-blocking store reads. const readTimeout = 10 * time.Second -// storeReadCtx is like storeRead but aborts if the context is cancelled or -// the read takes longer than readTimeout. -func storeReadCtx(ctx context.Context, s store.Store, name string) ([]byte, error) { - type result struct { - data []byte - err error - } - ch := make(chan result, 1) - go func() { - data, err := storeRead(s, name) - ch <- result{data, err} - }() - timer := time.NewTimer(readTimeout) - defer timer.Stop() - select { - case r := <-ch: - return r.data, r.err - case <-ctx.Done(): - return nil, ctx.Err() - case <-timer.C: - return nil, fmt.Errorf("read timeout") - } -} - -// storeBlockingRead opens an entry and performs a blocking read. -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 e.BlockingRead(ctx, base) -} - // Serve handles a single 9P connection. Each request is dispatched to its own // goroutine so blocking reads (e.g. *wait files) do not stall the serve loop. func (s *Server) Serve(conn net.Conn) { @@ -651,6 +779,19 @@ func (s *Server) pathType(path string) string { if path == "/" { return "dir" } + + // Check route table: directory match (e.g. "/a") or file match (e.g. "/a/foo"). + if st := s.routeDir(path); st != nil { + return "dir" + } + if st, name := s.routeStore(path); st != nil { + if _, err := st.Stat(name); err == nil { + return "file" + } + return "" + } + + // Session store. trimmed := strings.TrimPrefix(path, "/") parts := strings.SplitN(trimmed, "/", 3) switch { @@ -663,67 +804,12 @@ func (s *Server) pathType(path string) string { } return "file" } - case len(parts) == 1 && parts[0] == "a": - return "dir" - case len(parts) == 1 && parts[0] == "p": - return "dir" - case len(parts) == 1 && parts[0] == "m": - return "dir" - case len(parts) == 1 && parts[0] == "sk": - return "dir" - case len(parts) == 1 && parts[0] == "u": - return "dir" - case len(parts) == 2 && parts[0] == "u": - if _, err := s.utilStore.Stat(parts[1]); err == nil { - return "file" - } - case len(parts) == 1 && parts[0] == "x": - return "dir" - case len(parts) == 2 && parts[0] == "x": - if _, err := s.pluginStore.Stat(parts[1]); err == nil { - return "file" - } - case len(parts) == 1 && parts[0] == "tr": - return "dir" - case len(parts) == 2 && parts[0] == "tr": - if _, err := s.transcriptStore.Stat(parts[1]); err == nil { - return "file" - } - - case len(parts) == 1 && parts[0] == "tmp": - return "dir" - case len(parts) == 2 && parts[0] == "tmp": - if _, err := s.tmpStore.Stat(parts[1]); err == nil { - return "file" - } - case len(parts) == 1 && parts[0] == "backends": - return "file" - case len(parts) == 1 && parts[0] == "help": - return "file" - case len(parts) == 2 && parts[0] == "a": - if _, err := s.agentStore.Stat(parts[1]); err == nil { - return "file" - } - case len(parts) == 2 && parts[0] == "p": - if _, err := s.promptStore.Stat(parts[1]); err == nil { - return "file" - } - case len(parts) == 2 && parts[0] == "m": - if _, err := s.memStore.Stat(parts[1]); err == nil { - return "file" - } - case len(parts) == 2 && parts[0] == "sk": - if _, err := s.skillStore.Stat(parts[1]); err == nil { - return "file" - } case len(parts) == 3 && parts[0] == "s": if parts[2] == "t" { - // /s/{sid}/t is the tools directory if s.sessionStore.Session(parts[1]) != nil { return "dir" } } else if strings.HasPrefix(parts[2], "t/") { - // /s/{sid}/t/{rest} — tool file or subdir rel := strings.TrimPrefix(parts[2], "t/") if info, err := s.toolStore.Stat(rel); err == nil { if info.IsDir() { @@ -900,38 +986,24 @@ func (s *Server) create(cs *connState, fc *plan9.Fcall) *plan9.Fcall { newPath := pathJoin(f.path, fc.Name) s.log.Debug("Tcreate parent=%q name=%q", f.path, fc.Name) - // (e.g. touch) produces a real file. - switch f.path { - case "/a": - if err := s.agentStore.Create(fc.Name); err != nil { + + // Generic store-backed directories. + if st := s.routeDir(f.path); st != nil { + if err := st.Create(fc.Name); err != nil { return errFcall(fc, err.Error()) } - case "/m": - if err := s.memStore.Create(fc.Name); err != nil { - return errFcall(fc, err.Error()) - } - case "/tmp": - if err := s.tmpStore.Create(fc.Name); err != nil { - return errFcall(fc, err.Error()) - } - case "/sk": - if err := s.skillStore.Create(fc.Name); err != nil { - return errFcall(fc, err.Error()) - } - default: + } else if strings.HasPrefix(f.path, "/s/") { // Tool creates under session: /s/{sid}/t or /s/{sid}/t/... - if strings.HasPrefix(f.path, "/s/") { - parts := strings.SplitN(strings.TrimPrefix(f.path, "/s/"), "/", 3) - if (len(parts) == 2 && parts[1] == "t") || (len(parts) == 3 && parts[1] == "t") { - var rel string - if len(parts) == 2 { - rel = fc.Name - } else { - rel = parts[2] + "/" + fc.Name - } - if err := s.toolStore.Create(rel); err != nil { - return errFcall(fc, err.Error()) - } + parts := strings.SplitN(strings.TrimPrefix(f.path, "/s/"), "/", 3) + if (len(parts) == 2 && parts[1] == "t") || (len(parts) == 3 && parts[1] == "t") { + var rel string + if len(parts) == 2 { + rel = fc.Name + } else { + rel = parts[2] + "/" + fc.Name + } + if err := s.toolStore.Create(rel); err != nil { + return errFcall(fc, err.Error()) } } } @@ -957,197 +1029,68 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan if isDir { s.log.Debug("Tread dir path=%q offset=%d count=%d", path, fc.Offset, fc.Count) data := s.readDir(path, fc.Offset, fc.Count) - s.log.Debug("Rread dir path=%q len=%d", path, len(data)) return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data} } - // Fixed files directly under /s/ are served from the session store. - if isSessionStoreFile(path) { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.sessionStore, pathBase(path)) - if err != nil { - s.log.Debug("Rread path=%q err=%v", path, err) - return errFcall(fc, err.Error()) + s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) + entry, err := s.openEntry(path) + if err != nil { + return errFcall(fc, err.Error()) + } + + // OneShot entries yield data once per open; offset>0 means EOF. + if entry.OneShot() && fc.Offset > 0 { + return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0} + } + + // Blocking entries use BlockingRead with timeout and waitBase tracking. + if entry.IsBlocking() { + cs.mu.RLock() + f, fidOK := cs.fids[fc.Fid] + var base string + if fidOK { + base = f.waitBase } - s.log.Debug("Rread path=%q content_len=%d", path, len(content)) - return s.readSlice(fc, content) - } - - // backends is a static list of ollie-provided backends. - if path == "/backends" { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content := []byte(strings.Join(backend.Backends(), "\n") + "\n") - return s.readSlice(fc, content) - } - - // help is served from ~/.config/ollie/help.md. - if path == "/help" { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := os.ReadFile(s.helpPath()) + cs.mu.RUnlock() + waitCtx, waitCancel := context.WithTimeout(ctx, 5*time.Second) + defer waitCancel() + content, nextBase, err := entry.BlockingRead(waitCtx, base) if err != nil { return errFcall(fc, err.Error()) } - return s.readSlice(fc, content) - } - - // Agent config files are served from the agent store. - if strings.HasPrefix(path, "/a/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.agentStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Prompt files are served from the prompt store. - if strings.HasPrefix(path, "/p/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.promptStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Memory files are served from the memory store. - if strings.HasPrefix(path, "/m/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.memStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - - // Skill files are served from the skill store. - if strings.HasPrefix(path, "/sk/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.skillStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Util files are served from the util store. - if strings.HasPrefix(path, "/u/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.utilStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Plugin files are served from the plugin store. - if strings.HasPrefix(path, "/x/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.pluginStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Tmp files are served from the tmp store. - if strings.HasPrefix(path, "/tmp/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.tmpStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Transcript files are served from the transcript store. - if strings.HasPrefix(path, "/tr/") { - s.log.Debug("Tread path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.transcriptStore, pathBase(path)) - if err != nil { - return errFcall(fc, err.Error()) - } - return s.readSlice(fc, content) - } - - // Session files: /s/{id}/{file} - if strings.HasPrefix(path, "/s/") { - parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3) - if len(parts) == 3 { - // Tool files under session: /s/{id}/t/{tool} - if parts[2] == "t" || strings.HasPrefix(parts[2], "t/") { - rel := strings.TrimPrefix(parts[2], "t/") - if rel == "" { - // reading the directory itself as a file — shouldn't happen - return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0} - } - s.log.Debug("Tread tool path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - content, err := storeReadCtx(ctx, s.toolStore, rel) - if err != nil { - return errFcall(fc, err.Error()) - } - if rel == "idx" { - if allowed := s.sessionAllowTools(parts[1]); allowed != nil { - content = filterToolIndex(content, allowed) - } - } - return s.readSlice(fc, content) + if nextBase != "" { + cs.mu.Lock() + if f, ok := cs.fids[fc.Fid]; ok { + f.waitBase = nextBase } - s.log.Debug("Tread session file path=%q offset=%d count=%d", path, fc.Offset, fc.Count) - sfs, ok := s.sessionFileStore(parts[1]) - if !ok { - s.log.Debug("Rread session not found: %s", parts[1]) - return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0} - } - // fifo.out: non-zero offset is the trailing EOF read after a successful pop. - if parts[2] == "fifo.out" && fc.Offset > 0 { - return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0} - } - // *wait files block until a value changes; use connection context. - // One blocking read per open: a non-zero offset means the client - // already received data this open and is now polling for EOF. - if strings.HasSuffix(parts[2], "wait") { - if fc.Offset > 0 { - return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0} - } - cs.mu.RLock() - f, fidOK := cs.fids[fc.Fid] - var base string - if fidOK { - base = f.waitBase - } - cs.mu.RUnlock() - // Wrap with a short timeout so the FUSE read returns - // periodically, allowing pending signals (e.g. SIGINT) to - // be delivered to the blocked client process. - waitCtx, waitCancel := context.WithTimeout(ctx, 5*time.Second) - defer waitCancel() - content, nextBase, err := storeBlockingRead(sfs, parts[2], waitCtx, base) - if err != nil { - return errFcall(fc, err.Error()) - } - // Update the fid's baseline for subsequent reads. - if nextBase != "" { - cs.mu.Lock() - if f, ok := cs.fids[fc.Fid]; ok { - f.waitBase = nextBase - } - cs.mu.Unlock() - } - return s.readSlice(fc, content) - } - content, err := storeReadCtx(ctx, sfs, parts[2]) - if err != nil { - s.log.Debug("Rread session file err=%v", err) - return errFcall(fc, err.Error()) - } - s.log.Debug("Rread session file path=%q content_len=%d", path, len(content)) - return s.readSlice(fc, content) + cs.mu.Unlock() } + return s.readSlice(fc, content) + } + + // Normal read with timeout. + type result struct { + data []byte + err error + } + ch := make(chan result, 1) + go func() { + data, err := entry.Read() + ch <- result{data, err} + }() + timer := time.NewTimer(readTimeout) + defer timer.Stop() + select { + case r := <-ch: + if r.err != nil { + return errFcall(fc, r.err.Error()) + } + return s.readSlice(fc, r.data) + case <-ctx.Done(): + return errFcall(fc, ctx.Err().Error()) + case <-timer.C: + return errFcall(fc, "read timeout") } - s.log.Debug("Tread unhandled path=%q", path) - return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0} } // readSlice serves a byte slice at the requested offset/count. @@ -1225,35 +1168,26 @@ func (s *Server) wstat(cs *connState, fc *plan9.Fcall) *plan9.Fcall { return &plan9.Fcall{Type: plan9.Rwstat, Tag: fc.Tag} } - switch { - case strings.HasPrefix(f.path, "/a/"): - if err := s.agentStore.Rename(oldName, newDir.Name); err != nil { + // Generic store-backed files. + if st, _ := s.routeStore(f.path); st != nil { + if err := st.Rename(oldName, newDir.Name); err != nil { return errFcall(fc, err.Error()) } - cs.mu.Lock() - f.path = "/a/" + newDir.Name - f.qid.Path = qidPath(f.path) - cs.mu.Unlock() - - case strings.HasPrefix(f.path, "/m/"): - if err := s.memStore.Rename(oldName, newDir.Name); err != nil { - return errFcall(fc, err.Error()) + // Update fid path: find the prefix and replace the base. + for _, r := range s.storeRoutes() { + if strings.HasPrefix(f.path, r.prefix) { + cs.mu.Lock() + f.path = r.prefix + newDir.Name + f.qid.Path = qidPath(f.path) + cs.mu.Unlock() + break + } } - cs.mu.Lock() - f.path = "/m/" + newDir.Name - f.qid.Path = qidPath(f.path) - cs.mu.Unlock() + return &plan9.Fcall{Type: plan9.Rwstat, Tag: fc.Tag} + } - case strings.HasPrefix(f.path, "/sk/"): - if err := s.skillStore.Rename(oldName, newDir.Name); err != nil { - return errFcall(fc, err.Error()) - } - cs.mu.Lock() - f.path = "/sk/" + newDir.Name - f.qid.Path = qidPath(f.path) - cs.mu.Unlock() - - case strings.HasPrefix(f.path, "/s/"): + // Session paths. + if strings.HasPrefix(f.path, "/s/") { parts := strings.SplitN(strings.TrimPrefix(f.path, "/s/"), "/", 3) if len(parts) == 3 && parts[1] == "t" { // Tool rename under session: /s/{sid}/t/{rel} @@ -1273,9 +1207,6 @@ func (s *Server) wstat(cs *connState, fc *plan9.Fcall) *plan9.Fcall { return errFcall(fc, err.Error()) } } - - default: - return &plan9.Fcall{Type: plan9.Rwstat, Tag: fc.Tag} } return &plan9.Fcall{Type: plan9.Rwstat, Tag: fc.Tag} @@ -1301,7 +1232,11 @@ func (s *Server) clunk(cs *connState, fc *plan9.Fcall) *plan9.Fcall { if writable { s.log.Debug("Tclunk flush path=%q writeBuf=%d uname=%q", path, len(data), uname) input := strings.TrimSpace(string(data)) - if s.isAsyncWrite(path) { + async := false + if entry, err := s.openEntry(path); err == nil { + async = entry.Async() + } + if async { go s.handleWrite(path, input, uname) //nolint:errcheck } else if err := s.handleWrite(path, input, uname); err != nil { s.log.Debug("Tclunk handleWrite err=%v", err) @@ -1313,18 +1248,6 @@ func (s *Server) clunk(cs *connState, fc *plan9.Fcall) *plan9.Fcall { return &plan9.Fcall{Type: plan9.Rclunk, Tag: fc.Tag} } -// isAsyncWrite returns true for paths where writes may block (agent turns) -// and Rerror is not useful. -func (s *Server) isAsyncWrite(path string) bool { - if !strings.HasPrefix(path, "/s/") { - return false - } - switch pathBase(path) { - case "prompt", "fifo.in", "ctl": - return true - } - return false -} func (s *Server) remove(cs *connState, fc *plan9.Fcall) *plan9.Fcall { cs.mu.Lock() @@ -1339,27 +1262,20 @@ func (s *Server) remove(cs *connState, fc *plan9.Fcall) *plan9.Fcall { path := f.path s.log.Debug("Tremove path=%q", path) + + // Generic store-backed files. + if st, name := s.routeStore(path); st != nil { + if err := st.Delete(name); err != nil { + return errFcall(fc, err.Error()) + } + return &plan9.Fcall{Type: plan9.Rremove, Tag: fc.Tag} + } + var err error switch { - case strings.HasPrefix(path, "/a/"): - err = s.agentStore.Delete(pathBase(path)) - case strings.HasPrefix(path, "/m/"): - err = s.memStore.Delete(pathBase(path)) - case strings.HasPrefix(path, "/sk/"): - err = s.skillStore.Delete(pathBase(path)) - case strings.HasPrefix(path, "/tmp/"): - err = s.tmpStore.Delete(pathBase(path)) - case strings.HasPrefix(path, "/tr/"): - err = s.transcriptStore.Delete(pathBase(path)) case strings.HasPrefix(path, "/s/") && path != "/s/new": parts := strings.SplitN(strings.TrimPrefix(path, "/s/"), "/", 3) if len(parts) >= 2 && (parts[1] == "t" || strings.HasPrefix(parts[1], "t/")) { - // Tool remove: /s/{sid}/t/{tool} - // Reconstruct: parts[0]=sid, rest starts with t/... - // Actually with SplitN(..., 3): parts = [sid, "t", rest] or [sid, "t"] - // But we split on /s/ prefix first, so path="/s/sid/t/file" -> "sid/t/file" split into ["sid","t","file"] or ["sid","t/file"] - // Wait — SplitN("sid/t/file", "/", 3) = ["sid", "t", "file"] - // SplitN("sid/t/sub/file", "/", 3) = ["sid", "t", "sub/file"] if len(parts) == 2 && parts[1] == "t" { err = nil // can't remove the t/ dir itself } else if len(parts) == 3 && parts[1] == "t" { @@ -1368,7 +1284,6 @@ func (s *Server) remove(cs *connState, fc *plan9.Fcall) *plan9.Fcall { err = nil } } else if len(parts) >= 2 { - // Session file remove err = nil // synthetic file; let rm -r continue } else { err = s.sessionStore.Delete(parts[0]) @@ -1387,52 +1302,11 @@ func (s *Server) remove(cs *connState, fc *plan9.Fcall) *plan9.Fcall { // as a goroutine because they block for the entire agent turn). func (s *Server) handleWrite(path, input, uname string) error { s.log.Debug("handleWrite path=%q input_len=%d uname=%q", path, len(input), uname) - - if path == "/s/new" { - if input == "" { - return nil - } - return storeWrite(s.sessionStore, "new", []byte(input)) + st, name := s.routeStore(path) + if st == nil { + return nil } - - // Agent config writes go to the agent store. - if strings.HasPrefix(path, "/a/") { - return storeWrite(s.agentStore, pathBase(path), []byte(input)) - } - - // Memory file writes go to the memory store. - if strings.HasPrefix(path, "/m/") { - return storeWrite(s.memStore, pathBase(path), []byte(input)) - } - - // Skill file writes go to the skill store. - if strings.HasPrefix(path, "/sk/") { - return storeWrite(s.skillStore, pathBase(path), []byte(input)) - } - - // Tmp file writes go to the tmp store. - if strings.HasPrefix(path, "/tmp/") { - return storeWrite(s.tmpStore, pathBase(path), []byte(input)) - } - - // Session file writes: /s/{sessid}/{file} - if strings.HasPrefix(path, "/s/") { - parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3) - if len(parts) == 3 { - // Tool file writes: /s/{id}/t/{tool} - if strings.HasPrefix(parts[2], "t/") { - rel := strings.TrimPrefix(parts[2], "t/") - return storeWrite(s.toolStore, rel, []byte(input)) - } - sfs, ok := s.sessionFileStore(parts[1]) - if !ok { - return fmt.Errorf("session not found: %s", parts[1]) - } - return storeWrite(sfs, parts[2], []byte(input)) - } - } - - return nil + return storeWrite(st, name, []byte(input)) } // Shutdown kills all active sessions and batch jobs. @@ -1454,91 +1328,28 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte { if isDir { q.Type = QTDir } - return plan9.Dir{Qid: q, Mode: mode, Name: name, Uid: "ollie", Gid: "ollie", Muid: "ollie"} + uid, gid := s.fileOwnerGroup(fpath) + return plan9.Dir{Qid: q, Mode: mode, Name: name, Uid: uid, Gid: gid, Muid: uid} } if path == "/" { - dirs = append(dirs, makeDir("a", "/a", true, plan9.DMDIR|0755)) - dirs = append(dirs, makeDir("backends", "/backends", false, 0444)) - dirs = append(dirs, makeDir("help", "/help", false, 0444)) - dirs = append(dirs, makeDir("m", "/m", true, plan9.DMDIR|0755)) - dirs = append(dirs, makeDir("p", "/p", true, plan9.DMDIR|0750)) - dirs = append(dirs, makeDir("s", "/s", true, plan9.DMDIR|0555)) - dirs = append(dirs, makeDir("sk", "/sk", true, plan9.DMDIR|0555)) - dirs = append(dirs, makeDir("tmp", "/tmp", true, plan9.DMDIR|0755)) - dirs = append(dirs, makeDir("u", "/u", true, plan9.DMDIR|0755)) - dirs = append(dirs, makeDir("x", "/x", true, plan9.DMDIR|0555)) - dirs = append(dirs, makeDir("tr", "/tr", true, plan9.DMDIR|0555)) - } else if path == "/a" { - entries, _ := s.agentStore.List() + rootEntries := []string{"a", "backends", "help", "m", "p", "s", "sk", "tmp", "u", "x", "tr"} + for _, name := range rootEntries { + fpath := "/" + name + st := s.makeStat(fpath) + dirs = append(dirs, st) + } + } else if st := s.routeDir(path); st != nil { + entries, _ := st.List() for _, e := range entries { - if !e.IsDir() { - dirs = append(dirs, makeDir(e.Name(), "/a/"+e.Name(), false, 0666)) + fpath := path + "/" + e.Name() + d := s.makeStat(fpath) + if info, err := e.Info(); err == nil { + d.Atime = uint32(info.ModTime().Unix()) + d.Mtime = uint32(info.ModTime().Unix()) } + dirs = append(dirs, d) } - } else if path == "/p" { - entries, _ := s.promptStore.List() - for _, e := range entries { - if !e.IsDir() { - dirs = append(dirs, makeDir(e.Name(), "/p/"+e.Name(), false, 0640)) - } - } - } else if path == "/m" { - entries, _ := s.memStore.List() - for _, e := range entries { - if !e.IsDir() { - d := makeDir(e.Name(), "/m/"+e.Name(), false, 0666) - if info, err := e.Info(); err == nil { - d.Atime = uint32(info.ModTime().Unix()) - d.Mtime = uint32(info.ModTime().Unix()) - } - dirs = append(dirs, d) - } - } - } else if path == "/sk" { - entries, _ := s.skillStore.List() - for _, e := range entries { - mode := plan9.Perm(0666) - if e.Name() == "idx" { - mode = 0444 - } - dirs = append(dirs, makeDir(e.Name(), "/sk/"+e.Name(), false, mode)) - } - } else if path == "/tmp" { - entries, _ := s.tmpStore.List() - for _, e := range entries { - if !e.IsDir() { - d := makeDir(e.Name(), "/tmp/"+e.Name(), false, 0600) - if info, err := e.Info(); err == nil { - d.Atime = uint32(info.ModTime().Unix()) - d.Mtime = uint32(info.ModTime().Unix()) - } - dirs = append(dirs, d) - } - } - } else if path == "/tr" { - entries, _ := s.transcriptStore.List() - for _, e := range entries { - if !e.IsDir() { - d := makeDir(e.Name(), "/tr/"+e.Name(), false, 0444) - if info, err := e.Info(); err == nil { - d.Atime = uint32(info.ModTime().Unix()) - d.Mtime = uint32(info.ModTime().Unix()) - } - dirs = append(dirs, d) - } - } - } else if path == "/u" { - entries, _ := s.utilStore.List() - for _, e := range entries { - dirs = append(dirs, makeDir(e.Name(), "/u/"+e.Name(), false, 0555)) - } - } else if path == "/x" { - entries, _ := s.pluginStore.List() - for _, e := range entries { - dirs = append(dirs, makeDir(e.Name(), "/x/"+e.Name(), false, 0555)) - } - } else if path == "/s" { entries, _ := s.sessionStore.List() for _, e := range entries { @@ -1550,9 +1361,6 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte { dirs = append(dirs, makeDir(e.Name(), "/s/"+e.Name(), e.IsDir(), perm)) } } else if strings.HasPrefix(path, "/s/") { - // /s/{sid}/t — list tools - // /s/{sid}/t/{subdir} — list tool subdirectory - // /s/{sid} — list session files + t/ parts := strings.SplitN(strings.TrimPrefix(path, "/s/"), "/", 3) if len(parts) >= 2 && parts[1] == "t" { if len(parts) == 2 { @@ -1588,7 +1396,7 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte { dirs = append(dirs, makeDir(e.Name(), path+"/"+e.Name(), false, plan9.Perm(info.Mode()))) } } - dirs = append(dirs, makeDir("t", path+"/t", true, plan9.DMDIR|0777)) + dirs = append(dirs, makeDir("t", path+"/t", true, plan9.DMDIR|0500)) } } diff --git a/store/runfilestore.go b/store/runfilestore.go index 57df40c..f7d6a0d 100644 --- a/store/runfilestore.go +++ b/store/runfilestore.go @@ -8,12 +8,14 @@ import ( // 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) (content []byte, nextBase string, err error) // nil = not waitable - Size func() int64 // optional; if nil, len(Read()) + Name string + Mode os.FileMode + Read func() ([]byte, error) + Write func([]byte) error // nil = read-only + Wait func(ctx context.Context, base string) (content []byte, nextBase string, err error) // nil = not waitable + Size func() int64 // optional; if nil, len(Read()) + OneShot bool // true = yields data once per open + Async bool // true = writes dispatched asynchronously } // RunFileStore implements RunnableStore for any Runnable using a table of FileSpecs. @@ -101,5 +103,8 @@ func (rs *RunFileStore) open(name string) (StoreEntry, error) { ReadFn: spec.Read, WriteFn: writeFn, BlockingReadFn: waitFn, + OneShot_: spec.OneShot || spec.Wait != nil, + IsBlocking_: spec.Wait != nil, + Async_: spec.Async, }, nil } diff --git a/store/sessionfile.go b/store/sessionfile.go index cce821e..c185882 100644 --- a/store/sessionfile.go +++ b/store/sessionfile.go @@ -15,28 +15,30 @@ import ( // SessionFileList defines the fixed set of files in a session directory. var SessionFileList = []struct { - Name string - Mode os.FileMode + Name string + Mode os.FileMode + OneShot bool + Async bool }{ - {"plan", 0666}, - {"ctl", 0200}, - {"prompt", 0200}, - {"fifo.in", 0200}, - {"fifo.out", 0444}, - {"chat", 0666}, - {"offset", 0444}, - {"cfg", 0666}, - {"state", 0444}, - {"statewait", 0444}, - {"usage", 0444}, - {"cost", 0444}, - {"ctxsz", 0444}, - {"models", 0444}, - {"systemprompt", 0444}, - {"env", 0444}, - {"tail", 0555}, - {"prompt.prev", 0444}, - {"context", 0444}, + {"plan", 0666, false, false}, + {"ctl", 0200, false, true}, + {"prompt", 0200, false, true}, + {"fifo.in", 0200, false, true}, + {"fifo.out", 0444, true, false}, + {"chat", 0666, false, false}, + {"offset", 0444, false, false}, + {"cfg", 0666, false, false}, + {"state", 0444, false, false}, + {"statewait", 0444, false, false}, + {"usage", 0444, false, false}, + {"cost", 0444, false, false}, + {"ctxsz", 0444, false, false}, + {"models", 0444, false, false}, + {"systemprompt", 0444, false, false}, + {"env", 0444, false, false}, + {"tail", 0555, false, false}, + {"prompt.prev", 0444, false, false}, + {"context", 0444, false, false}, } // SessionFileStore is a RunFileStore for a session directory. @@ -47,6 +49,8 @@ func NewSessionFileStore(sess *Session, log *olog.Logger, kill func(), rename fu specs := make([]FileSpec, len(SessionFileList)) for i, f := range SessionFileList { specs[i] = h.fileSpec(f.Name, f.Mode) + specs[i].OneShot = f.OneShot + specs[i].Async = f.Async } return NewRunFileStore(sess, specs) } diff --git a/store/store.go b/store/store.go index 224e735..121215e 100644 --- a/store/store.go +++ b/store/store.go @@ -16,6 +16,12 @@ type StoreEntry interface { Read() ([]byte, error) Write(data []byte) error BlockingRead(ctx context.Context, base string) (content []byte, nextBase string, err error) + // OneShot returns true if the entry yields data once per open (offset>0 → EOF). + OneShot() bool + // IsBlocking returns true if Read blocks until data is available (use BlockingRead). + IsBlocking() bool + // Async returns true if writes should be dispatched asynchronously. + Async() bool } // Store is a named collection of entries. @@ -49,6 +55,9 @@ type EntryConfig struct { ReadFn func() ([]byte, error) WriteFn func([]byte) error BlockingReadFn func(context.Context, string) ([]byte, string, error) + OneShot_ bool + IsBlocking_ bool + Async_ bool } func (e *EntryConfig) Stat() (os.FileInfo, error) { return e.StatFn() } @@ -57,6 +66,9 @@ func (e *EntryConfig) Write(data []byte) error { func (e *EntryConfig) BlockingRead(ctx context.Context, base string) ([]byte, string, error) { return e.BlockingReadFn(ctx, base) } +func (e *EntryConfig) OneShot() bool { return e.OneShot_ } +func (e *EntryConfig) IsBlocking() bool { return e.IsBlocking_ } +func (e *EntryConfig) Async() bool { return e.Async_ } // storeConfig implements Store via function pointers. type storeConfig struct {