server: use Open instead of JobStore/SessionFileStore
- FlatDirStore alias points to BlobStore - Server fields use BlobStore for flat stores - Replace JobStore/SessionFileStore with Open - Type-assert for Wait/LogInfo where needed - Rename store variable to sfs to avoid package shadowing
This commit is contained in:
parent
d5bc214bd0
commit
ee1710b2cb
66
server.go
66
server.go
|
|
@ -30,12 +30,10 @@ type (
|
|||
ReadWriteStore = store.ReadWriteStore
|
||||
BlobStore = store.BlobStore
|
||||
Store = store.Store
|
||||
FlatDirStore = store.Store
|
||||
FlatDirStore = store.BlobStore
|
||||
BatchStore = store.BatchStore
|
||||
BatchJobStore = store.BatchJobStore
|
||||
Session = store.Session
|
||||
SessionStore = store.SessionStore
|
||||
SessionFileStore = store.SessionFileStore
|
||||
|
||||
syntheticFileInfo = store.SyntheticFileInfo
|
||||
)
|
||||
|
|
@ -44,7 +42,7 @@ func NewFlatDirStore(dir string, perm os.FileMode) FlatDirStore {
|
|||
return store.NewFlatDir(dir, perm)
|
||||
}
|
||||
|
||||
func NewSkillStore() Store {
|
||||
func NewSkillStore() BlobStore {
|
||||
return store.NewSkillStore()
|
||||
}
|
||||
|
||||
|
|
@ -207,17 +205,17 @@ type Server struct {
|
|||
log *olog.Logger
|
||||
sink *olog.Sink
|
||||
agentsDir string // kept for session creation wiring
|
||||
agentStore Store
|
||||
agentStore BlobStore
|
||||
promptStore ReadableStore
|
||||
memStore Store
|
||||
toolStore Store
|
||||
utilStore Store
|
||||
pluginStore Store
|
||||
skillStore Store
|
||||
memStore BlobStore
|
||||
toolStore BlobStore
|
||||
utilStore BlobStore
|
||||
pluginStore BlobStore
|
||||
skillStore BlobStore
|
||||
sessionStore *SessionStore
|
||||
batchStore *BatchStore
|
||||
transcriptStore Store
|
||||
tmpStore Store
|
||||
transcriptStore BlobStore
|
||||
tmpStore BlobStore
|
||||
}
|
||||
|
||||
// New creates a new Server.
|
||||
|
|
@ -301,9 +299,13 @@ func defaultTmpDir() string {
|
|||
return paths.DataDir() + "/tmp"
|
||||
}
|
||||
|
||||
// sessionFileStore returns a SessionFileStore for the given session ID.
|
||||
func (s *Server) sessionFileStore(sessID string) (*SessionFileStore, bool) {
|
||||
return s.sessionStore.SessionFileStore(sessID)
|
||||
// sessionFileStore returns a Store for the given session ID.
|
||||
func (s *Server) sessionFileStore(sessID string) (Store, bool) {
|
||||
st, err := s.sessionStore.Open(sessID)
|
||||
if err != nil {
|
||||
return nil, false
|
||||
}
|
||||
return st, true
|
||||
}
|
||||
|
||||
// Serve handles a single 9P connection. Each request is dispatched to its own
|
||||
|
|
@ -498,7 +500,7 @@ func (s *Server) pathType(path string) string {
|
|||
}
|
||||
}
|
||||
case len(parts) == 3 && parts[0] == "b":
|
||||
if js, ok := s.batchStore.JobStore(parts[1]); ok {
|
||||
if js, err := s.batchStore.Open(parts[1]); err == nil {
|
||||
if _, err := js.Stat(parts[2]); err == nil {
|
||||
return "file"
|
||||
}
|
||||
|
|
@ -534,8 +536,8 @@ func (s *Server) pathType(path string) string {
|
|||
return "file"
|
||||
}
|
||||
case len(parts) == 3 && parts[0] == "s":
|
||||
if store, ok := s.sessionFileStore(parts[1]); ok {
|
||||
if _, err := store.Stat(parts[2]); err == nil {
|
||||
if sfs, ok := s.sessionFileStore(parts[1]); ok {
|
||||
if _, err := sfs.Stat(parts[2]); err == nil {
|
||||
return "file"
|
||||
}
|
||||
}
|
||||
|
|
@ -730,8 +732,8 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
|
||||
if len(parts) == 3 && parts[0] == "b" {
|
||||
s.log.Debug("Tread batch file path=%q offset=%d count=%d", path, fc.Offset, fc.Count)
|
||||
js, ok := s.batchStore.JobStore(parts[1])
|
||||
if !ok {
|
||||
js, err := s.batchStore.Open(parts[1])
|
||||
if err != nil {
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
}
|
||||
if parts[2] == "statewait" {
|
||||
|
|
@ -747,7 +749,7 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
cs.mu.RUnlock()
|
||||
waitCtx, waitCancel := context.WithTimeout(ctx, 5*time.Second)
|
||||
defer waitCancel()
|
||||
content, err := js.Wait(waitCtx, parts[2], base)
|
||||
content, err := js.(*store.BatchJobStore).Wait(waitCtx, parts[2], base)
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
|
|
@ -893,7 +895,7 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
|
||||
if len(parts) == 3 {
|
||||
s.log.Debug("Tread session file path=%q offset=%d count=%d", path, fc.Offset, fc.Count)
|
||||
store, ok := s.sessionFileStore(parts[1])
|
||||
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}
|
||||
|
|
@ -921,7 +923,7 @@ 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 := store.Wait(waitCtx, parts[2], base)
|
||||
content, err := sfs.(*store.SessionFileStore).Wait(waitCtx, parts[2], base)
|
||||
if err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
|
|
@ -935,7 +937,7 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall, ctx context.Context) *plan
|
|||
}
|
||||
return s.readSlice(fc, content)
|
||||
}
|
||||
content, err := store.Get(parts[2])
|
||||
content, err := sfs.Get(parts[2])
|
||||
if err != nil {
|
||||
s.log.Debug("Rread session file err=%v", err)
|
||||
return errFcall(fc, err.Error())
|
||||
|
|
@ -1220,11 +1222,11 @@ func (s *Server) handleWrite(path, input string) error {
|
|||
if strings.HasPrefix(path, "/s/") {
|
||||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
|
||||
if len(parts) == 3 {
|
||||
store, ok := s.sessionFileStore(parts[1])
|
||||
sfs, ok := s.sessionFileStore(parts[1])
|
||||
if !ok {
|
||||
return fmt.Errorf("session not found: %s", parts[1])
|
||||
}
|
||||
return store.Put(parts[2], []byte(input))
|
||||
return sfs.Put(parts[2], []byte(input))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1364,7 +1366,7 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte {
|
|||
}
|
||||
} else if strings.HasPrefix(path, "/b/") {
|
||||
// /b/{id} directory listing
|
||||
if js, ok := s.batchStore.JobStore(pathBase(path)); ok {
|
||||
if js, err := s.batchStore.Open(pathBase(path)); err == nil {
|
||||
entries, _ := js.List()
|
||||
for _, e := range entries {
|
||||
info, _ := e.Info()
|
||||
|
|
@ -1384,8 +1386,8 @@ func (s *Server) readDir(path string, offset uint64, count uint32) []byte {
|
|||
} else {
|
||||
// Session subdirectory: /s/{sessid}
|
||||
sessID := pathBase(path)
|
||||
if store, ok := s.sessionFileStore(sessID); ok {
|
||||
entries, _ := store.List()
|
||||
if sfs, ok := s.sessionFileStore(sessID); ok {
|
||||
entries, _ := sfs.List()
|
||||
for _, e := range entries {
|
||||
info, _ := e.Info()
|
||||
dirs = append(dirs, makeDir(e.Name(), path+"/"+e.Name(), false, plan9.Perm(info.Mode())))
|
||||
|
|
@ -1515,9 +1517,9 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
if strings.HasPrefix(path, "/b/") {
|
||||
parts := strings.SplitN(strings.TrimPrefix(path, "/"), "/", 3)
|
||||
if len(parts) == 3 && parts[0] == "b" {
|
||||
if js, ok := s.batchStore.JobStore(parts[1]); ok {
|
||||
if js, err := s.batchStore.Open(parts[1]); err == nil {
|
||||
if base == "log" {
|
||||
length, vers := js.LogInfo()
|
||||
length, vers := js.(*store.BatchJobStore).LogInfo()
|
||||
dir.Length = uint64(length)
|
||||
dir.Qid.Vers = vers
|
||||
}
|
||||
|
|
@ -1595,7 +1597,7 @@ func (s *Server) makeStat(path string) plan9.Dir {
|
|||
case strings.HasPrefix(path, "/b/"):
|
||||
parts := strings.SplitN(strings.TrimPrefix(path, "/b/"), "/", 2)
|
||||
if len(parts) == 2 {
|
||||
if js, ok := s.batchStore.JobStore(parts[0]); ok {
|
||||
if js, err := s.batchStore.Open(parts[0]); err == nil {
|
||||
if info, err := js.Stat(parts[1]); err == nil {
|
||||
dir.Length = uint64(info.Size())
|
||||
}
|
||||
|
|
|
|||
Reference in New Issue