diff --git a/pkg/blob/file.go b/pkg/blob/file.go new file mode 100644 index 00000000..fc57688f --- /dev/null +++ b/pkg/blob/file.go @@ -0,0 +1,229 @@ +package blob + +import ( + "context" + "errors" + "fmt" + "net/url" + "os" + "path/filepath" + "strings" + + "github.com/google/uuid" +) + +// FileStore is a blob.Store backed by the local filesystem rooted at a +// configured directory. Writes go to a hidden staging area (root + +// stagingDir) and are atomically renamed to their final key on Complete. +type FileStore struct { + root string + stagingDir string +} + +// stagingDir is the per-store hidden subdir under root that holds +// in-progress writes. The leading dot keeps it out of plain blob key +// space — keys with a leading dot are technically allowed, but uncommon. +const stagingDir = ".staging" + +// NewFileStore returns a FileStore rooted at root. The root and the +// staging directory inside it are created on construction if they +// don't already exist. +func NewFileStore(root string) (*FileStore, error) { + abs, err := filepath.Abs(root) + if err != nil { + return nil, fmt.Errorf("blob: abs root %q: %w", root, err) + } + staging := filepath.Join(abs, stagingDir) + if err := os.MkdirAll(staging, 0o755); err != nil { + return nil, fmt.Errorf("blob: mkdir staging %q: %w", staging, err) + } + return &FileStore{root: abs, stagingDir: staging}, nil +} + +// Root returns the absolute path of the FileStore's root, useful for +// callers that need to construct keys from absolute paths. +func (s *FileStore) Root() string { return s.root } + +func (s *FileStore) URL(key string) string { + return "file://" + s.absPath(key) +} + +func (s *FileStore) absPath(key string) string { + return filepath.Join(s.root, filepath.FromSlash(key)) +} + +func (s *FileStore) Open(ctx context.Context, key string) (Reader, error) { + path := s.absPath(key) + f, err := os.Open(path) + if err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil, fmt.Errorf("%w: %s", ErrNotFound, path) + } + return nil, fmt.Errorf("blob file open %s: %w", path, err) + } + st, err := f.Stat() + if err != nil { + _ = f.Close() + return nil, fmt.Errorf("blob file stat %s: %w", path, err) + } + return &fileReader{f: f, size: st.Size()}, nil +} + +func (s *FileStore) NewWriter(ctx context.Context, key, contentType string) (Writer, error) { + if err := s.ensureKeyDir(key); err != nil { + return nil, err + } + id, err := uuid.NewV7() + if err != nil { + return nil, fmt.Errorf("blob file: generate staging id: %w", err) + } + staging := filepath.Join(s.stagingDir, id.String()) + f, err := os.OpenFile(staging, os.O_CREATE|os.O_WRONLY|os.O_EXCL, 0o644) + if err != nil { + return nil, fmt.Errorf("blob file: open staging %s: %w", staging, err) + } + return &fileWriter{ + f: f, + staging: staging, + finalKey: key, + finalAbs: s.absPath(key), + }, nil +} + +// ensureKeyDir creates any parent directories implied by key so the +// final rename has somewhere to land. +func (s *FileStore) ensureKeyDir(key string) error { + dir := filepath.Dir(s.absPath(key)) + if err := os.MkdirAll(dir, 0o755); err != nil { + return fmt.Errorf("blob file: mkdir %s: %w", dir, err) + } + return nil +} + +func (s *FileStore) Move(ctx context.Context, srcKey, dstKey string) error { + if err := s.ensureKeyDir(dstKey); err != nil { + return err + } + src, dst := s.absPath(srcKey), s.absPath(dstKey) + err := os.Rename(src, dst) + if err != nil { + if errors.Is(err, os.ErrNotExist) { + // Idempotency: if src is gone, maybe a previous Move + // already succeeded. Check whether the destination exists; + // if it does, we're in the post-condition the caller wants. + if _, statErr := os.Stat(dst); statErr == nil { + return nil + } + return fmt.Errorf("%w: %s", ErrNotFound, src) + } + return fmt.Errorf("blob file rename %s -> %s: %w", src, dst, err) + } + return nil +} + +func (s *FileStore) Delete(ctx context.Context, key string) error { + path := s.absPath(key) + if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("blob file delete %s: %w", path, err) + } + return nil +} + +// ParseLocation accepts either a plain absolute path under root or a +// file:// URL, and returns the corresponding key. Locations outside +// root are rejected with ok=false. +func (s *FileStore) ParseLocation(location string) (string, bool) { + abs := location + if strings.HasPrefix(location, "file://") { + u, err := url.Parse(location) + if err != nil { + return "", false + } + abs = u.Path + } + if !filepath.IsAbs(abs) { + // Treat as already-relative key; just sanity-check it doesn't + // escape via "..". + clean := filepath.Clean(abs) + if strings.HasPrefix(clean, "..") { + return "", false + } + return filepath.ToSlash(clean), true + } + rel, err := filepath.Rel(s.root, abs) + if err != nil { + return "", false + } + if strings.HasPrefix(rel, "..") { + return "", false + } + return filepath.ToSlash(rel), true +} + +// fileReader implements blob.Reader over an open *os.File. +type fileReader struct { + f *os.File + size int64 +} + +func (r *fileReader) ReadAt(p []byte, off int64) (int, error) { return r.f.ReadAt(p, off) } +func (r *fileReader) Close() error { return r.f.Close() } +func (r *fileReader) Size() int64 { return r.size } + +// fileWriter implements blob.Writer atop a staging temp file that gets +// renamed on Complete. +type fileWriter struct { + f *os.File + staging string + finalKey string + finalAbs string + finalized bool +} + +func (w *fileWriter) Write(p []byte) (int, error) { + if w.finalized { + return 0, fmt.Errorf("blob file: write after Complete/Abort on %s", w.finalKey) + } + return w.f.Write(p) +} + +func (w *fileWriter) Complete() error { + if w.finalized { + return fmt.Errorf("blob file: Complete called twice on %s", w.finalKey) + } + if err := w.f.Sync(); err != nil { + return fmt.Errorf("blob file: fsync %s: %w", w.staging, err) + } + if err := w.f.Close(); err != nil { + return fmt.Errorf("blob file: close staging %s: %w", w.staging, err) + } + if err := os.Rename(w.staging, w.finalAbs); err != nil { + // Best-effort cleanup so an aborted rename doesn't leave the + // staging file lying around. + _ = os.Remove(w.staging) + return fmt.Errorf("blob file: rename staging -> %s: %w", w.finalAbs, err) + } + w.finalized = true + return nil +} + +func (w *fileWriter) Abort() error { + if w.finalized { + return nil + } + w.finalized = true + _ = w.f.Close() + if err := os.Remove(w.staging); err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("blob file: remove staging %s: %w", w.staging, err) + } + return nil +} + +func (w *fileWriter) Close() error { return w.Abort() } + +// Compile-time interface checks. +var ( + _ Store = (*FileStore)(nil) + _ Reader = (*fileReader)(nil) + _ Writer = (*fileWriter)(nil) +) diff --git a/pkg/blob/file_test.go b/pkg/blob/file_test.go new file mode 100644 index 00000000..119ed7bd --- /dev/null +++ b/pkg/blob/file_test.go @@ -0,0 +1,207 @@ +package blob + +import ( + "context" + "errors" + "io" + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" +) + +func TestFileStore_WriteCompleteReadDeleteMove(t *testing.T) { + ctx := context.Background() + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + + const key = "vod/abc/def.fmp4" + payload := []byte("hello blob storage world") + + w, err := store.NewWriter(ctx, key, "video/mp4") + require.NoError(t, err) + n, err := w.Write(payload) + require.NoError(t, err) + require.Equal(t, len(payload), n) + // Before Complete, the final key is not visible. + _, err = store.Open(ctx, key) + require.ErrorIs(t, err, ErrNotFound) + + require.NoError(t, w.Complete()) + // Double-complete is rejected. + require.Error(t, w.Complete()) + // Close after Complete is a no-op. + require.NoError(t, w.Close()) + + // Random-access read. + r, err := store.Open(ctx, key) + require.NoError(t, err) + require.Equal(t, int64(len(payload)), r.Size()) + + buf := make([]byte, 5) + got, err := r.ReadAt(buf, 6) + require.NoError(t, err) + require.Equal(t, 5, got) + require.Equal(t, "blob ", string(buf)) + require.NoError(t, r.Close()) + + // Move atomically renames; the source key disappears. + const dstKey = "vod/xyz/moved.fmp4" + require.NoError(t, store.Move(ctx, key, dstKey)) + _, err = store.Open(ctx, key) + require.ErrorIs(t, err, ErrNotFound) + r2, err := store.Open(ctx, dstKey) + require.NoError(t, err) + require.NoError(t, r2.Close()) + + // Move is idempotent: re-running with src already gone returns nil + // as long as the destination still exists. + require.NoError(t, store.Move(ctx, key, dstKey)) + + // Delete removes; Open afterwards is ErrNotFound. + require.NoError(t, store.Delete(ctx, dstKey)) + _, err = store.Open(ctx, dstKey) + require.ErrorIs(t, err, ErrNotFound) + // Delete of missing blob is also nil. + require.NoError(t, store.Delete(ctx, dstKey)) +} + +func TestFileStore_AbortLeavesNoFinalBlob(t *testing.T) { + ctx := context.Background() + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + + const key = "vod/aborted.fmp4" + w, err := store.NewWriter(ctx, key, "") + require.NoError(t, err) + _, err = w.Write([]byte("incoming bytes")) + require.NoError(t, err) + require.NoError(t, w.Abort()) + + _, err = store.Open(ctx, key) + require.ErrorIs(t, err, ErrNotFound) +} + +func TestFileStore_CloseWithoutCompleteAborts(t *testing.T) { + ctx := context.Background() + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + + const key = "vod/forgotten.fmp4" + w, err := store.NewWriter(ctx, key, "") + require.NoError(t, err) + _, err = w.Write([]byte("oops never finalized")) + require.NoError(t, err) + require.NoError(t, w.Close()) + + _, err = store.Open(ctx, key) + require.ErrorIs(t, err, ErrNotFound) + + // And the staging directory should have no residue. + stagingEntries, err := os.ReadDir(filepath.Join(store.Root(), stagingDir)) + require.NoError(t, err) + require.Empty(t, stagingEntries, "staging dir should be empty after Close-without-Complete") +} + +func TestFileStore_ParseLocation(t *testing.T) { + root := t.TempDir() + store, err := NewFileStore(root) + require.NoError(t, err) + + t.Run("absolute path inside root", func(t *testing.T) { + key, ok := store.ParseLocation(filepath.Join(root, "uploads/abc")) + require.True(t, ok) + require.Equal(t, "uploads/abc", key) + }) + t.Run("file:// URL inside root", func(t *testing.T) { + key, ok := store.ParseLocation("file://" + filepath.Join(root, "uploads/abc")) + require.True(t, ok) + require.Equal(t, "uploads/abc", key) + }) + t.Run("relative key", func(t *testing.T) { + key, ok := store.ParseLocation("uploads/abc") + require.True(t, ok) + require.Equal(t, "uploads/abc", key) + }) + t.Run("absolute path outside root", func(t *testing.T) { + _, ok := store.ParseLocation("/etc/passwd") + require.False(t, ok) + }) + t.Run("relative path attempting escape", func(t *testing.T) { + _, ok := store.ParseLocation("../etc/passwd") + require.False(t, ok) + }) +} + +func TestFileStore_URL(t *testing.T) { + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + url := store.URL("vod/abc.fmp4") + require.Contains(t, url, "file://") + require.Contains(t, url, "vod/abc.fmp4") +} + +func TestFileStore_ConcurrentReads(t *testing.T) { + // Multiple Readers against the same blob should be able to operate + // independently, since each Open returns a fresh *os.File. + ctx := context.Background() + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + + const key = "concurrent.bin" + payload := make([]byte, 4096) + for i := range payload { + payload[i] = byte(i % 251) + } + w, err := store.NewWriter(ctx, key, "") + require.NoError(t, err) + _, err = w.Write(payload) + require.NoError(t, err) + require.NoError(t, w.Complete()) + + r1, err := store.Open(ctx, key) + require.NoError(t, err) + defer r1.Close() + r2, err := store.Open(ctx, key) + require.NoError(t, err) + defer r2.Close() + + buf1 := make([]byte, 1024) + buf2 := make([]byte, 512) + + n1, err := r1.ReadAt(buf1, 100) + require.NoError(t, err) + require.Equal(t, 1024, n1) + + n2, err := r2.ReadAt(buf2, 3000) + require.NoError(t, err) + require.Equal(t, 512, n2) + + require.Equal(t, payload[100:1124], buf1) + require.Equal(t, payload[3000:3512], buf2) +} + +func TestFileStore_OpenMissingReturnsErrNotFound(t *testing.T) { + ctx := context.Background() + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + + _, err = store.Open(ctx, "does/not/exist") + require.True(t, errors.Is(err, ErrNotFound)) +} + +func TestFileStore_WriterCloseAndAbortIdempotent(t *testing.T) { + ctx := context.Background() + store, err := NewFileStore(t.TempDir()) + require.NoError(t, err) + + w, err := store.NewWriter(ctx, "x", "") + require.NoError(t, err) + _, err = io.WriteString(w, "hi") + require.NoError(t, err) + require.NoError(t, w.Abort()) + // Second Abort is a no-op (returns nil). + require.NoError(t, w.Abort()) + require.NoError(t, w.Close()) +} diff --git a/pkg/blob/s3.go b/pkg/blob/s3.go new file mode 100644 index 00000000..caf79015 --- /dev/null +++ b/pkg/blob/s3.go @@ -0,0 +1,203 @@ +package blob + +import ( + "context" + "errors" + "fmt" + "strings" + + "github.com/aws/aws-sdk-go-v2/aws" + awss3 "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" + "github.com/google/uuid" + + s3pkg "stream.place/streamplace/pkg/s3" +) + +// S3Store is a blob.Store backed by an S3-compatible object store. +// +// Writes go to a hidden staging prefix (.staging/) so that +// in-progress multipart uploads can't collide with the final +// content-addressed key. Complete renames staging -> the configured +// key via CopyObject + DeleteObject. +type S3Store struct { + client *awss3.Client + bucket string +} + +// stagingPrefix is the in-bucket prefix used for in-progress writes. +// Matches the FileStore convention (.staging at the root). +const s3StagingPrefix = ".staging/" + +// NewS3Store wraps an existing S3 client + bucket as a blob.Store. +func NewS3Store(client *awss3.Client, bucket string) *S3Store { + return &S3Store{client: client, bucket: bucket} +} + +func (s *S3Store) URL(key string) string { return "s3://" + s.bucket + "/" + key } + +func (s *S3Store) Bucket() string { return s.bucket } + +func (s *S3Store) Open(ctx context.Context, key string) (Reader, error) { + ra, err := s3pkg.NewReaderAt(ctx, s.client, s.bucket, key) + if err != nil { + // The AWS SDK returns *types.NoSuchKey wrapped in opaque + // generic errors; sniff the message rather than relying on + // type assertions that change between SDK versions. + if isS3NotFound(err) { + return nil, fmt.Errorf("%w: s3://%s/%s", ErrNotFound, s.bucket, key) + } + return nil, err + } + return ra, nil +} + +func (s *S3Store) NewWriter(ctx context.Context, key, contentType string) (Writer, error) { + id, err := uuid.NewV7() + if err != nil { + return nil, fmt.Errorf("blob s3: generate staging id: %w", err) + } + stagingKey := s3StagingPrefix + id.String() + + mw, err := s3pkg.NewMultipartWriter(ctx, s.client, s.bucket, stagingKey, contentType) + if err != nil { + return nil, fmt.Errorf("blob s3: create staging multipart: %w", err) + } + return &s3Writer{ + store: s, + ctx: ctx, + mw: mw, + stagingKey: stagingKey, + finalKey: key, + }, nil +} + +func (s *S3Store) Move(ctx context.Context, srcKey, dstKey string) error { + _, err := s.client.CopyObject(ctx, &awss3.CopyObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(dstKey), + CopySource: aws.String(s.bucket + "/" + srcKey), + }) + if err != nil { + if isS3NotFound(err) { + // Idempotency: maybe a previous Move already renamed + // source -> dest. If dest exists, we're done. + if _, headErr := s.client.HeadObject(ctx, &awss3.HeadObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(dstKey), + }); headErr == nil { + return nil + } + return fmt.Errorf("%w: s3://%s/%s", ErrNotFound, s.bucket, srcKey) + } + return fmt.Errorf("blob s3: copy %s -> %s: %w", srcKey, dstKey, err) + } + if _, err := s.client.DeleteObject(ctx, &awss3.DeleteObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(srcKey), + }); err != nil { + // Non-fatal: the dst exists; staging is leftover but + // eventually garbage-collected. + return fmt.Errorf("blob s3: delete src %s after copy to %s: %w", srcKey, dstKey, err) + } + return nil +} + +func (s *S3Store) Delete(ctx context.Context, key string) error { + _, err := s.client.DeleteObject(ctx, &awss3.DeleteObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(key), + }) + if err != nil && !isS3NotFound(err) { + return fmt.Errorf("blob s3 delete %s: %w", key, err) + } + return nil +} + +// ParseLocation accepts "s3://bucket/key" URLs. The bucket must match +// the Store's bucket; otherwise we return ok=false to avoid silently +// reading from a different bucket than the Store is configured for. +func (s *S3Store) ParseLocation(location string) (string, bool) { + if !strings.HasPrefix(location, "s3://") { + return "", false + } + rest := strings.TrimPrefix(location, "s3://") + slash := strings.IndexByte(rest, '/') + if slash <= 0 || slash == len(rest)-1 { + return "", false + } + bucket, key := rest[:slash], rest[slash+1:] + if bucket != s.bucket { + return "", false + } + return key, true +} + +// s3Writer adapts the existing s3.MultipartWriter (which writes to one +// fixed key) to the blob.Writer interface that needs Complete to +// publish at the originally-requested final key. It does this with a +// staging-then-Move pattern that mirrors what pkg/vod already did +// inline. +type s3Writer struct { + store *S3Store + ctx context.Context + mw *s3pkg.MultipartWriter + stagingKey string + finalKey string + finalized bool +} + +func (w *s3Writer) Write(p []byte) (int, error) { + if w.finalized { + return 0, fmt.Errorf("blob s3: write after Complete/Abort on %s", w.finalKey) + } + return w.mw.Write(p) +} + +func (w *s3Writer) Complete() error { + if w.finalized { + return fmt.Errorf("blob s3: Complete called twice on %s", w.finalKey) + } + if err := w.mw.Complete(); err != nil { + return err + } + if err := w.store.Move(w.ctx, w.stagingKey, w.finalKey); err != nil { + // We have a completed staging upload but the rename failed. + // Try to remove the staging blob so it doesn't dangle; if that + // fails the staging janitor (if/when one exists) will pick it + // up later. + _ = w.store.Delete(w.ctx, w.stagingKey) + return fmt.Errorf("blob s3: move staging -> %s: %w", w.finalKey, err) + } + w.finalized = true + return nil +} + +func (w *s3Writer) Abort() error { + if w.finalized { + return nil + } + w.finalized = true + return w.mw.Abort() +} + +func (w *s3Writer) Close() error { return w.Abort() } + +// isS3NotFound reports whether err looks like an S3 "the object you +// asked about doesn't exist" condition. The aws-sdk-go-v2 returns +// these as typed errors that vary across the API surface (NoSuchKey +// for GetObject, NotFound for HeadObject); string matching is the +// least painful way to cover both without dropping into reflect. +func isS3NotFound(err error) bool { + var nsk *types.NoSuchKey + if errors.As(err, &nsk) { + return true + } + msg := err.Error() + return strings.Contains(msg, "NoSuchKey") || strings.Contains(msg, "NotFound") || strings.Contains(msg, "status code: 404") +} + +var ( + _ Store = (*S3Store)(nil) + _ Writer = (*s3Writer)(nil) +) diff --git a/pkg/blob/store.go b/pkg/blob/store.go new file mode 100644 index 00000000..0ca48196 --- /dev/null +++ b/pkg/blob/store.go @@ -0,0 +1,102 @@ +// Package blob defines a content-agnostic blob storage interface and +// supplies file-system and S3 implementations. +// +// The interface is deliberately small — Open / NewWriter / Move / +// Delete — so it can serve as a building block for higher-level +// stores: an S3 archive backed by a local disk cache backed by an +// in-memory cache, for instance. Implementations all assume random +// access on reads (every Reader is an io.ReaderAt with a known Size) +// because the upstream consumers (gstreamer demuxers, range requests +// from playback clients) need it. +// +// Keys are forward-slash-separated strings interpreted relative to the +// store's root (a directory for FileStore, a bucket for S3Store). +// Implementations must accept arbitrary nested paths; the FileStore +// auto-creates parent directories at write time. +package blob + +import ( + "context" + "errors" + "io" +) + +// Store is a content-agnostic random-access blob storage backend. +// +// Implementations are safe for concurrent use from multiple goroutines. +// Reads via Open are completely independent of each other and of any +// in-flight writes to the same key. Behavior when a write to a key +// races a delete or a move of the same key is implementation-defined +// — callers should serialize those operations externally. +type Store interface { + // URL returns a human-readable URL for the given key, useful for + // logging and error messages. Format is implementation-specific + // (e.g. "file:///abs/path", "s3://bucket/key"). + URL(key string) string + + // Open returns a Reader for the blob at key. The supplied context + // scopes any underlying network requests (S3) and is checked when + // closing the reader (the reader holds it). The returned Reader's + // Size method returns the byte length of the blob. + Open(ctx context.Context, key string) (Reader, error) + + // NewWriter starts a streaming write to key. Bytes written go to + // implementation-private staging until Complete is called; on + // Complete the blob becomes visible at key atomically. Abort + // (or Close without Complete) discards the staged bytes. + // + // contentType is advisory and may be ignored by implementations + // that don't track it (FileStore). + NewWriter(ctx context.Context, key, contentType string) (Writer, error) + + // Move relocates the blob from srcKey to dstKey atomically (where + // the underlying storage permits — POSIX rename on FileStore, + // CopyObject+DeleteObject on S3Store). If dstKey already exists, + // it is overwritten. Returns nil if srcKey does not exist after a + // successful Move (idempotency for retried renames). + Move(ctx context.Context, srcKey, dstKey string) error + + // Delete removes the blob at key. Returns nil if the blob does + // not exist (the desired post-condition is the same either way). + Delete(ctx context.Context, key string) error + + // ParseLocation translates a backend-specific URL or path (as + // stored in legacy fields like the upload row's Location column) + // into a Store-relative key. Returns ok=false if the location + // doesn't belong to this Store (wrong scheme, wrong bucket, path + // outside the configured root). + ParseLocation(location string) (key string, ok bool) +} + +// Reader is the random-access read half of a blob. Implementations +// must support concurrent ReadAt calls (file ReaderAt is safe; the +// S3-backed implementation serializes internally). +type Reader interface { + io.ReaderAt + io.Closer + // Size returns the byte length of the blob, known at Open time. + Size() int64 +} + +// Writer is the streaming-write half of a blob. Writes go to staging +// until Complete (atomic publish to the configured key) or Abort +// (discard). Close runs Abort if Complete hasn't been called, so a +// `defer w.Close()` after NewWriter is the recommended error-path +// guard. +type Writer interface { + io.Writer + // Complete publishes the staged bytes to the configured key. After + // a successful Complete the blob exists at key; the writer cannot + // be reused. + Complete() error + // Abort discards staged bytes. Safe to call multiple times and + // safe to call after Complete (in which case it's a no-op). + Abort() error + // Close calls Abort if Complete has not been called. + io.Closer +} + +// ErrNotFound is returned by Open when the requested key doesn't +// exist. Implementations should wrap their backend's not-found error +// (os.ErrNotExist, s3 NoSuchKey) so callers can test with errors.Is. +var ErrNotFound = errors.New("blob: not found") diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 670b21bd..eae900f2 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -42,8 +42,13 @@ import ( "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/storage" "stream.place/streamplace/pkg/upload" + "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/vod" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/credentials" + awss3 "github.com/aws/aws-sdk-go-v2/service/s3" + _ "github.com/go-gst/go-glib/glib" _ "github.com/go-gst/go-gst/gst" "stream.place/streamplace/pkg/api" @@ -349,8 +354,12 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu if err != nil { return err } + vodStore, err := makeVODStore(ctx, cli) + if err != nil { + return fmt.Errorf("make vod store: %w", err) + } state.SetVODProcessor(func(ctx context.Context, t statedb.VODProcessTask) (string, error) { - return vod.ProcessVOD(ctx, cli, vod.Input{ + return vod.ProcessVOD(ctx, vodStore, vod.Input{ UploadID: t.UploadID, RepoDID: t.RepoDID, MimeType: t.MimeType, @@ -556,6 +565,41 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu return group.Wait() } +// makeVODStore picks the blob.Store backing VOD output for this +// process. S3 if it's configured (production / multi-node); otherwise +// the local DataDir. Either way the same Store is used to read the +// user upload AND write the content-addressed VOD output — uploads +// land under "uploads/" and VOD output lands under "vod/" within the +// Store's namespace. +// +// FileStore is rooted at DataDir so it can see the upload manager's +// "uploads/" tree alongside its own "vod/" tree; S3Store is rooted +// at the configured bucket for the same reason. +// +// Mirrors upload.New's multi-node-requires-S3 invariant: in single-node +// file mode the upload and the produced VOD share local disk; in +// multi-node S3 mode any station can pick up a queued VOD task and +// land output in shared storage. +func makeVODStore(ctx context.Context, cli *config.CLI) (blob.Store, error) { + if cli.S3Configured() { + s3client := awss3.New(awss3.Options{ + Region: cli.S3Region, + Credentials: credentials.NewStaticCredentialsProvider( + cli.S3AccessKeyID, + cli.S3SecretAccessKey, + "", + ), + BaseEndpoint: aws.String(cli.S3Endpoint), + UsePathStyle: true, + }) + log.Log(ctx, "VOD store: S3", "bucket", cli.S3Bucket) + return blob.NewS3Store(s3client, cli.S3Bucket), nil + } + root := cli.DataFilePath(nil) + log.Log(ctx, "VOD store: file", "root", root) + return blob.NewFileStore(root) +} + var ErrCaughtSignal = errors.New("caught signal") func handleSignals(ctx context.Context) error { diff --git a/pkg/vod/process.go b/pkg/vod/process.go index 862ea920..011bbe7d 100644 --- a/pkg/vod/process.go +++ b/pkg/vod/process.go @@ -1,15 +1,16 @@ // Package vod implements the post-upload processing pipeline for a VOD // upload: read the user's raw file from wherever the upload manager // stored it, run it through gstreamer (parsebin -> mp4mux -> muxl -// concatenator), and write the resulting fMP4 to S3 under a key derived -// from the BLAKE3-based BDASL CID of the final bytes. +// concatenator), and write the resulting fMP4 to a content-addressed +// key under a blob.Store. // -// The full pipeline is streaming end-to-end — bytes flow from the source +// The pipeline is streaming end-to-end — bytes flow from the source // (file or ranged S3 GETs), through gstreamer's appsrc -> parsebin -> // fdkaacenc/h264parse -> mp4mux -> appsink, into the muxl wasm -// concatenator, and out to an S3 multipart upload. The bdasl hasher -// runs as a tee on the way out; the final CID is known only when the -// last byte is written. +// concatenator, and out to a blob.Store writer. A bdasl.Writer runs +// as a tee on the way out; the final CID is known only when the last +// byte is written, at which point we Move the staging blob to the +// content-addressed key. package vod import ( @@ -17,23 +18,18 @@ import ( "errors" "fmt" "io" - "os" "time" - "github.com/aws/aws-sdk-go-v2/aws" - "github.com/aws/aws-sdk-go-v2/credentials" - awss3 "github.com/aws/aws-sdk-go-v2/service/s3" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/trace" "stream.place/streamplace/pkg/bdasl" - "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/muxl" - s3pkg "stream.place/streamplace/pkg/s3" "stream.place/streamplace/pkg/spmetrics" ) @@ -45,8 +41,8 @@ const ( stageOpenSource = "open_source" stageStaging = "start_staging" stagePipeline = "gstreamer_pipeline" - stageStagingComplete = "s3_complete" - stageContentAddressCopy = "content_address_copy" + stageStagingComplete = "store_complete" + stageContentAddressCopy = "content_address_move" ) // Input is the per-upload state needed to drive ProcessVOD. It's a @@ -59,33 +55,34 @@ type Input struct { MimeType string Filename string Size int64 - // Backend is the storage tier the user upload lives on: "file" or "s3". + // Backend describes the upload source backend ("file" or "s3"), + // matching what the upload manager stored on the row. Used only + // for metrics labeling; the actual reading goes through the + // blob.Store passed to ProcessVOD. Backend string - // Location is the path (file backend) or s3:// URL (s3 backend) of - // the user upload. + // Location is the backend-specific URL/path (a file path or + // "s3://bucket/key" URL) the upload landed at. Translated to a + // Store key via Store.ParseLocation. Location string } -// Backend constants mirror pkg/upload.BackendFile / BackendS3 without -// importing the upload package (which would pull statedb in transitively). -const ( - BackendFile = "file" - BackendS3 = "s3" -) - -// StagingPrefix is where in-progress VOD outputs land in S3 before being -// renamed to their content-addressed key. We park them under a dedicated -// prefix so a periodic janitor can sweep abandoned uploads. +// StagingPrefix is where in-progress VOD outputs land before being +// renamed to their content-addressed key. We park them under a +// dedicated prefix so a periodic janitor can sweep abandoned uploads. const StagingPrefix = "vod-staging/" // ContentPrefix is the prefix for the final, content-addressed object. const ContentPrefix = "vod/" // ProcessVOD runs the streaming pipeline for one VOD upload and returns -// the BDASL CID of the resulting fMP4. The output object is written at -// ContentPrefix+.fmp4 in the configured S3 bucket; staging objects -// are cleaned up on success or failure. -func ProcessVOD(ctx context.Context, cli *config.CLI, in Input) (string, error) { +// the BDASL CID of the resulting fMP4. Reads come from `in` via the +// supplied Store; the output blob lands at ContentPrefix+.fmp4 in +// the same Store; staging blobs are cleaned up on success or failure. +// +// The Store is the storage layer; it can be either a FileStore (single- +// node deployments) or an S3Store (production). Future Stores can mix +// caches and archives behind the same interface. +func ProcessVOD(ctx context.Context, store blob.Store, in Input) (string, error) { ctx = log.WithLogValues(ctx, "func", "ProcessVOD", "uploadId", in.UploadID, "did", in.RepoDID) ctx, span := vodTracer.Start(ctx, "vod.ProcessVOD", trace.WithAttributes( attribute.String("upload_id", in.UploadID), @@ -102,16 +99,14 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, in Input) (string, error) spmetrics.VODProcessDurationMS.Observe(float64(time.Since(startTime).Milliseconds())) }() - log.Log(ctx, "starting VOD processing", "backend", in.Backend, "size", in.Size, "mimeType", in.MimeType) - - if !cli.S3Configured() { - err := errors.New("vod processing requires S3 to be configured") - recordErr(span, "config", err) - return "", err - } - s3client := newS3Client(cli) + log.Log(ctx, "starting VOD processing", + "backend", in.Backend, + "size", in.Size, + "mimeType", in.MimeType, + "location", in.Location, + ) - src, size, closer, err := openSource(ctx, cli, s3client, in) + src, size, closer, err := openSource(ctx, store, in) if err != nil { recordErr(span, stageOpenSource, err) return "", fmt.Errorf("open source: %w", err) @@ -122,7 +117,7 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, in Input) (string, error) stagingKey := StagingPrefix + in.UploadID + ".fmp4" span.SetAttributes(attribute.String("staging_key", stagingKey)) - staging, err := s3pkg.NewMultipartWriter(ctx, s3client, cli.S3Bucket, stagingKey, "video/mp4") + staging, err := store.NewWriter(ctx, stagingKey, "video/mp4") if err != nil { recordErr(span, stageStaging, err) return "", fmt.Errorf("start staging upload: %w", err) @@ -151,9 +146,10 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, in Input) (string, error) span.SetAttributes( attribute.String("cid", finalCID), attribute.String("content_key", contentKey), + attribute.String("content_url", store.URL(contentKey)), ) - if err := finalizeUpload(ctx, s3client, cli.S3Bucket, stagingKey, contentKey); err != nil { + if err := finalizeMove(ctx, store, stagingKey, contentKey); err != nil { recordErr(span, stageContentAddressCopy, err) return "", fmt.Errorf("finalize: %w", err) } @@ -162,7 +158,7 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, in Input) (string, error) span.SetStatus(codes.Ok, "") log.Log(ctx, "VOD processed", "cid", finalCID, - "key", contentKey, + "url", store.URL(contentKey), "input_size", size, "output_size", counter.n, "duration_ms", time.Since(startTime).Milliseconds(), @@ -170,10 +166,10 @@ func ProcessVOD(ctx context.Context, cli *config.CLI, in Input) (string, error) return finalCID, nil } -// completeStaging wraps the multipart Complete call in its own span so -// we can see how much of total processing time is spent waiting for S3 -// to acknowledge the upload. -func completeStaging(ctx context.Context, staging *s3pkg.MultipartWriter, stagingKey string) error { +// completeStaging wraps the writer Complete call in its own span so +// we can see how much of total processing time is spent waiting for +// the storage layer to acknowledge the upload. +func completeStaging(ctx context.Context, staging blob.Writer, stagingKey string) error { _, span := vodTracer.Start(ctx, "vod.completeStaging", trace.WithAttributes( attribute.String("staging_key", stagingKey), )) @@ -186,6 +182,28 @@ func completeStaging(ctx context.Context, staging *s3pkg.MultipartWriter, stagin return nil } +// finalizeMove atomically promotes the staged blob at stagingKey to the +// content-addressed contentKey. For S3Store this is a CopyObject + +// DeleteObject; for FileStore it's an os.Rename. +func finalizeMove(ctx context.Context, store blob.Store, stagingKey, contentKey string) error { + ctx, span := vodTracer.Start(ctx, "vod.finalizeMove", trace.WithAttributes( + attribute.String("staging_key", stagingKey), + attribute.String("content_key", contentKey), + )) + defer span.End() + moveStart := time.Now() + if err := store.Move(ctx, stagingKey, contentKey); err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, "move") + return err + } + span.SetAttributes(attribute.Int64("move_duration_ms", time.Since(moveStart).Milliseconds())) + log.Debug(ctx, "moved staging to content-addressed key", + "duration_ms", time.Since(moveStart).Milliseconds(), + ) + return nil +} + // countingWriter is an io.Writer that tallies bytes written. Used to // observe output size without buffering or hashing it twice. type countingWriter struct{ n int64 } @@ -346,101 +364,35 @@ func consumeConcatTraced(ctx context.Context, c *muxl.Concatenator, dst io.Write return nil } -// openSource returns a ReaderAt + size for the upload, regardless of -// backend. The returned closer must be called when the caller is done +// openSource returns a Reader + size for the upload via the supplied +// Store. The returned closer must be called when the caller is done // with the source. -func openSource(ctx context.Context, cli *config.CLI, s3client *awss3.Client, in Input) (io.ReaderAt, int64, func(), error) { +// +// The upload's Location field (a backend-specific URL/path) is +// translated to a Store-relative key via store.ParseLocation. If the +// configured Store can't claim the Location — e.g. the upload was +// stored in S3 but the deployment is now serving file-only — the open +// fails with a clear error. +func openSource(ctx context.Context, store blob.Store, in Input) (io.ReaderAt, int64, func(), error) { ctx, span := vodTracer.Start(ctx, "vod.openSource", trace.WithAttributes( attribute.String("backend", in.Backend), attribute.String("location", in.Location), )) defer span.End() - switch in.Backend { - case BackendFile: - f, err := os.Open(in.Location) - if err != nil { - span.RecordError(err) - return nil, 0, nil, fmt.Errorf("open file upload %q: %w", in.Location, err) - } - st, err := f.Stat() - if err != nil { - span.RecordError(err) - _ = f.Close() - return nil, 0, nil, fmt.Errorf("stat file upload %q: %w", in.Location, err) - } - span.SetAttributes(attribute.Int64("size_bytes", st.Size())) - log.Debug(ctx, "opened file upload", "path", in.Location, "size", st.Size()) - return f, st.Size(), func() { _ = f.Close() }, nil - case BackendS3: - bucket, key, err := s3pkg.ParseURL(in.Location) - if err != nil { - span.RecordError(err) - return nil, 0, nil, fmt.Errorf("parse s3 location %q: %w", in.Location, err) - } - span.SetAttributes(attribute.String("bucket", bucket), attribute.String("key", key)) - ra, err := s3pkg.NewReaderAt(ctx, s3client, bucket, key) - if err != nil { - span.RecordError(err) - return nil, 0, nil, fmt.Errorf("open s3 upload: %w", err) - } - span.SetAttributes(attribute.Int64("size_bytes", ra.Size())) - log.Debug(ctx, "opened s3 upload", "bucket", bucket, "key", key, "size", ra.Size()) - return ra, ra.Size(), func() { _ = ra.Close() }, nil - default: - err := fmt.Errorf("unknown upload backend %q", in.Backend) + key, ok := store.ParseLocation(in.Location) + if !ok { + err := fmt.Errorf("store does not own upload location %q (backend %s)", in.Location, in.Backend) span.RecordError(err) return nil, 0, nil, err } -} - -// finalizeUpload renames the staging object to its content-addressed -// key. S3 has no rename: we CopyObject server-side, then DeleteObject -// the staging key. If the content key already exists (duplicate upload -// of identical content), we still proceed — copy is idempotent and we -// still want to drop the staging copy. -func finalizeUpload(ctx context.Context, c *awss3.Client, bucket, stagingKey, contentKey string) error { - ctx, span := vodTracer.Start(ctx, "vod.finalizeUpload", trace.WithAttributes( - attribute.String("bucket", bucket), - attribute.String("staging_key", stagingKey), - attribute.String("content_key", contentKey), - )) - defer span.End() - - copyStart := time.Now() - _, err := c.CopyObject(ctx, &awss3.CopyObjectInput{ - Bucket: aws.String(bucket), - Key: aws.String(contentKey), - CopySource: aws.String(bucket + "/" + stagingKey), - }) + span.SetAttributes(attribute.String("key", key)) + rdr, err := store.Open(ctx, key) if err != nil { span.RecordError(err) - span.SetStatus(codes.Error, "copy") - return fmt.Errorf("copy staging -> %s: %w", contentKey, err) + return nil, 0, nil, fmt.Errorf("open upload: %w", err) } - span.SetAttributes(attribute.Int64("copy_duration_ms", time.Since(copyStart).Milliseconds())) - log.Debug(ctx, "copied staging to content-addressed key", "duration_ms", time.Since(copyStart).Milliseconds()) - - if _, err := c.DeleteObject(ctx, &awss3.DeleteObjectInput{ - Bucket: aws.String(bucket), - Key: aws.String(stagingKey), - }); err != nil { - // Non-fatal: the content key is in place; staging will be swept - // later. Log + continue, and record it on the span for visibility. - span.RecordError(err) - log.Warn(ctx, "failed to delete staging object", "key", stagingKey, "error", err) - } - return nil -} - -func newS3Client(cli *config.CLI) *awss3.Client { - return awss3.New(awss3.Options{ - Region: cli.S3Region, - Credentials: credentials.NewStaticCredentialsProvider( - cli.S3AccessKeyID, - cli.S3SecretAccessKey, - "", - ), - BaseEndpoint: aws.String(cli.S3Endpoint), - UsePathStyle: true, - }) + size := rdr.Size() + span.SetAttributes(attribute.Int64("size_bytes", size)) + log.Debug(ctx, "opened upload via Store", "url", store.URL(key), "size", size) + return rdr, size, func() { _ = rdr.Close() }, nil }