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
36 changes: 36 additions & 0 deletions component/attr_cache/attr_cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1420,6 +1420,42 @@ func (suite *attrCacheTestSuite) TestTruncateFile() {
suite.assert.True(checkItem.exists())
}

// Attributes returned by GetAttr are read by callers after the cache lock is released,
// so later cache updates must not modify them in place
func (suite *attrCacheTestSuite) TestGetAttrResultNotModifiedByCacheUpdates() {
defer suite.cleanupTest()
path := "a"
suite.addPathToCache(path)

attr, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: path})
suite.assert.NoError(err)
before := *attr

readerDone := make(chan struct{})
go func() {
defer close(readerDone)
for range 1000 {
_ = attr.Size
_ = attr.Mtime
_ = attr.Mode
}
}()

truncateOptions := internal.TruncateFileOptions{Name: path, NewSize: 1234}
suite.mock.EXPECT().TruncateFile(truncateOptions).Return(nil)
suite.assert.NoError(suite.attrCache.TruncateFile(truncateOptions))
chmodOptions := internal.ChmodOptions{Name: path, Mode: 0600}
suite.mock.EXPECT().Chmod(chmodOptions).Return(nil)
suite.assert.NoError(suite.attrCache.Chmod(chmodOptions))
<-readerDone

suite.assert.Equal(before, *attr)
updated, err := suite.attrCache.GetAttr(internal.GetAttrOptions{Name: path})
suite.assert.NoError(err)
suite.assert.EqualValues(1234, updated.Size)
suite.assert.Equal(os.FileMode(0600), updated.Mode.Perm())
}

// Tests CopyFromFile
func (suite *attrCacheTestSuite) TestCopyFromFileError() {
defer suite.cleanupTest()
Expand Down
10 changes: 10 additions & 0 deletions component/attr_cache/cacheMap.go
Original file line number Diff line number Diff line change
Expand Up @@ -265,7 +265,15 @@ func (value *attrCacheItem) markInCloud(inCloud bool) {
}
}

// cloneAttr replaces the item's attributes with a private copy so they can be modified.
// Callers may still be reading the previous copy after the cache lock was released.
func (value *attrCacheItem) cloneAttr() {
attr := *value.attr
value.attr = &attr
}

