Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ func (t *SingleMountReadsTestSuite) runAppendAndReadTest(verifyFunc readAndVerif
for i := range numAppends {
// Wait for a minute for stat to return the correct file size, which is needed by appendToFile.
if i > 0 {
time.Sleep(operations.WaitDurationAfterFlushZB)
operations.WaitForSizeUpdate(setup.IsZonalBucketRun() || setup.IsPirloBucketRun(), operations.WaitDurationAfterFlushRapid)
}

t.appendToFile(appendFileHandle, setup.GenerateRandomString(appendSize))
Expand Down Expand Up @@ -89,7 +89,7 @@ func (t *DualMountReadsTestSuite) runAppendAndReadTest(verifyFunc readAndVerifyF
// Wait for metadata cache to expire to fetch the latest size for the next read.
// Metadata update for appends in current iteration itself takes a minute, so the
// cached size will expire in ttl-60 secs from now, so wait accordingly.
time.Sleep(time.Duration(metadataCacheTTLSecs*time.Second - operations.WaitDurationAfterFlushZB))
time.Sleep(metadataCacheTTLSecs*time.Second - operations.WaitDurationAfterFlushRapid)
// Expect read up to the latest file size which is the size after the append.
verifyFunc(t.T(), readPath, []byte(t.fileContent[:sizeAfterAppend]))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -140,7 +140,8 @@ func (s *BaseSymlinkSuite) createGCSSymlinkObject(linkName, target string) {
_, err := w.Write(content)
s.Require().NoError(err)
s.Require().NoError(w.Close())
operations.WaitForSizeUpdate(setup.IsZonalBucketRun(), operations.WaitDurationAfterCloseZB)
isUnfinalizedObject := w.Append && !w.FinalizeOnClose
operations.WaitForSizeUpdate(isUnfinalizedObject, operations.WaitDurationAfterCloseRapid)
}

////////////////////////////////////////////////////////////////////////
Expand Down
6 changes: 4 additions & 2 deletions tools/integration_tests/util/client/storage_client.go
Original file line number Diff line number Diff line change
Expand Up @@ -292,7 +292,8 @@ func CreateObjectWithOptions(ctx context.Context, client *storage.Client, object
if err := wc.Close(); err != nil {
return fmt.Errorf("wc.Close failed for object %q: %w", object, err)
}
operations.WaitForSizeUpdate(setup.IsZonalBucketRun(), operations.WaitDurationAfterCloseZB)
isUnfinalizedObject := wc.Append && !wc.FinalizeOnClose
operations.WaitForSizeUpdate(isUnfinalizedObject, operations.WaitDurationAfterCloseRapid)
return nil
}

Expand Down Expand Up @@ -406,7 +407,8 @@ func UploadGcsObjectWithPreconditions(ctx context.Context, client *storage.Clien
if err := w.Close(); err != nil {
log.Printf("Failed to close GCS object gs://%s/%s: %v", bucketName, objectName, err)
}
operations.WaitForSizeUpdate(setup.IsZonalBucketRun(), operations.WaitDurationAfterCloseZB)
isUnfinalizedObject := w.Append && !w.FinalizeOnClose
operations.WaitForSizeUpdate(isUnfinalizedObject, operations.WaitDurationAfterCloseRapid)
}()

filePathToUpload := localPath
Expand Down
20 changes: 13 additions & 7 deletions tools/integration_tests/util/operations/file_operations.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,9 @@ const (
TimeSlop = 25 * time.Millisecond
// TmpDirectory specifies the directory where temporary files will be created.
// In this case, we are using the system's default temporary directory.
TmpDirectory = "/tmp"
WaitDurationAfterFlushZB = time.Minute
WaitDurationAfterCloseZB = time.Second
TmpDirectory = "/tmp"
WaitDurationAfterFlushRapid = time.Minute
WaitDurationAfterCloseRapid = time.Second
)

func copyFile(srcFileName, dstFileName string, allowOverwrite bool) (err error) {
Expand Down Expand Up @@ -174,15 +174,19 @@ func CloseFiles(t *testing.T, files []*os.File) {
err := file.Close()
assert.NoError(t, err)
}
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterCloseZB)
// Pirlo creates finalized objects on close by default, so we don't wait for it here.
// Update this condition if this method is used to create unfinalized Pirlo objects in the future.
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterCloseRapid)
Comment thread
vipnydav marked this conversation as resolved.
}

// Deprecated: please use CloseFileShouldNotThrowError instead.
func CloseFile(file *os.File) {
if err := file.Close(); err != nil {
log.Fatalf("error in closing: %v", err)
}
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterCloseZB)
// Pirlo creates finalized objects on close by default, so we don't wait for it here.
// Update this condition if this method is used to create unfinalized Pirlo objects in the future.
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterCloseRapid)
}

func RemoveFile(filePath string) {
Expand Down Expand Up @@ -581,7 +585,9 @@ func WriteAt(content string, offset int64, fh *os.File, t testing.TB) {
func CloseFileShouldNotThrowError(t testing.TB, file *os.File) {
err := file.Close()
assert.NoError(t, err)
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterCloseZB)
// Pirlo creates finalized objects on close by default, so we don't wait for it here.
// Update this condition if this method is used to create unfinalized Pirlo objects in the future.
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterCloseRapid)
}

func CloseFileShouldThrowError(t *testing.T, file *os.File) {
Expand All @@ -598,7 +604,7 @@ func SyncFile(fh *os.File, t *testing.T) {
if err != nil {
t.Fatalf("%s.Sync(): %v", fh.Name(), err)
}
WaitForSizeUpdate(setup.IsZonalBucketRun(), WaitDurationAfterFlushZB)
WaitForSizeUpdate(setup.IsZonalBucketRun() || setup.IsPirloBucketRun(), WaitDurationAfterFlushRapid)
}

func SyncFiles(files []*os.File, t *testing.T) {
Expand Down
4 changes: 2 additions & 2 deletions tools/integration_tests/util/operations/operations.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,8 +81,8 @@ func ExecuteGcloudCommand(command string) ([]byte, error) {

// WaitForSizeUpdate waits for a specified time duration to ensure that stat()
// call returns correct size for unfinalized object.
func WaitForSizeUpdate(isZonal bool, duration time.Duration) {
if isZonal {
func WaitForSizeUpdate(isUnfinalized bool, duration time.Duration) {
if isUnfinalized {
time.Sleep(duration)
}
}
Loading