Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions internal/apiserver/file/file_handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,10 @@ func (c *errStoreFileClient) Retrieve(ctx context.Context, fileName, folderName
return nil, nil, errors.New("not implemented")
}

func (c *errStoreFileClient) RetrieveRange(_ context.Context, _, _ string, _ int64, _ int64) (io.ReadCloser, error) {
return nil, errors.New("not implemented")
}

func (c *errStoreFileClient) Delete(ctx context.Context, fileName, folderName string) error {
return nil
}
Expand Down
3 changes: 3 additions & 0 deletions internal/files_store/api/files.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,9 @@ type BatchFilesClient interface {
// Retrieve retrieves a file from the files storage.
Retrieve(ctx context.Context, fileName, folderName string) (reader io.ReadCloser, fileMd *BatchFileMetadata, err error)

// RetrieveRange retrieves a byte range [offset, offset+length) from a file in storage.
RetrieveRange(ctx context.Context, fileName, folderName string, offset int64, length int64) (io.ReadCloser, error)

// Delete deletes the file in the specified location.
Delete(ctx context.Context, fileName, folderName string) (err error)
}
27 changes: 27 additions & 0 deletions internal/files_store/fs/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ limitations under the License.
package fs

import (
"bytes"
"context"
"errors"
"fmt"
Expand Down Expand Up @@ -213,6 +214,32 @@ func (c *Client) Retrieve(ctx context.Context, fileName, folderName string) (io.
return file, metadata, nil
}

// RetrieveRange retrieves a byte range [offset, offset+length) from a file in the filesystem.
func (c *Client) RetrieveRange(ctx context.Context, fileName, folderName string, offset int64, length int64) (io.ReadCloser, error) {
relPath, err := c.resolvePath(folderName, fileName)
if err != nil {
return nil, err
}

file, err := c.root.Open(relPath)
if err != nil {
return nil, err
}
defer file.Close()

buf := make([]byte, length)
Comment thread
acardace marked this conversation as resolved.
n, err := file.ReadAt(buf, offset)
if err != nil && (err != io.EOF || int64(n) != length) {
return nil, fmt.Errorf("read range at offset %d length %d: %w", offset, length, err)
}
buf = buf[:n]

logr.FromContextOrDiscard(ctx).V(logging.TRACE).Info("Range retrieved",
"path", relPath, "offset", offset, "length", length)

return io.NopCloser(bytes.NewReader(buf)), nil
}

// Delete deletes a file from the filesystem.
// After removing the file, it attempts to remove the parent directory if empty.
// The directory cleanup is best-effort and does not fail the operation. (e.g. if the directory is not empty, it will not be removed)
Expand Down
35 changes: 35 additions & 0 deletions internal/files_store/fs/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,41 @@ func TestRetrieve(t *testing.T) {

}

func TestRetrieveRange(t *testing.T) {
ctx := context.Background()

t.Run("retrieves byte range", func(t *testing.T) {
client := newTestClient(t)
content := []byte("hello world range test")

if _, err := client.Store(ctx, "range.txt", testFolder, 1024, 0, bytes.NewReader(content)); err != nil {
t.Fatalf("failed to store: %v", err)
}

rc, err := client.RetrieveRange(ctx, "range.txt", testFolder, 6, 5)
if err != nil {
t.Fatalf("expected no error, got %v", err)
}
defer rc.Close()

data, _ := io.ReadAll(rc)
if string(data) != "world" {
t.Errorf("expected %q, got %q", "world", string(data))
}
})

t.Run("returns error for non-existent file", func(t *testing.T) {
client := newTestClient(t)

_, err := client.RetrieveRange(ctx, "nonexistent.txt", testFolder, 0, 10)
if err != nil {
// Expect an error (file not found).
return
}
t.Error("expected error for non-existent file")
})
}

func TestDelete(t *testing.T) {
ctx := context.Background()

Expand Down
21 changes: 21 additions & 0 deletions internal/files_store/mock/mock_files_client.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package mock

import (
"bufio"
"bytes"
"context"
"fmt"
"io"
Expand Down Expand Up @@ -143,6 +144,26 @@ func (m *MockBatchFilesClient) Retrieve(ctx context.Context, fileName, folderNam
}, nil
}

// RetrieveRange retrieves a byte range [offset, offset+length) from a file.
func (m *MockBatchFilesClient) RetrieveRange(_ context.Context, fileName, folderName string, offset int64, length int64) (io.ReadCloser, error) {
filePath := filepath.Join(m.rootDir, folderName, fileName)

file, err := os.Open(filePath)
if err != nil {
return nil, fmt.Errorf("failed to open file: %w", err)
}
defer file.Close()

buf := make([]byte, length)
n, err := file.ReadAt(buf, offset)
if err != nil && (err != io.EOF || int64(n) != length) {
return nil, fmt.Errorf("failed to read range at offset %d length %d: %w", offset, length, err)
}
buf = buf[:n]

return io.NopCloser(bytes.NewReader(buf)), nil
}

// List lists the files in the specified location.
func (m *MockBatchFilesClient) List(ctx context.Context, location string) ([]api.BatchFileMetadata, error) {
// Use /tmp as root folder
Expand Down
27 changes: 27 additions & 0 deletions internal/files_store/retryclient/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,33 @@ func (c *Client) Retrieve(ctx context.Context, fileName, folderName string) (io.
return rc, meta, err
}

func (c *Client) RetrieveRange(ctx context.Context, fileName, folderName string, offset int64, length int64) (io.ReadCloser, error) {
var rc io.ReadCloser
attempts, err := retry.Do(ctx, &c.cfg, func(attempt int) error {
Comment thread
acardace marked this conversation as resolved.
if attempt > 1 {
if rc != nil {
_ = rc.Close()
}
recordRetry("retrieve_range", c.component)
logr.FromContextOrDiscard(ctx).Info("Retrying file retrieve range",
"file", fileName, "offset", offset, "length", length,
"attempt", attempt, "maxRetries", c.cfg.MaxRetries)
}

var retrieveErr error
rc, retrieveErr = c.inner.RetrieveRange(ctx, fileName, folderName, offset, length)
return retrieveErr
})
if err != nil {
if attempts > c.cfg.MaxRetries {
recordExhausted("retrieve_range", c.component)
}
} else {
recordSuccess("retrieve_range", c.component)
}
return rc, err
}

func (c *Client) Delete(ctx context.Context, fileName, folderName string) error {
attempts, err := retry.Do(ctx, &c.cfg, func(attempt int) error {
if attempt > 1 {
Expand Down
26 changes: 26 additions & 0 deletions internal/files_store/retryclient/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,14 @@ func (m *mockFilesClient) Delete(_ context.Context, _, _ string) error {
return nil
}

func (m *mockFilesClient) RetrieveRange(_ context.Context, _, _ string, _ int64, _ int64) (io.ReadCloser, error) {
m.retrieveCalls++
if m.retrieveCalls <= m.failUntil {
return nil, errors.New("transient retrieve range error")
}
return io.NopCloser(bytes.NewReader([]byte("data"))), nil
}

func (m *mockFilesClient) Close() error { return nil }

func retryCfg() retry.Config {
Expand Down Expand Up @@ -196,6 +204,24 @@ func TestRetrieve_SucceedsAfterRetry(t *testing.T) {
}
}

func TestRetrieveRange_SucceedsAfterRetry(t *testing.T) {
mock := &mockFilesClient{failUntil: 1}
c := New(mock, retryCfg(), "test")

rc, err := c.RetrieveRange(context.Background(), "f.txt", "folder", 0, 4)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
defer rc.Close()
data, _ := io.ReadAll(rc)
if string(data) != "data" {
t.Fatalf("expected %q, got %q", "data", string(data))
}
if mock.retrieveCalls != 2 {
t.Fatalf("expected 2 retrieve calls, got %d", mock.retrieveCalls)
}
}

func TestDelete_SucceedsAfterRetry(t *testing.T) {
mock := &mockFilesClient{failUntil: 2}
c := New(mock, retryCfg(), "test")
Expand Down
25 changes: 25 additions & 0 deletions internal/files_store/s3/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -274,6 +274,31 @@ func (c *Client) Retrieve(ctx context.Context, fileName, folderName string) (io.
return out.Body, metadata, nil
}

// RetrieveRange retrieves a byte range [offset, offset+length) from a file in S3.
func (c *Client) RetrieveRange(ctx context.Context, fileName, folderName string, offset int64, length int64) (io.ReadCloser, error) {
key := c.buildKey(folderName, fileName)
rangeHeader := fmt.Sprintf("bytes=%d-%d", offset, offset+length-1)

out, err := c.s3Client.GetObject(ctx, &s3.GetObjectInput{
Bucket: aws.String(c.bucket),
Key: aws.String(key),
Range: aws.String(rangeHeader),
})
if err != nil {
var noSuchKey *types.NoSuchKey
var noSuchBucket *types.NoSuchBucket
if errors.As(err, &noSuchKey) || errors.As(err, &noSuchBucket) {
return nil, os.ErrNotExist
}
return nil, err
}

logr.FromContextOrDiscard(ctx).V(logging.TRACE).Info("Range retrieved",
"bucket", c.bucket, "key", key, "offset", offset, "length", length)

return out.Body, nil
}

// Delete deletes a file from S3.
// The folderName parameter is used as a key prefix for tenant isolation.
func (c *Client) Delete(ctx context.Context, fileName, folderName string) error {
Expand Down
51 changes: 49 additions & 2 deletions internal/files_store/s3/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"bytes"
"context"
"errors"
"fmt"
"io"
"os"
"testing"
Expand Down Expand Up @@ -98,9 +99,20 @@ func (m *mockS3Client) GetObject(_ context.Context, params *s3.GetObjectInput, _
if !ok {
return nil, &types.NoSuchKey{}
}
data := obj.data
if params.Range != nil {
// Parse "bytes=start-end" range header.
var start, end int64
if _, err := fmt.Sscanf(*params.Range, "bytes=%d-%d", &start, &end); err == nil {
if end >= int64(len(data)) {
end = int64(len(data)) - 1
}
data = data[start : end+1]
}
}
return &s3.GetObjectOutput{
Body: io.NopCloser(bytes.NewReader(obj.data)),
ContentLength: aws.Int64(int64(len(obj.data))),
Body: io.NopCloser(bytes.NewReader(data)),
ContentLength: aws.Int64(int64(len(data))),
LastModified: aws.Time(obj.lastModTime),
}, nil
}
Expand Down Expand Up @@ -337,6 +349,41 @@ func TestRetrieve(t *testing.T) {
})
}

func TestRetrieveRange(t *testing.T) {
ctx := context.Background()

t.Run("retrieves byte range", func(t *testing.T) {
mock := newMockS3Client()
client := newTestClient(mock)
content := []byte("hello world range test")

if _, err := client.Store(ctx, "range.txt", testBucketName, 1024, 0, bytes.NewReader(content)); err != nil {
t.Fatalf("failed to store: %v", err)
}

rc, err := client.RetrieveRange(ctx, "range.txt", testBucketName, 6, 5)
if err != nil {
t.Fatalf("expected no error, got %v", err)
}
defer rc.Close()

data, _ := io.ReadAll(rc)
if string(data) != "world" {
t.Errorf("expected %q, got %q", "world", string(data))
}
})

t.Run("returns error for non-existent file", func(t *testing.T) {
mock := newMockS3Client()
client := newTestClient(mock)

_, err := client.RetrieveRange(ctx, "nonexistent.txt", testBucketName, 0, 10)
if !errors.Is(err, os.ErrNotExist) {
t.Errorf("expected os.ErrNotExist, got %v", err)
}
})
}

func TestDelete(t *testing.T) {
ctx := context.Background()

Expand Down
20 changes: 20 additions & 0 deletions internal/files_store/tracing/tracing.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,26 @@ func (r *tracedReadCloser) Close() error {
return err
}

func (c *Client) RetrieveRange(ctx context.Context, fileName, folderName string, offset int64, length int64) (io.ReadCloser, error) {
_, span := uotel.StartSpan(ctx, "storage.RetrieveRange")
span.SetAttributes(
c.backend,
attribute.String("storage.file_name", fileName),
attribute.String("storage.folder", folderName),
attribute.Int64("storage.offset", offset),
attribute.Int64("storage.length", length),
)

reader, err := c.inner.RetrieveRange(ctx, fileName, folderName, offset, length)
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, "retrieve range failed")
span.End()
return nil, err
}
return &tracedReadCloser{ReadCloser: reader, span: span}, nil
}

func (c *Client) Delete(ctx context.Context, fileName, folderName string) error {
ctx, span := uotel.StartSpan(ctx, "storage.Delete")
defer span.End()
Expand Down
Loading
Loading