fix(model): close fd immediately in copier (#10079)
This commit is contained in:
@@ -1344,7 +1344,7 @@ func (f *sendReceiveFolder) copierRoutine(in <-chan copyBlocksState, pullChan ch
|
|||||||
|
|
||||||
buf = protocol.BufferPool.Upgrade(buf, int(block.Size))
|
buf = protocol.BufferPool.Upgrade(buf, int(block.Size))
|
||||||
|
|
||||||
found := false
|
copied := false
|
||||||
blocks, _ := f.model.sdb.AllLocalBlocksWithHash(block.Hash)
|
blocks, _ := f.model.sdb.AllLocalBlocksWithHash(block.Hash)
|
||||||
innerBlocks:
|
innerBlocks:
|
||||||
for _, e := range blocks {
|
for _, e := range blocks {
|
||||||
@@ -1355,47 +1355,19 @@ func (f *sendReceiveFolder) copierRoutine(in <-chan copyBlocksState, pullChan ch
|
|||||||
for folderID, files := range res {
|
for folderID, files := range res {
|
||||||
ffs := folderFilesystems[folderID]
|
ffs := folderFilesystems[folderID]
|
||||||
for _, fi := range files {
|
for _, fi := range files {
|
||||||
fd, err := ffs.Open(fi.Name)
|
copied, err = f.copyBlock(fi.Name, e.Offset, dstFd, ffs, block, buf)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
continue
|
state.fail(err)
|
||||||
}
|
|
||||||
defer fd.Close()
|
|
||||||
|
|
||||||
_, err = fd.ReadAt(buf, e.Offset)
|
|
||||||
if err != nil {
|
|
||||||
fd.Close()
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
// Hash is not SHA256 as it's an encrypted hash token. In that
|
|
||||||
// case we can't verify the block integrity so we'll take it on
|
|
||||||
// trust. (The other side can and will verify.)
|
|
||||||
if f.Type != config.FolderTypeReceiveEncrypted {
|
|
||||||
if err := f.verifyBuffer(buf, block); err != nil {
|
|
||||||
l.Debugln("Finder failed to verify buffer", err)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if f.CopyRangeMethod != config.CopyRangeMethodStandard {
|
|
||||||
err = f.withLimiter(func() error {
|
|
||||||
dstFd.mut.Lock()
|
|
||||||
defer dstFd.mut.Unlock()
|
|
||||||
return fs.CopyRange(f.CopyRangeMethod.ToFS(), fd, dstFd.fd, e.Offset, block.Offset, int64(block.Size))
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
err = f.limitedWriteAt(dstFd, buf, block.Offset)
|
|
||||||
}
|
|
||||||
if err != nil {
|
|
||||||
state.fail(fmt.Errorf("dst write: %w", err))
|
|
||||||
break innerBlocks
|
break innerBlocks
|
||||||
}
|
}
|
||||||
|
if !copied {
|
||||||
|
continue
|
||||||
|
}
|
||||||
if fi.Name == state.file.Name {
|
if fi.Name == state.file.Name {
|
||||||
state.copiedFromOrigin(block.Size)
|
state.copiedFromOrigin(block.Size)
|
||||||
} else {
|
} else {
|
||||||
state.copiedFromElsewhere(block.Size)
|
state.copiedFromElsewhere(block.Size)
|
||||||
}
|
}
|
||||||
found = true
|
|
||||||
break innerBlocks
|
break innerBlocks
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1405,7 +1377,7 @@ func (f *sendReceiveFolder) copierRoutine(in <-chan copyBlocksState, pullChan ch
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
if !found {
|
if !copied {
|
||||||
state.pullStarted()
|
state.pullStarted()
|
||||||
ps := pullBlockState{
|
ps := pullBlockState{
|
||||||
sharedPullerState: state.sharedPullerState,
|
sharedPullerState: state.sharedPullerState,
|
||||||
@@ -1421,6 +1393,47 @@ func (f *sendReceiveFolder) copierRoutine(in <-chan copyBlocksState, pullChan ch
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Returns true, if the block was successfully copied to buf. The returned
|
||||||
|
// errors is *only* non-nil for errors that are fatal for pulling this entire
|
||||||
|
// file (i.e. writing to dstFd fails). For other errors that mean the block
|
||||||
|
// wasn't copied, but we can continue trying, the error is nil and the bool
|
||||||
|
// false.
|
||||||
|
// The buffer must be big enough to hold the block. It's only purpose is
|
||||||
|
// efficiency resp. reuse, the caller must not use/rely on its contents.
|
||||||
|
func (f *sendReceiveFolder) copyBlock(srcName string, srcOffset int64, dstFd *lockedWriterAt, ffs fs.Filesystem, block protocol.BlockInfo, buf []byte) (bool, error) {
|
||||||
|
fd, err := ffs.Open(srcName)
|
||||||
|
if err != nil {
|
||||||
|
l.Debugf("Failed to open file %v trying to copy block %v (folderID %v): %v", srcName, block.Hash, f.folderID, err)
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
defer fd.Close()
|
||||||
|
|
||||||
|
_, err = fd.ReadAt(buf, srcOffset)
|
||||||
|
if err != nil {
|
||||||
|
l.Debugf("Failed to read block from file %v in copier (folderID: %v, hash: %v): %v", srcName, f.folderID, block.Hash, err)
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Hash is not SHA256 as it's an encrypted hash token. In that
|
||||||
|
// case we can't verify the block integrity so we'll take it on
|
||||||
|
// trust. (The other side can and will verify.)
|
||||||
|
if f.Type != config.FolderTypeReceiveEncrypted {
|
||||||
|
if err := f.verifyBuffer(buf, block); err != nil {
|
||||||
|
l.Debugf("Failed to verify buffer in copier (folderID: %v): %v", f.folderID, err)
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if f.CopyRangeMethod != config.CopyRangeMethodStandard {
|
||||||
|
return true, f.withLimiter(func() error {
|
||||||
|
dstFd.mut.Lock()
|
||||||
|
defer dstFd.mut.Unlock()
|
||||||
|
return fs.CopyRange(f.CopyRangeMethod.ToFS(), fd, dstFd.fd, srcOffset, block.Offset, int64(block.Size))
|
||||||
|
})
|
||||||
|
}
|
||||||
|
return true, f.limitedWriteAt(dstFd, buf, block.Offset)
|
||||||
|
}
|
||||||
|
|
||||||
func (*sendReceiveFolder) verifyBuffer(buf []byte, block protocol.BlockInfo) error {
|
func (*sendReceiveFolder) verifyBuffer(buf []byte, block protocol.BlockInfo) error {
|
||||||
if len(buf) != int(block.Size) {
|
if len(buf) != int(block.Size) {
|
||||||
return fmt.Errorf("length mismatch %d != %d", len(buf), block.Size)
|
return fmt.Errorf("length mismatch %d != %d", len(buf), block.Size)
|
||||||
|
|||||||
Reference in New Issue
Block a user