Skip to content

Commit ed9a4ba

Browse files
committed
fix(storage): bound reconciliation and cleanup memory usage
1 parent 41db781 commit ed9a4ba

18 files changed

Lines changed: 1457 additions & 618 deletions

internal/cache/service.go

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ const (
4343
)
4444

4545
// MaxBlockListEntries is the Azure block blob protocol limit.
46-
const MaxBlockListEntries = 50_000
46+
const MaxBlockListEntries = storage.MaxIndexedObjects
4747

4848
var (
4949
ErrNoWriteScope = errors.New("no scope with write permission found")
@@ -452,13 +452,17 @@ func (s *Service) CompleteUpload(ctx context.Context, key, version string, scope
452452
currentUpload.FinishedPartUploadCount,
453453
)
454454
}
455-
parts, inspectErr := s.storage.InspectFolder(activityCtx, partsFolderName(currentUpload.FolderName))
456-
if inspectErr != nil {
457-
return fmt.Errorf("inspect finalized cache parts: %w", inspectErr)
458-
}
459-
sizeBytes, validationErr = parts.LogicalIndexedSize(partCount)
455+
sizeBytes, validationErr = s.storage.InspectIndexedFolder(
456+
activityCtx,
457+
partsFolderName(currentUpload.FolderName),
458+
partCount,
459+
)
460460
if validationErr != nil {
461-
return fmt.Errorf("%w: %v", ErrPartCountMismatch, validationErr)
461+
if errors.Is(validationErr, storage.ErrIndexedObjectMissing) ||
462+
errors.Is(validationErr, storage.ErrIndexedObjectLimitExceeded) {
463+
return fmt.Errorf("%w: %v", ErrPartCountMismatch, validationErr)
464+
}
465+
return fmt.Errorf("inspect finalized cache parts: %w", validationErr)
462466
}
463467
return nil
464468
})

internal/cache/service_test.go

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -222,6 +222,45 @@ func TestCompleteUploadKeepsLegacyUploadAfterTransientStorageError(t *testing.T)
222222
require.NoError(t, err)
223223
}
224224

225+
func TestCompleteUploadKeepsUploadAfterTransientIndexedInspectionError(t *testing.T) {
226+
ctx, client, filesystem := newTestServiceDeps(t)
227+
adapter := &storageCallTrackingAdapter{Adapter: filesystem, inspectIndexedErr: errInjectedStorageFailure}
228+
service := NewService(Options{DB: client, Storage: adapter})
229+
scope := writableScope()
230+
231+
upload, err := service.CreateUpload(ctx, "key", "version", scope)
232+
require.NoError(t, err)
233+
require.NoError(t, service.UploadPart(ctx, upload.UploadID, bytes.NewBufferString("data")))
234+
235+
_, err = service.CompleteUpload(ctx, "key", "version", scope)
236+
require.ErrorIs(t, err, errInjectedStorageFailure)
237+
_, err = client.Upload.Get(ctx, upload.UploadID)
238+
require.NoError(t, err)
239+
}
240+
241+
func TestCompleteUploadRejectsAndDeletesOverLimitRestoredSession(t *testing.T) {
242+
ctx, client, filesystem := newTestServiceDeps(t)
243+
service := NewService(Options{DB: client, Storage: filesystem})
244+
scope := writableScope()
245+
const uploadID = int64(42)
246+
client.Upload.Create().
247+
SetID(uploadID).
248+
SetKey("restored-key").
249+
SetVersion("version").
250+
SetScope(scope.Scopes[0].Scope).
251+
SetRepoId(scope.RepoID).
252+
SetCreatedAt(time.Now().UnixMilli()).
253+
SetFinishedPartUploadCount(storage.MaxIndexedObjects + 1).
254+
SetCommittedPartCount(storage.MaxIndexedObjects + 1).
255+
SetFolderName("restored-upload").
256+
SaveX(ctx)
257+
258+
_, err := service.CompleteUpload(ctx, "restored-key", "version", scope)
259+
require.ErrorIs(t, err, ErrPartCountMismatch)
260+
_, err = client.Upload.Get(ctx, uploadID)
261+
require.True(t, ent.IsNotFound(err))
262+
}
263+
225264
func TestCompleteUploadRejectsBlocksWithoutCommittedBlockList(t *testing.T) {
226265
ctx, client, filesystem := newTestServiceDeps(t)
227266
service := NewService(Options{DB: client, Storage: filesystem})
@@ -1287,6 +1326,7 @@ type storageCallTrackingAdapter struct {
12871326
downloadCalls int
12881327
objectExistsCalls []string
12891328
countErr error
1329+
inspectIndexedErr error
12901330
}
12911331

12921332
func (s *storageCallTrackingAdapter) CountFilesInFolder(ctx context.Context, folderName string) (int, error) {
@@ -1302,6 +1342,13 @@ func (s *storageCallTrackingAdapter) CreateDownloadStream(ctx context.Context, o
13021342
return s.Adapter.CreateDownloadStream(ctx, objectName)
13031343
}
13041344

1345+
func (s *storageCallTrackingAdapter) InspectIndexedFolder(ctx context.Context, folderName string, expectedObjects int) (int64, error) {
1346+
if s.inspectIndexedErr != nil {
1347+
return 0, s.inspectIndexedErr
1348+
}
1349+
return s.Adapter.InspectIndexedFolder(ctx, folderName, expectedObjects)
1350+
}
1351+
13051352
func (s *storageCallTrackingAdapter) ObjectExists(ctx context.Context, objectName string) (bool, error) {
13061353
s.objectExistsCalls = append(s.objectExistsCalls, objectName)
13071354
return s.Adapter.ObjectExists(ctx, objectName)

internal/cache/upload_activity_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -331,10 +331,10 @@ type afterInspectStorage struct {
331331
once sync.Once
332332
}
333333

334-
func (s *afterInspectStorage) InspectFolder(ctx context.Context, folderName string) (storage.FolderContents, error) {
335-
contents, err := s.Adapter.InspectFolder(ctx, folderName)
334+
func (s *afterInspectStorage) InspectIndexedFolder(ctx context.Context, folderName string, expectedObjects int) (int64, error) {
335+
sizeBytes, err := s.Adapter.InspectIndexedFolder(ctx, folderName, expectedObjects)
336336
if err == nil {
337337
s.once.Do(s.after)
338338
}
339-
return contents, err
339+
return sizeBytes, err
340340
}

0 commit comments

Comments
 (0)