generic read/write dispatch via store entry flags
- Add OneShot, IsBlocking, Async to StoreEntry interface - FileSpec gains OneShot/Async fields; session entries annotated - read() is fully generic: openEntry → check flags → dispatch - handleWrite() reduced to 5 lines via routeStore - clunk uses entry.Async() instead of hardcoded isAsyncWrite - Root store for /backends and /help (consistent semantics) - ToolStore.OpenFiltered for per-session allowTools on idx - Remove storeReadCtx, storeBlockingRead, isAsyncWrite
This commit is contained in:
parent
000ad64673
commit
8f0dfdef37
690
server.go
690
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,27 +986,14 @@ 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
|
||||
|
|
@ -934,7 +1007,6 @@ func (s *Server) create(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
qid := plan9.Qid{Path: qidPath(newPath)}
|
||||
f.path = newPath
|
||||
|
|
@ -957,160 +1029,22 @@ 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("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())
|
||||
entry, err := s.openEntry(path)
|
||||
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)
|
||||
}
|
||||
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 {
|
||||
// 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
|
||||
|
|
@ -1118,16 +1052,12 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
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)
|
||||
content, nextBase, err := entry.BlockingRead(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 {
|
||||
|
|
@ -1137,17 +1067,30 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
}
|
||||
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())
|
||||
|
||||
// Normal read with timeout.
|
||||
type result struct {
|
||||
data []byte
|
||||
err error
|
||||
}
|
||||
s.log.Debug("Rread session file path=%q content_len=%d", path, len(content))
|
||||
return s.readSlice(fc, content)
|
||||
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())
|
||||
}
|
||||
// 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 = "/a/" + newDir.Name
|
||||
f.path = r.prefix + 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())
|
||||
break
|
||||
}
|
||||
cs.mu.Lock()
|
||||
f.path = "/m/" + newDir.Name
|
||||
f.qid.Path = qidPath(f.path)
|
||||
cs.mu.Unlock()
|
||||
|
||||
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()
|
||||
return &plan9.Fcall{Type: plan9.Rwstat, Tag: fc.Tag}
|
||||
}
|
||||
|
||||
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 == "" {
|
||||
st, name := s.routeStore(path)
|
||||
if st == nil {
|
||||
return nil
|
||||
}
|
||||
return storeWrite(s.sessionStore, "new", []byte(input))
|
||||
}
|
||||
|
||||
// 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))
|
||||
}
|
||||
}
|
||||
} 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)
|
||||
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 == "/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))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,6 +14,8 @@ type FileSpec struct {
|
|||
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -17,26 +17,28 @@ import (
|
|||
var SessionFileList = []struct {
|
||||
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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Reference in New Issue