Skip to content

Commit f66944f

Browse files
committed
Implement per-object locks for better parallelism
1 parent 0c333c1 commit f66944f

2 files changed

Lines changed: 346 additions & 17 deletions

File tree

internal/backend/backend_test.go

Lines changed: 266 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -374,6 +374,272 @@ func TestBucketDuplication(t *testing.T) {
374374
})
375375
}
376376

377+
func TestParallelUploadsDifferentObjects(t *testing.T) {
378+
testForStorageBackends(t, func(t *testing.T, storage Storage) {
379+
const bucketName = "parallel-test-bucket"
380+
const numObjects = 10
381+
382+
err := storage.CreateBucket(bucketName, BucketAttrs{})
383+
noError(t, err)
384+
385+
// Use a channel to synchronize goroutines starting together
386+
start := make(chan struct{})
387+
errCh := make(chan error, numObjects)
388+
389+
for i := 0; i < numObjects; i++ {
390+
go func(idx int) {
391+
<-start // Wait for signal to start
392+
objectName := fmt.Sprintf("object-%d", idx)
393+
content := []byte(fmt.Sprintf("content for object %d", idx))
394+
obj := StreamingObject{
395+
ObjectAttrs: ObjectAttrs{
396+
BucketName: bucketName,
397+
Name: objectName,
398+
},
399+
Content: newStreamingContent(content),
400+
}
401+
created, err := storage.CreateObject(obj, NoConditions{})
402+
if err != nil {
403+
errCh <- fmt.Errorf("failed to create object %d: %w", idx, err)
404+
return
405+
}
406+
created.Close()
407+
errCh <- nil
408+
}(i)
409+
}
410+
411+
// Start all goroutines at once
412+
close(start)
413+
414+
// Wait for all to complete
415+
for i := 0; i < numObjects; i++ {
416+
if err := <-errCh; err != nil {
417+
t.Error(err)
418+
}
419+
}
420+
421+
// Verify all objects were created
422+
objects, err := storage.ListObjects(bucketName, "", false)
423+
noError(t, err)
424+
if len(objects) != numObjects {
425+
t.Errorf("expected %d objects, got %d", numObjects, len(objects))
426+
}
427+
})
428+
}
429+
430+
func TestParallelDownloadsSameObject(t *testing.T) {
431+
testForStorageBackends(t, func(t *testing.T, storage Storage) {
432+
const bucketName = "parallel-download-bucket"
433+
const objectName = "shared-object"
434+
const numReaders = 10
435+
content := []byte("shared content for parallel reads")
436+
437+
err := storage.CreateBucket(bucketName, BucketAttrs{})
438+
noError(t, err)
439+
440+
obj := StreamingObject{
441+
ObjectAttrs: ObjectAttrs{
442+
BucketName: bucketName,
443+
Name: objectName,
444+
},
445+
Content: newStreamingContent(content),
446+
}
447+
created, err := storage.CreateObject(obj, NoConditions{})
448+
noError(t, err)
449+
created.Close()
450+
451+
// Use a channel to synchronize goroutines starting together
452+
start := make(chan struct{})
453+
errCh := make(chan error, numReaders)
454+
455+
for i := 0; i < numReaders; i++ {
456+
go func(idx int) {
457+
<-start // Wait for signal to start
458+
retrieved, err := storage.GetObject(bucketName, objectName)
459+
if err != nil {
460+
errCh <- fmt.Errorf("reader %d failed to get object: %w", idx, err)
461+
return
462+
}
463+
defer retrieved.Close()
464+
465+
data, err := io.ReadAll(retrieved.Content)
466+
if err != nil {
467+
errCh <- fmt.Errorf("reader %d failed to read content: %w", idx, err)
468+
return
469+
}
470+
if !bytes.Equal(data, content) {
471+
errCh <- fmt.Errorf("reader %d got wrong content: %q", idx, data)
472+
return
473+
}
474+
errCh <- nil
475+
}(i)
476+
}
477+
478+
// Start all goroutines at once
479+
close(start)
480+
481+
// Wait for all to complete
482+
for i := 0; i < numReaders; i++ {
483+
if err := <-errCh; err != nil {
484+
t.Error(err)
485+
}
486+
}
487+
})
488+
}
489+
490+
func TestParallelUploadAndDownloadDifferentObjects(t *testing.T) {
491+
testForStorageBackends(t, func(t *testing.T, storage Storage) {
492+
const bucketName = "parallel-mixed-bucket"
493+
const existingObject = "existing-object"
494+
const newObject = "new-object"
495+
existingContent := []byte("existing content")
496+
newContent := []byte("new content being uploaded")
497+
498+
err := storage.CreateBucket(bucketName, BucketAttrs{})
499+
noError(t, err)
500+
501+
// Create an existing object to download
502+
obj := StreamingObject{
503+
ObjectAttrs: ObjectAttrs{
504+
BucketName: bucketName,
505+
Name: existingObject,
506+
},
507+
Content: newStreamingContent(existingContent),
508+
}
509+
created, err := storage.CreateObject(obj, NoConditions{})
510+
noError(t, err)
511+
created.Close()
512+
513+
// Use a channel to synchronize goroutines starting together
514+
start := make(chan struct{})
515+
errCh := make(chan error, 2)
516+
517+
// Goroutine 1: Download existing object
518+
go func() {
519+
<-start
520+
retrieved, err := storage.GetObject(bucketName, existingObject)
521+
if err != nil {
522+
errCh <- fmt.Errorf("download failed: %w", err)
523+
return
524+
}
525+
defer retrieved.Close()
526+
527+
data, err := io.ReadAll(retrieved.Content)
528+
if err != nil {
529+
errCh <- fmt.Errorf("read failed: %w", err)
530+
return
531+
}
532+
if !bytes.Equal(data, existingContent) {
533+
errCh <- fmt.Errorf("wrong content: %q", data)
534+
return
535+
}
536+
errCh <- nil
537+
}()
538+
539+
// Goroutine 2: Upload new object
540+
go func() {
541+
<-start
542+
newObj := StreamingObject{
543+
ObjectAttrs: ObjectAttrs{
544+
BucketName: bucketName,
545+
Name: newObject,
546+
},
547+
Content: newStreamingContent(newContent),
548+
}
549+
created, err := storage.CreateObject(newObj, NoConditions{})
550+
if err != nil {
551+
errCh <- fmt.Errorf("upload failed: %w", err)
552+
return
553+
}
554+
created.Close()
555+
errCh <- nil
556+
}()
557+
558+
// Start both goroutines at once
559+
close(start)
560+
561+
// Wait for both to complete
562+
for i := 0; i < 2; i++ {
563+
if err := <-errCh; err != nil {
564+
t.Error(err)
565+
}
566+
}
567+
568+
// Verify both objects exist
569+
objects, err := storage.ListObjects(bucketName, "", false)
570+
noError(t, err)
571+
if len(objects) != 2 {
572+
t.Errorf("expected 2 objects, got %d", len(objects))
573+
}
574+
})
575+
}
576+
577+
func TestParallelUploadsSameObject(t *testing.T) {
578+
testForStorageBackends(t, func(t *testing.T, storage Storage) {
579+
const bucketName = "parallel-same-object-bucket"
580+
const objectName = "contested-object"
581+
const numWriters = 10
582+
583+
err := storage.CreateBucket(bucketName, BucketAttrs{})
584+
noError(t, err)
585+
586+
// Use a channel to synchronize goroutines starting together
587+
start := make(chan struct{})
588+
errCh := make(chan error, numWriters)
589+
590+
for i := 0; i < numWriters; i++ {
591+
go func(idx int) {
592+
<-start // Wait for signal to start
593+
content := []byte(fmt.Sprintf("content from writer %d", idx))
594+
obj := StreamingObject{
595+
ObjectAttrs: ObjectAttrs{
596+
BucketName: bucketName,
597+
Name: objectName,
598+
},
599+
Content: newStreamingContent(content),
600+
}
601+
created, err := storage.CreateObject(obj, NoConditions{})
602+
if err != nil {
603+
errCh <- fmt.Errorf("writer %d failed: %w", idx, err)
604+
return
605+
}
606+
created.Close()
607+
errCh <- nil
608+
}(i)
609+
}
610+
611+
// Start all goroutines at once
612+
close(start)
613+
614+
// Wait for all to complete - all should succeed (last write wins)
615+
for i := 0; i < numWriters; i++ {
616+
if err := <-errCh; err != nil {
617+
t.Error(err)
618+
}
619+
}
620+
621+
// Verify exactly one object exists
622+
objects, err := storage.ListObjects(bucketName, "", false)
623+
noError(t, err)
624+
if len(objects) != 1 {
625+
t.Errorf("expected 1 object, got %d", len(objects))
626+
}
627+
})
628+
}
629+
630+
// newStreamingContent creates a ReadSeekCloser from a byte slice
631+
func newStreamingContent(data []byte) io.ReadSeekCloser {
632+
return &bytesReadSeekCloser{Reader: bytes.NewReader(data)}
633+
}
634+
635+
type bytesReadSeekCloser struct {
636+
*bytes.Reader
637+
}
638+
639+
func (b *bytesReadSeekCloser) Close() error {
640+
return nil
641+
}
642+
377643
func compareStreamingObjects(o1, o2 StreamingObject) error {
378644
if o1.BucketName != o2.BucketName {
379645
return fmt.Errorf("bucket name differs:\nmain %q\narg %q", o1.BucketName, o2.BucketName)

0 commit comments

Comments
 (0)