func (value *attrCacheItem) setSize(size int64, changedAt time.Time) {
value.cloneAttr()
value.attr.Mtime = changedAt
value.attr.Ctime = changedAt
value.attr.Size = size
Expand All @@ -276,6 +284,7 @@ func (value *attrCacheItem) touchModifyAndChangeTimes(changedAt time.Time) {
if value == nil || !value.exists() {
return
}
value.cloneAttr()
value.attr.Mtime = changedAt
value.attr.Ctime = changedAt
value.cachedAt = changedAt
Expand All @@ -287,6 +296,7 @@ func (value *attrCacheItem) setMode(mode os.FileMode) {
currentType = mode & os.ModeType
}
modeBits := mode & (os.ModePerm | os.ModeSetuid | os.ModeSetgid | os.ModeSticky)
value.cloneAttr()
value.attr.Mode = currentType | modeBits
value.attr.Flags.Clear(internal.PropFlagModeDefault)
value.attr.Ctime = time.Now()
Expand Down
2 changes: 1 addition & 1 deletion component/azstorage/block_blob_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -356,7 +356,7 @@ func generateContainerName() string {

func createTestContainerWithRetry(create func() error) error {
var err error
for i := 0; i < 5; i++ {
for i := range 5 {
err = create()
if err == nil {
return nil
Expand Down
25 changes: 12 additions & 13 deletions component/block_cache/block_cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -1190,6 +1190,7 @@ func (bc *BlockCache) lineupDownload(handle *handlemap.Handle, block *Block, pre
failCnt: 0,
upload: false,
ETag: Etag,
fileSize: handle.Size,
}

// Remove this block from free block list and add to in-process list
Expand Down Expand Up @@ -1258,12 +1259,12 @@ func (bc *BlockCache) download(item *workItem) {
}

if numberOfBytes != int(bc.blockSize) &&
item.block.offset+uint64(numberOfBytes) != uint64(item.handle.Size) {
item.block.offset+uint64(numberOfBytes) != uint64(item.fileSize) {
log.Err(
"BlockCache::download : Local data retrieved from disk size mismatch, Expected %v, OnDisk %v, fileSize %v",
bc.getBlockSize(uint64(item.handle.Size), item.block),
bc.getBlockSize(uint64(item.fileSize), item.block),
numberOfBytes,
item.handle.Size,
item.fileSize,
)
successfulRead = false
f.Close()
Expand Down Expand Up @@ -1305,8 +1306,7 @@ func (bc *BlockCache) download(item *workItem) {
if item.failCnt > MAX_FAIL_CNT {
// If we failed to read the data 3 times then just give up
log.Err(
"BlockCache::download : 3 attempts to download a block have failed %v=>%s (index %v, offset %v)",
item.handle.ID,
"BlockCache::download : 3 attempts to download a block have failed %s (index %v, offset %v)",
item.handle.Path,
item.block.id,
item.block.offset,
Expand All @@ -1319,8 +1319,7 @@ func (bc *BlockCache) download(item *workItem) {
if err != nil && err != io.EOF {
// Fail to read the data so just reschedule this request
log.Err(
"BlockCache::download : Failed to read %v=>%s from offset %v [%s]",
item.handle.ID,
"BlockCache::download : Failed to read %s from offset %v [%s]",
item.handle.Path,
item.block.id,
err.Error(),
Expand All @@ -1331,8 +1330,7 @@ func (bc *BlockCache) download(item *workItem) {
} else if n == 0 {
// No data read so just reschedule this request
log.Err(
"BlockCache::download : Failed to read %v=>%s from offset %v [0 bytes read]",
item.handle.ID,
"BlockCache::download : Failed to read %s from offset %v [0 bytes read]",
item.handle.Path,
item.block.id,
)
Expand All @@ -1345,8 +1343,7 @@ func (bc *BlockCache) download(item *workItem) {
if etag != "" {
if item.ETag != "" && item.ETag != etag {
log.Err(
"BlockCache::download : Blob has changed for %v=>%s (index %v, offset %v)",
item.handle.ID,
"BlockCache::download : Blob has changed for %s (index %v, offset %v)",
item.handle.Path,
item.block.id,
item.block.offset,
Expand Down Expand Up @@ -1600,7 +1597,8 @@ func (bc *BlockCache) getOrCreateBlock(handle *handlemap.Handle, offset uint64)
block = node.(*Block)

// If the block was staged earlier then we are overwriting it here so move it back to cooking queue
if block.flags.IsSet(BlockFlagSynced) {
// The upload worker sets Synced before it signals completion, so wait for it if it is still uploading
if block.flags.IsSet(BlockFlagSynced) && !block.flags.IsSet(BlockFlagUploading) {
log.Debug(
"BlockCache::getOrCreateBlock : Overwriting back to staged block %v for %v=>%s",
block.id,
Expand Down Expand Up @@ -1777,6 +1775,7 @@ func (bc *BlockCache) lineupUpload(
failCnt: 0,
upload: true,
blockId: id,
fileSize: handle.Size,
}

block.Uploading()
Expand Down Expand Up @@ -1854,7 +1853,7 @@ func (bc *BlockCache) upload(item *workItem) {
flock := bc.fileLocks.Get(fileName)
flock.Lock()
defer flock.Unlock()
blockSize := bc.getBlockSize(uint64(item.handle.Size), item.block)
blockSize := bc.getBlockSize(uint64(item.fileSize), item.block)
// This block is updated so we need to stage it now
err := bc.NextComponent().StageData(internal.StageDataOptions{
Name: item.handle.Path,
Expand Down
118 changes: 117 additions & 1 deletion component/block_cache/block_cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1445,7 +1445,8 @@ func (suite *blockCacheTestSuite) TestZZZZLazyWrite() {
suite.assert.True(handle.Dirty())

_ = tobj.blockCache.ReleaseFile(internal.ReleaseFileOptions{Handle: handle})
time.Sleep(1 * time.Second)
// wait for the async close before turning lazy write back off
tobj.blockCache.fileCloseOpt.Wait()
tobj.blockCache.lazyWrite = false

// As lazy write is enabled flush shall not upload the file
Expand Down Expand Up @@ -3320,6 +3321,121 @@ func (suite *blockCacheTestSuite) TestReadCommittedLastBlocksOverwrite() {
suite.assert.Equal(h.Size, int64((15*_1MB)+(_1MB/2)))
}

// An overwrite must wait for an in-flight upload of the same block. The upload worker marks
// the block Synced before it clears the dirty bit and signals completion, so a writer that
// trusts Synced alone can have its dirty bit wiped and its data never uploaded.
func (suite *blockCacheTestSuite) TestOverwriteWaitsForInFlightUpload() {
cfg := "block_cache:\n block-size-mb: 1\n mem-size-mb: 20\n prefetch: 12\n parallelism: 10"
tobj, err := setupPipeline(cfg)
defer tobj.cleanupPipeline()
suite.assert.NoError(err)

path := getTestFileName(suite.T().Name())
h, err := tobj.blockCache.CreateFile(internal.CreateFileOptions{Name: path, Mode: 0777})
suite.assert.NoError(err)

_, err = tobj.blockCache.WriteFile(
&internal.WriteFileOptions{Handle: h, Offset: 0, Data: dataBuff[:10]},
)
suite.assert.NoError(err)
node, found := h.GetValue("0")
suite.assert.True(found)
block := node.(*Block)

// Put the block in the state an upload worker leaves it in just before it finishes:
// queued for upload and already marked Synced, but still dirty and not yet signalled.
h.Lock()
block.Uploading()
block.flags.Set(BlockFlagUploading)
block.flags.Set(BlockFlagSynced)
tobj.blockCache.addToCooked(h, block)
h.Unlock()

done := make(chan error, 1)
go func() {
_, werr := tobj.blockCache.WriteFile(
&internal.WriteFileOptions{Handle: h, Offset: 0, Data: dataBuff[10:30]},
)
done <- werr
}()

select {
case <-done:
suite.assert.Fail("overwrite did not wait for the in-flight upload to finish")
return
case <-time.After(200 * time.Millisecond):
}

// Let the simulated upload worker finish
block.NoMoreDirty()
block.Ready(BlockStatusUploaded)

select {
case err = <-done:
suite.assert.NoError(err)
case <-time.After(5 * time.Second):
suite.assert.Fail("overwrite never completed")
}
suite.assert.True(block.IsDirty(), "overwrite lost its dirty bit")

err = tobj.blockCache.ReleaseFile(internal.ReleaseFileOptions{Handle: h})
suite.assert.NoError(err)

data, err := os.ReadFile(filepath.Join(tobj.fake_storage_path, path))
suite.assert.NoError(err)
suite.assert.Equal(dataBuff[10:30], data)
}

// An upload must stage the block at the size recorded in the block list when it was lined up.
// If it reads the live handle size instead, a concurrent write that extends the file makes it
// stage a full block while the block list still records the short size, and the commit then
// pads it with a filler block, corrupting the file.
func (suite *blockCacheTestSuite) TestUploadRacingWriteThatExtendsFile() {
cfg := "block_cache:\n block-size-mb: 1\n mem-size-mb: 20\n prefetch: 12\n parallelism: 10"
tobj, err := setupPipeline(cfg)
defer tobj.cleanupPipeline()
suite.assert.NoError(err)

for i := range 20 {
path := fmt.Sprintf("%s_%d", getTestFileName(suite.T().Name()), i)
h, err := tobj.blockCache.CreateFile(internal.CreateFileOptions{Name: path, Mode: 0777})
suite.assert.NoError(err)

_, err = tobj.blockCache.WriteFile(
&internal.WriteFileOptions{Handle: h, Offset: 0, Data: dataBuff[:10]},
)
suite.assert.NoError(err)

h.Lock()
err = tobj.blockCache.stageBlocks(h, 1)
lst, _ := h.GetValue("blockList")
staged := *lst.(map[int64]*blockInfo)[0]
h.Unlock()
suite.assert.NoError(err)

// extend the file while block 0 may still be uploading
_, err = tobj.blockCache.WriteFile(
&internal.WriteFileOptions{Handle: h, Offset: int64(2 * _1MB), Data: dataBuff[:10]},
)
suite.assert.NoError(err)

// wait for the upload of block 0, then compare what was staged with what was recorded
h.Lock()
tobj.blockCache.waitAndFreeUploadedBlocks(h, 1)
h.Unlock()
stagedPath := filepath.Join(tobj.fake_storage_path, path) + "_" +
strings.ReplaceAll(staged.id, "/", "_")
fi, err := os.Stat(stagedPath)
suite.assert.NoError(err)
if err == nil {
suite.assert.Equal(int64(staged.size), fi.Size(), "iteration %d", i)
}

err = tobj.blockCache.ReleaseFile(internal.ReleaseFileOptions{Handle: h})
suite.assert.NoError(err)
}
}

// In order for 'go test' to run this suite, we need to create
// a normal test function and pass our suite to suite.Run
func TestBlockCacheTestSuite(t *testing.T) {
Expand Down
1 change: 1 addition & 0 deletions component/block_cache/threadpool.go
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ type workItem struct {
upload bool // Flag marking this is a upload request or not
blockId string // BlockId of the block
ETag string // Etag of the file before scheduling.
fileSize int64 // Size of the file when this item was scheduled
Comment thread
jfantinhardesty marked this conversation as resolved.
}

// Reason for storing Etag in workitem struct:
Expand Down
12 changes: 7 additions & 5 deletions component/block_cache/threadpool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -126,8 +126,9 @@ func (suite *threadPoolTestSuite) TestPrioritySchedule() {
tp.Schedule(i < 20, &workItem{failCnt: 5})
}

time.Sleep(100 * time.Millisecond)
suite.assert.Equal(int32(100), callbackCnt)
suite.assert.Eventually(func() bool {
return atomic.LoadInt32(&callbackCnt) == 100
}, 5*time.Second, 10*time.Millisecond)
tp.Stop()
}

Expand Down Expand Up @@ -158,9 +159,10 @@ func (suite *threadPoolTestSuite) TestPriorityScheduleWithWriter() {
tp.Schedule(i < 20, &workItem{failCnt: 5, upload: true, blockId: "test"})
}

time.Sleep(100 * time.Millisecond)
suite.assert.Equal(int32(100), callbackWCnt)
suite.assert.Equal(int32(0), callbackRCnt)
suite.assert.Eventually(func() bool {
return atomic.LoadInt32(&callbackWCnt) == 100
}, 5*time.Second, 10*time.Millisecond)
suite.assert.Equal(int32(0), atomic.LoadInt32(&callbackRCnt))
tp.Stop()
}

Expand Down
Loading
Loading