153 lines
4.5 KiB
Go
153 lines
4.5 KiB
Go
package local
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
|
|
"git.apinb.com/ops/files/internal/storage"
|
|
)
|
|
|
|
func (c *Client) Open(ctx context.Context, objectKey string, byteRange *storage.ByteRange) (io.ReadCloser, error) {
|
|
path, err := c.resolveObjectPath(objectKey)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
file, err := os.Open(path)
|
|
if err != nil {
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil, storage.ErrObjectNotFound
|
|
}
|
|
return nil, fmt.Errorf("打开本地文件: %w", err)
|
|
}
|
|
if byteRange == nil {
|
|
return &contextReadCloser{ctx: ctx, reader: file, closer: file}, nil
|
|
}
|
|
info, err := file.Stat()
|
|
if err != nil {
|
|
_ = file.Close()
|
|
return nil, err
|
|
}
|
|
if byteRange.Start < 0 || byteRange.End < byteRange.Start || byteRange.End >= info.Size() {
|
|
_ = file.Close()
|
|
return nil, fmt.Errorf("文件读取范围无效")
|
|
}
|
|
if _, err := file.Seek(byteRange.Start, io.SeekStart); err != nil {
|
|
_ = file.Close()
|
|
return nil, err
|
|
}
|
|
return &contextReadCloser{ctx: ctx, reader: io.LimitReader(file, byteRange.Length()), closer: file}, nil
|
|
}
|
|
|
|
func (c *Client) Stat(ctx context.Context, objectKey string) (storage.ObjectInfo, error) {
|
|
path, err := c.resolveObjectPath(objectKey)
|
|
if err != nil {
|
|
return storage.ObjectInfo{}, err
|
|
}
|
|
file, err := os.Open(path)
|
|
if err != nil {
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return storage.ObjectInfo{}, storage.ErrObjectNotFound
|
|
}
|
|
return storage.ObjectInfo{}, fmt.Errorf("打开本地文件: %w", err)
|
|
}
|
|
defer file.Close()
|
|
info, err := file.Stat()
|
|
if err != nil {
|
|
return storage.ObjectInfo{}, err
|
|
}
|
|
hasher := sha256.New()
|
|
if _, err := io.Copy(hasher, &contextReader{ctx: ctx, reader: file}); err != nil {
|
|
return storage.ObjectInfo{}, fmt.Errorf("计算本地文件摘要: %w", err)
|
|
}
|
|
return storage.ObjectInfo{
|
|
Size: info.Size(),
|
|
StorageETag: "sha256:" + hex.EncodeToString(hasher.Sum(nil)),
|
|
LastModified: info.ModTime().UTC(),
|
|
}, nil
|
|
}
|
|
|
|
func (c *Client) Delete(_ context.Context, objectKey string) error {
|
|
path, err := c.resolveObjectPath(objectKey)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
|
|
return fmt.Errorf("删除本地文件: %w", err)
|
|
}
|
|
c.removeEmptyParents(filepath.Dir(path))
|
|
return nil
|
|
}
|
|
|
|
func (c *Client) removeEmptyParents(directory string) {
|
|
boundary := filepath.Clean(c.objectsPath)
|
|
for current := filepath.Clean(directory); current != boundary; current = filepath.Dir(current) {
|
|
relative, err := filepath.Rel(boundary, current)
|
|
if err != nil || relative == ".." || strings.HasPrefix(relative, ".."+string(filepath.Separator)) {
|
|
return
|
|
}
|
|
if err := os.Remove(current); err != nil {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func (c *Client) Status(_ context.Context) (storage.RuntimeStatus, error) {
|
|
diskStatus, err := c.diskStatus()
|
|
if err != nil {
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container()}, err
|
|
}
|
|
probe, err := os.CreateTemp(c.stagingPath, ".write-probe-")
|
|
if err != nil {
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container(), Disk: diskStatus}, fmt.Errorf("本地存储目录不可写: %w", err)
|
|
}
|
|
probePath := probe.Name()
|
|
removeProbe := func() { _ = os.Remove(probePath) }
|
|
defer removeProbe()
|
|
if err := probe.Chmod(0o640); err != nil {
|
|
_ = probe.Close()
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container(), Disk: diskStatus}, err
|
|
}
|
|
if _, err := probe.Write([]byte("files-storage-probe")); err != nil {
|
|
_ = probe.Close()
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container(), Disk: diskStatus}, err
|
|
}
|
|
if err := probe.Sync(); err != nil {
|
|
_ = probe.Close()
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container(), Disk: diskStatus}, err
|
|
}
|
|
if err := probe.Close(); err != nil {
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container(), Disk: diskStatus}, err
|
|
}
|
|
if err := os.Remove(probePath); err != nil {
|
|
return storage.RuntimeStatus{Provider: c.Provider(), Container: c.Container(), Disk: diskStatus}, err
|
|
}
|
|
return storage.RuntimeStatus{
|
|
Provider: c.Provider(),
|
|
Container: c.Container(),
|
|
Writable: diskStatus.Level != "critical",
|
|
Disk: diskStatus,
|
|
}, nil
|
|
}
|
|
|
|
type contextReadCloser struct {
|
|
ctx context.Context
|
|
reader io.Reader
|
|
closer io.Closer
|
|
}
|
|
|
|
func (r *contextReadCloser) Read(buffer []byte) (int, error) {
|
|
if err := r.ctx.Err(); err != nil {
|
|
return 0, err
|
|
}
|
|
return r.reader.Read(buffer)
|
|
}
|
|
|
|
func (r *contextReadCloser) Close() error { return r.closer.Close() }
|