some fixes
This commit is contained in:
parent
13a64cecfe
commit
1075b3215c
4 changed files with 24 additions and 7 deletions
|
|
@ -49,7 +49,7 @@ func main() {
|
||||||
flag.IntVar(&max, "max", -1, "max parallel preloads")
|
flag.IntVar(&max, "max", -1, "max parallel preloads")
|
||||||
flag.IntVar(&bs, "bs", -1, "tune preload buffer size in bytes (default: 8192)")
|
flag.IntVar(&bs, "bs", -1, "tune preload buffer size in bytes (default: 8192)")
|
||||||
flag.Int64Var("a, "quota", 1, "max disk usage quota for the dst cache in GiB")
|
flag.Int64Var("a, "quota", 1, "max disk usage quota for the dst cache in GiB")
|
||||||
flag.BoolVar(&block, "block", false, "block until sftp is connected")
|
flag.BoolVar(&block, "block", true, "block until sftp is connected")
|
||||||
flag.Parse()
|
flag.Parse()
|
||||||
|
|
||||||
log = klogr.New().WithName("main")
|
log = klogr.New().WithName("main")
|
||||||
|
|
@ -82,6 +82,7 @@ func main() {
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
go func() {
|
go func() {
|
||||||
|
log.Info("starting main server", "addr", listenHttp)
|
||||||
err := s.ListenAndServe()
|
err := s.ListenAndServe()
|
||||||
if err != nil && err != http.ErrServerClosed {
|
if err != nil && err != http.ErrServerClosed {
|
||||||
klog.Fatalf("init failed: %s", err)
|
klog.Fatalf("init failed: %s", err)
|
||||||
|
|
@ -92,6 +93,7 @@ func main() {
|
||||||
cache *http.Server
|
cache *http.Server
|
||||||
)
|
)
|
||||||
if listenWebdav != "" {
|
if listenWebdav != "" {
|
||||||
|
log.Info("starting webdav server", "addr", listenWebdav)
|
||||||
dav = &http.Server{
|
dav = &http.Server{
|
||||||
Addr: listenWebdav,
|
Addr: listenWebdav,
|
||||||
Handler: srv.NewWebDavServer(
|
Handler: srv.NewWebDavServer(
|
||||||
|
|
@ -107,6 +109,7 @@ func main() {
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
if listenCache != "" {
|
if listenCache != "" {
|
||||||
|
log.Info("starting cache server", "addr", listenCache)
|
||||||
cache = &http.Server{
|
cache = &http.Server{
|
||||||
Addr: listenCache,
|
Addr: listenCache,
|
||||||
Handler: srv.NewCacheServer(
|
Handler: srv.NewCacheServer(
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ package fs
|
||||||
import (
|
import (
|
||||||
"cachefs/pkg/chunk"
|
"cachefs/pkg/chunk"
|
||||||
"cachefs/pkg/provider"
|
"cachefs/pkg/provider"
|
||||||
|
"cachefs/pkg/provider/sftp"
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
|
|
@ -80,19 +81,20 @@ func MetadataGenerator(file string, dst provider.FS, fstat, encrypted bool) erro
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
size := fi.Size()
|
size := fi.Size()
|
||||||
|
path = "/" + strings.TrimPrefix(path, dst.Root())
|
||||||
if fstat {
|
if fstat {
|
||||||
blocks, err := dst.Fstat(f.Fd())
|
blocks, err := dst.Fstat(f.Fd())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
if size <= blocks*512 {
|
if size <= blocks*512 {
|
||||||
md[strings.TrimPrefix(path, dst.Root())] = &Metadata{
|
md[path] = &Metadata{
|
||||||
Size: size,
|
Size: size,
|
||||||
Chunks: chunk.Chunks{{0, size}},
|
Chunks: chunk.Chunks{{0, size}},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
md[strings.TrimPrefix(path, dst.Root())] = &Metadata{
|
md[path] = &Metadata{
|
||||||
Size: size,
|
Size: size,
|
||||||
Chunks: chunk.Chunks{{0, size}},
|
Chunks: chunk.Chunks{{0, size}},
|
||||||
}
|
}
|
||||||
|
|
@ -332,7 +334,7 @@ func MapError(err error) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
func IsIOErr(err error) bool {
|
func IsIOErr(err error) bool {
|
||||||
return MapError(err) == IOErr
|
return MapError(err) == IOErr || errors.Is(err, sftp.ErrNotConnected)
|
||||||
}
|
}
|
||||||
|
|
||||||
type Metadata struct {
|
type Metadata struct {
|
||||||
|
|
|
||||||
|
|
@ -53,8 +53,12 @@ func NewPreloadHandler(ctx context.Context, fs *FS, max int, log logr.Logger) (p
|
||||||
func (ph *PreloadHandler) CancelPreload(name string) {
|
func (ph *PreloadHandler) CancelPreload(name string) {
|
||||||
ph.mu.Lock()
|
ph.mu.Lock()
|
||||||
defer ph.mu.Unlock()
|
defer ph.mu.Unlock()
|
||||||
if p, ok := ph.pm[name]; ok && p.Running {
|
p, ok := ph.pm[name]
|
||||||
p.cancel()
|
if ok {
|
||||||
|
if p.Running {
|
||||||
|
p.cancel()
|
||||||
|
}
|
||||||
|
delete(ph.pm, name)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -5,12 +5,15 @@ import (
|
||||||
"io"
|
"io"
|
||||||
"io/fs"
|
"io/fs"
|
||||||
"os"
|
"os"
|
||||||
|
"sync"
|
||||||
)
|
)
|
||||||
|
|
||||||
var _ provider.File = (*file)(nil)
|
var _ provider.File = (*file)(nil)
|
||||||
|
|
||||||
type file struct {
|
type file struct {
|
||||||
provider.File
|
provider.File
|
||||||
|
rmu sync.Mutex
|
||||||
|
wmu sync.Mutex
|
||||||
fs *FS
|
fs *FS
|
||||||
r *reader
|
r *reader
|
||||||
w *writer
|
w *writer
|
||||||
|
|
@ -140,10 +143,12 @@ func (f *file) Seek(offset int64, whence int) (n int64, err error) {
|
||||||
f.r.nonce = f.nonce
|
f.r.nonce = f.nonce
|
||||||
f.w.nonce = f.nonce
|
f.w.nonce = f.nonce
|
||||||
f.r.off = roff
|
f.r.off = roff
|
||||||
|
f.r.cn = cn
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *file) Read(p []byte) (n int, err error) {
|
func (f *file) Read(p []byte) (n int, err error) {
|
||||||
|
f.r.cn = -1
|
||||||
n, err = f.r.Read(p)
|
n, err = f.r.Read(p)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -154,6 +159,8 @@ func (f *file) Write(p []byte) (n int, err error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *file) ReadAt(p []byte, pos int64) (n int, err error) {
|
func (f *file) ReadAt(p []byte, pos int64) (n int, err error) {
|
||||||
|
f.rmu.Lock()
|
||||||
|
defer f.rmu.Unlock()
|
||||||
cn, _, roff := align(pos)
|
cn, _, roff := align(pos)
|
||||||
var last bool
|
var last bool
|
||||||
if f.r.cn != cn {
|
if f.r.cn != cn {
|
||||||
|
|
@ -165,7 +172,6 @@ func (f *file) ReadAt(p []byte, pos int64) (n int, err error) {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
f.r.cn = cn
|
|
||||||
}
|
}
|
||||||
n = copy(p, f.r.unread[roff:])
|
n = copy(p, f.r.unread[roff:])
|
||||||
if last && len(f.r.unread) == n+int(roff) {
|
if last && len(f.r.unread) == n+int(roff) {
|
||||||
|
|
@ -175,6 +181,8 @@ func (f *file) ReadAt(p []byte, pos int64) (n int, err error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *file) WriteAt(data []byte, pos int64) (n int, err error) {
|
func (f *file) WriteAt(data []byte, pos int64) (n int, err error) {
|
||||||
|
f.wmu.Lock()
|
||||||
|
defer f.wmu.Unlock()
|
||||||
cn, off, woff := align(pos)
|
cn, off, woff := align(pos)
|
||||||
c, ok := f.w.cm[cn]
|
c, ok := f.w.cm[cn]
|
||||||
if !ok {
|
if !ok {
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue