Files
files/internal/storage/local/object.go
2026-09-09 16:42:21 +08:00

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() }