initial commit
This commit is contained in:
commit
3053f89410
39 changed files with 4567 additions and 0 deletions
52
pkg/fs/chunk.go
Normal file
52
pkg/fs/chunk.go
Normal file
|
|
@ -0,0 +1,52 @@
|
|||
// Copyright (C) 2022 Marius Schellenberger
|
||||
|
||||
package fs
|
||||
|
||||
import "sort"
|
||||
|
||||
type Chunk [2]int64
|
||||
|
||||
type Chunks []Chunk
|
||||
|
||||
func (cs Chunks) Len() int { return len(cs) }
|
||||
func (cs Chunks) Less(i, j int) bool { return cs[i][0] < cs[j][0] }
|
||||
func (cs Chunks) Swap(i, j int) { cs[i], cs[j] = cs[j], cs[i] }
|
||||
|
||||
func (cs Chunks) Exists(off int64, n int, size int64) bool {
|
||||
end := off + int64(n)
|
||||
if end > size {
|
||||
end = size
|
||||
}
|
||||
for _, c := range cs {
|
||||
if c[0] <= off && c[1] > off && c[1] >= end {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func (cs *Chunks) Add(off int64, n int) {
|
||||
*cs = append(*cs, Chunk{off, off + int64(n)})
|
||||
cs.merge()
|
||||
}
|
||||
|
||||
func (cs *Chunks) merge() {
|
||||
c := *cs
|
||||
if len(c) < 2 {
|
||||
return
|
||||
}
|
||||
sort.Sort(*cs)
|
||||
for i := 0; i < len(c); i++ {
|
||||
if i+1 == len(c) {
|
||||
break
|
||||
}
|
||||
if c[i][1] >= c[i+1][0] {
|
||||
if c[i+1][1] > c[i][1] {
|
||||
c[i][1] = c[i+1][1]
|
||||
}
|
||||
c = append(c[:i+1], c[i+2:]...)
|
||||
i--
|
||||
}
|
||||
}
|
||||
*cs = c
|
||||
}
|
||||
106
pkg/fs/file.go
Normal file
106
pkg/fs/file.go
Normal file
|
|
@ -0,0 +1,106 @@
|
|||
// Copyright (C) 2022 Marius Schellenberger
|
||||
|
||||
package fs
|
||||
|
||||
import (
|
||||
"io"
|
||||
"os"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
)
|
||||
|
||||
type File struct {
|
||||
log logr.Logger
|
||||
f *os.File
|
||||
md *Metadata
|
||||
offset int64
|
||||
}
|
||||
|
||||
type preload struct {
|
||||
f *File
|
||||
skipped int64
|
||||
written int64
|
||||
}
|
||||
|
||||
func (p *preload) Read(data []byte) (n int, err error) {
|
||||
if p.f.offset >= p.f.size() {
|
||||
return 0, io.EOF
|
||||
}
|
||||
n = len(data)
|
||||
if p.f.hasChunk(n) {
|
||||
p.skipped++
|
||||
_, err = p.f.Seek(int64(n), io.SeekCurrent)
|
||||
return
|
||||
}
|
||||
n, err = p.f.readWithoutCache(data)
|
||||
p.written++
|
||||
return
|
||||
}
|
||||
|
||||
func (f *File) Preload(unlock func()) {
|
||||
f.log.V(2).Info("preload started")
|
||||
p := &preload{f: f}
|
||||
_, err := io.Copy(io.Discard, p)
|
||||
if err != nil && err != io.EOF {
|
||||
f.log.Error(err, "error preloading file")
|
||||
}
|
||||
f.log.V(2).Info("preload finished", "skipped", p.skipped, "written", p.written)
|
||||
unlock()
|
||||
}
|
||||
|
||||
func (f *File) Read(p []byte) (n int, err error) {
|
||||
if f.hasChunk(len(p)) {
|
||||
n, err = f.md.ReadAt(p, f.offset)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_, err = f.Seek(int64(n), io.SeekCurrent)
|
||||
return
|
||||
}
|
||||
return f.readWithoutCache(p)
|
||||
}
|
||||
|
||||
func (f *File) size() int64 {
|
||||
return f.md.Size
|
||||
}
|
||||
|
||||
func (f *File) hasChunk(n int) bool {
|
||||
return f.md.HasChunk(f.offset, n)
|
||||
}
|
||||
|
||||
func (f *File) readWithoutCache(p []byte) (n int, err error) {
|
||||
n, err = f.f.Read(p)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
f.md.WriteAt(p, f.offset)
|
||||
f.md.AddChunk(f.offset, n)
|
||||
f.offset += int64(n)
|
||||
return
|
||||
}
|
||||
|
||||
func (f *File) Seek(offset int64, whence int) (n int64, err error) {
|
||||
n, err = f.f.Seek(offset, whence)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
switch whence {
|
||||
case io.SeekStart:
|
||||
f.offset = offset
|
||||
case io.SeekCurrent:
|
||||
f.offset += offset
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (f *File) Readdir(count int) ([]os.FileInfo, error) {
|
||||
return f.f.Readdir(count)
|
||||
}
|
||||
|
||||
func (f *File) Stat() (os.FileInfo, error) {
|
||||
return f.f.Stat()
|
||||
}
|
||||
|
||||
func (f *File) Close() error {
|
||||
return f.f.Close()
|
||||
}
|
||||
225
pkg/fs/fs.go
Normal file
225
pkg/fs/fs.go
Normal file
|
|
@ -0,0 +1,225 @@
|
|||
// Copyright (C) 2022 Marius Schellenberger
|
||||
|
||||
package fs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
stdfs "io/fs"
|
||||
"net/http"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
)
|
||||
|
||||
type FS struct {
|
||||
mu sync.RWMutex
|
||||
log logr.Logger
|
||||
NoCache http.Handler
|
||||
src string
|
||||
dst string
|
||||
mdf *os.File
|
||||
jenc *json.Encoder
|
||||
mm map[string]*Metadata
|
||||
cancel func()
|
||||
}
|
||||
|
||||
func NewFS(src, dst, metadata string, log logr.Logger) (fs *FS, err error) {
|
||||
if !filepath.IsAbs(src) {
|
||||
return nil, errors.New("src path is not absolute")
|
||||
}
|
||||
if !filepath.IsAbs(dst) {
|
||||
return nil, errors.New("dst path is not absolute")
|
||||
}
|
||||
if !filepath.IsAbs(metadata) {
|
||||
return nil, errors.New("metadata path is not absolute")
|
||||
}
|
||||
if src == dst {
|
||||
return nil, errors.New("src and dst path can not be equal")
|
||||
}
|
||||
mdf, err := os.OpenFile(metadata, os.O_RDWR|os.O_CREATE, 0o640)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
fs = &FS{
|
||||
log: log,
|
||||
NoCache: http.FileServer(http.Dir(src)),
|
||||
src: src,
|
||||
dst: dst,
|
||||
mdf: mdf,
|
||||
mm: make(map[string]*Metadata),
|
||||
jenc: json.NewEncoder(mdf),
|
||||
}
|
||||
err = json.NewDecoder(mdf).Decode(&fs.mm)
|
||||
if err == io.EOF {
|
||||
err = nil
|
||||
}
|
||||
for k := range fs.mm {
|
||||
_, err = stat(fs.dst, k)
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
delete(fs.mm, k)
|
||||
}
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
fs.cancel = cancel
|
||||
go fs.flusher(ctx)
|
||||
return fs, err
|
||||
}
|
||||
|
||||
func (fs *FS) Stat(name string) (stdfs.FileInfo, error) {
|
||||
return stat(fs.src, name)
|
||||
}
|
||||
|
||||
func (fs *FS) Open(name string) (http.File, error) {
|
||||
return fs.open(name)
|
||||
}
|
||||
|
||||
func (fs *FS) Preload(name string) {
|
||||
f, err := fs.open(name)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if f.md.Preload() {
|
||||
fs.log.WithValues("file", name).Error(errors.New("preload is running"), "error starting preload")
|
||||
return
|
||||
}
|
||||
unlock := func() { f.md.UnlockPreload() }
|
||||
go f.Preload(unlock)
|
||||
}
|
||||
|
||||
func skipLog(name string) bool {
|
||||
return strings.HasSuffix(name, "index.html")
|
||||
}
|
||||
|
||||
func (fs *FS) open(name string) (f *File, err error) {
|
||||
log := fs.log.WithValues("file", name)
|
||||
rp := filepath.Join(fs.src, filepath.FromSlash(path.Clean("/"+name)))
|
||||
rf, err := os.Open(rp)
|
||||
if err != nil {
|
||||
if !skipLog(name) {
|
||||
log.Error(err, "error opening source file")
|
||||
}
|
||||
return
|
||||
}
|
||||
rfi, err := rf.Stat()
|
||||
if err != nil {
|
||||
log.Error(err, "error stat source file")
|
||||
return
|
||||
}
|
||||
if rfi.IsDir() {
|
||||
return &File{f: rf}, nil
|
||||
}
|
||||
mp := filepath.Join(fs.dst, filepath.FromSlash(path.Clean("/"+name)))
|
||||
var mf *os.File
|
||||
i := strings.LastIndex(mp, "/")
|
||||
if i > 0 {
|
||||
dir := mp[:i]
|
||||
_, err = os.Stat(dir)
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
err = os.MkdirAll(dir, 0o755)
|
||||
if err != nil {
|
||||
log.Error(err, "error creating cache dir", "dir", mp)
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
}
|
||||
_, err = os.Stat(mp)
|
||||
truncate := errors.Is(err, os.ErrNotExist)
|
||||
if !fs.isOpen(name) {
|
||||
mf, err = os.OpenFile(mp, os.O_RDWR|os.O_CREATE, 0o644)
|
||||
if err != nil {
|
||||
log.Error(err, "error opening cache file")
|
||||
return nil, err
|
||||
}
|
||||
if truncate {
|
||||
err = mf.Truncate(rfi.Size())
|
||||
if err != nil {
|
||||
log.Error(err, "error truncating cache file")
|
||||
return nil, err
|
||||
}
|
||||
err = mf.Sync()
|
||||
if err != nil {
|
||||
log.Error(err, "error syncing cache file")
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
}
|
||||
return &File{
|
||||
log: log,
|
||||
f: rf,
|
||||
md: fs.metadata(name, rfi.Size(), mf),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (fs *FS) isOpen(name string) bool {
|
||||
fs.mu.RLock()
|
||||
defer fs.mu.RUnlock()
|
||||
md, ok := fs.mm[name]
|
||||
return ok && md.f != nil
|
||||
}
|
||||
|
||||
func (fs *FS) metadata(name string, size int64, f *os.File) (md *Metadata) {
|
||||
fs.mu.Lock()
|
||||
defer fs.mu.Unlock()
|
||||
if md, ok := fs.mm[name]; ok {
|
||||
if md.f == nil {
|
||||
md.f = f
|
||||
}
|
||||
return md
|
||||
}
|
||||
md = &Metadata{Size: size, f: f}
|
||||
fs.mm[name] = md
|
||||
return
|
||||
}
|
||||
|
||||
func (fs *FS) flusher(ctx context.Context) {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(time.Minute * 5):
|
||||
}
|
||||
fs.flushMetadata()
|
||||
}
|
||||
}
|
||||
|
||||
func (fs *FS) flushMetadata() {
|
||||
log := fs.log
|
||||
fs.mu.Lock()
|
||||
defer fs.mu.Unlock()
|
||||
_, err := fs.mdf.Seek(0, io.SeekStart)
|
||||
if err != nil {
|
||||
log.Error(err, "failure seeking metadata file")
|
||||
return
|
||||
}
|
||||
err = fs.mdf.Truncate(0)
|
||||
if err != nil {
|
||||
log.Error(err, "failure truncating metadata file")
|
||||
return
|
||||
}
|
||||
err = fs.jenc.Encode(fs.mm)
|
||||
if err != nil {
|
||||
log.Error(err, "failure flushing metadata file")
|
||||
}
|
||||
}
|
||||
|
||||
func (fs *FS) Close() {
|
||||
fs.cancel()
|
||||
for _, md := range fs.mm {
|
||||
md.f.Sync()
|
||||
md.f.Close()
|
||||
}
|
||||
fs.flushMetadata()
|
||||
}
|
||||
|
||||
func stat(prefix, name string) (stdfs.FileInfo, error) {
|
||||
p := filepath.Join(prefix, filepath.FromSlash(path.Clean("/"+name)))
|
||||
return os.Stat(p)
|
||||
}
|
||||
52
pkg/fs/metadata.go
Normal file
52
pkg/fs/metadata.go
Normal file
|
|
@ -0,0 +1,52 @@
|
|||
// Copyright (C) 2022 Marius Schellenberger
|
||||
|
||||
package fs
|
||||
|
||||
import (
|
||||
"os"
|
||||
"sync"
|
||||
)
|
||||
|
||||
type Metadata struct {
|
||||
mu sync.RWMutex `json:"-"`
|
||||
f *os.File `json:"-"`
|
||||
Size int64 `json:"s"`
|
||||
Chunks Chunks `json:"c"`
|
||||
preload bool `json:"-"`
|
||||
}
|
||||
|
||||
func (md *Metadata) HasChunk(off int64, n int) bool {
|
||||
md.mu.RLock()
|
||||
defer md.mu.RUnlock()
|
||||
return md.Chunks.Exists(off, n, md.Size)
|
||||
}
|
||||
|
||||
func (md *Metadata) AddChunk(off int64, n int) {
|
||||
md.mu.Lock()
|
||||
defer md.mu.Unlock()
|
||||
md.Chunks.Add(off, n)
|
||||
}
|
||||
|
||||
func (md *Metadata) ReadAt(p []byte, pos int64) (int, error) {
|
||||
return md.f.ReadAt(p, pos)
|
||||
}
|
||||
|
||||
func (md *Metadata) WriteAt(data []byte, pos int64) (int, error) {
|
||||
return md.f.WriteAt(data, pos)
|
||||
}
|
||||
|
||||
func (md *Metadata) Preload() bool {
|
||||
md.mu.Lock()
|
||||
defer md.mu.Unlock()
|
||||
if md.preload {
|
||||
return true
|
||||
}
|
||||
md.preload = true
|
||||
return false
|
||||
}
|
||||
|
||||
func (md *Metadata) UnlockPreload() {
|
||||
md.mu.Lock()
|
||||
defer md.mu.Unlock()
|
||||
md.preload = false
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue