You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
351 lines
11 KiB
351 lines
11 KiB
9 years ago
|
/*
|
||
|
* Minio Cloud Storage, (C) 2016 Minio, Inc.
|
||
|
*
|
||
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
||
|
* you may not use this file except in compliance with the License.
|
||
|
* You may obtain a copy of the License at
|
||
|
*
|
||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||
|
*
|
||
|
* Unless required by applicable law or agreed to in writing, software
|
||
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||
|
* See the License for the specific language governing permissions and
|
||
|
* limitations under the License.
|
||
|
*/
|
||
|
|
||
9 years ago
|
package cmd
|
||
9 years ago
|
|
||
9 years ago
|
import (
|
||
|
"encoding/hex"
|
||
|
"errors"
|
||
9 years ago
|
"io"
|
||
9 years ago
|
"sync"
|
||
9 years ago
|
|
||
|
"github.com/klauspost/reedsolomon"
|
||
9 years ago
|
"github.com/minio/minio/pkg/bpool"
|
||
9 years ago
|
)
|
||
9 years ago
|
|
||
9 years ago
|
// isSuccessDecodeBlocks - do we have all the blocks to be
|
||
|
// successfully decoded?. Input encoded blocks ordered matrix.
|
||
|
func isSuccessDecodeBlocks(enBlocks [][]byte, dataBlocks int) bool {
|
||
9 years ago
|
// Count number of data and parity blocks that were read.
|
||
|
var successDataBlocksCount = 0
|
||
|
var successParityBlocksCount = 0
|
||
9 years ago
|
for index := range enBlocks {
|
||
|
if enBlocks[index] == nil {
|
||
9 years ago
|
continue
|
||
|
}
|
||
9 years ago
|
// block index lesser than data blocks, update data block count.
|
||
9 years ago
|
if index < dataBlocks {
|
||
|
successDataBlocksCount++
|
||
|
continue
|
||
9 years ago
|
} // else { // update parity block count.
|
||
9 years ago
|
successParityBlocksCount++
|
||
|
}
|
||
9 years ago
|
// Returns true if we have atleast dataBlocks parity.
|
||
|
return successDataBlocksCount == dataBlocks || successDataBlocksCount+successParityBlocksCount >= dataBlocks
|
||
9 years ago
|
}
|
||
|
|
||
|
// isSuccessDataBlocks - do we have all the data blocks?
|
||
9 years ago
|
// Input encoded blocks ordered matrix.
|
||
|
func isSuccessDataBlocks(enBlocks [][]byte, dataBlocks int) bool {
|
||
9 years ago
|
// Count number of data blocks that were read.
|
||
|
var successDataBlocksCount = 0
|
||
9 years ago
|
for index := range enBlocks[:dataBlocks] {
|
||
|
if enBlocks[index] == nil {
|
||
9 years ago
|
continue
|
||
|
}
|
||
9 years ago
|
// block index lesser than data blocks, update data block count.
|
||
9 years ago
|
if index < dataBlocks {
|
||
|
successDataBlocksCount++
|
||
|
}
|
||
|
}
|
||
9 years ago
|
// Returns true if we have atleast the dataBlocks.
|
||
9 years ago
|
return successDataBlocksCount >= dataBlocks
|
||
|
}
|
||
|
|
||
9 years ago
|
// Return readable disks slice from which we can read parallelly.
|
||
9 years ago
|
func getReadDisks(orderedDisks []StorageAPI, index int, dataBlocks int) (readDisks []StorageAPI, nextIndex int, err error) {
|
||
|
readDisks = make([]StorageAPI, len(orderedDisks))
|
||
|
dataDisks := 0
|
||
|
parityDisks := 0
|
||
|
// Count already read data and parity chunks.
|
||
|
for i := 0; i < index; i++ {
|
||
|
if orderedDisks[i] == nil {
|
||
|
continue
|
||
|
}
|
||
|
if i < dataBlocks {
|
||
|
dataDisks++
|
||
|
} else {
|
||
|
parityDisks++
|
||
|
}
|
||
|
}
|
||
|
|
||
|
// Sanity checks - we should never have this situation.
|
||
|
if dataDisks == dataBlocks {
|
||
|
return nil, 0, errUnexpected
|
||
|
}
|
||
9 years ago
|
if dataDisks+parityDisks >= dataBlocks {
|
||
9 years ago
|
return nil, 0, errUnexpected
|
||
|
}
|
||
|
|
||
|
// Find the disks from which next set of parallel reads should happen.
|
||
|
for i := index; i < len(orderedDisks); i++ {
|
||
|
if orderedDisks[i] == nil {
|
||
|
continue
|
||
|
}
|
||
|
if i < dataBlocks {
|
||
|
dataDisks++
|
||
|
} else {
|
||
|
parityDisks++
|
||
|
}
|
||
|
readDisks[i] = orderedDisks[i]
|
||
|
if dataDisks == dataBlocks {
|
||
|
return readDisks, i + 1, nil
|
||
9 years ago
|
} else if dataDisks+parityDisks == dataBlocks {
|
||
9 years ago
|
return readDisks, i + 1, nil
|
||
|
}
|
||
|
}
|
||
|
return nil, 0, errXLReadQuorum
|
||
|
}
|
||
|
|
||
9 years ago
|
// parallelRead - reads chunks in parallel from the disks specified in []readDisks.
|
||
9 years ago
|
func parallelRead(volume, path string, readDisks []StorageAPI, orderedDisks []StorageAPI, enBlocks [][]byte, blockOffset int64, curChunkSize int64, bitRotVerify func(diskIndex int) bool, pool *bpool.BytePool) {
|
||
9 years ago
|
// WaitGroup to synchronise the read go-routines.
|
||
|
wg := &sync.WaitGroup{}
|
||
|
|
||
|
// Read disks in parallel.
|
||
|
for index := range readDisks {
|
||
|
if readDisks[index] == nil {
|
||
|
continue
|
||
|
}
|
||
|
wg.Add(1)
|
||
|
// Reads chunk from readDisk[index] in routine.
|
||
|
go func(index int) {
|
||
|
defer wg.Done()
|
||
|
|
||
|
// Verify bit rot for the file on this disk.
|
||
|
if !bitRotVerify(index) {
|
||
|
// So that we don't read from this disk for the next block.
|
||
|
orderedDisks[index] = nil
|
||
|
return
|
||
|
}
|
||
|
|
||
9 years ago
|
buf, err := pool.Get()
|
||
9 years ago
|
if err != nil {
|
||
9 years ago
|
errorIf(err, "unable to get buffer from byte pool")
|
||
9 years ago
|
orderedDisks[index] = nil
|
||
|
return
|
||
|
}
|
||
9 years ago
|
buf = buf[:curChunkSize]
|
||
9 years ago
|
|
||
9 years ago
|
_, err = readDisks[index].ReadFile(volume, path, blockOffset, buf)
|
||
|
if err != nil {
|
||
|
orderedDisks[index] = nil
|
||
|
return
|
||
|
}
|
||
|
enBlocks[index] = buf
|
||
9 years ago
|
}(index)
|
||
|
}
|
||
|
|
||
|
// Waiting for first routines to finish.
|
||
|
wg.Wait()
|
||
|
}
|
||
|
|
||
9 years ago
|
// erasureReadFile - read bytes from erasure coded files and writes to given writer.
|
||
|
// Erasure coded files are read block by block as per given erasureInfo and data chunks
|
||
9 years ago
|
// are decoded into a data block. Data block is trimmed for given offset and length,
|
||
|
// then written to given writer. This function also supports bit-rot detection by
|
||
9 years ago
|
// verifying checksum of individual block's checksum.
|
||
9 years ago
|
func erasureReadFile(writer io.Writer, disks []StorageAPI, volume string, path string, offset int64, length int64, totalLength int64, blockSize int64, dataBlocks int, parityBlocks int, checkSums []string, algo string, pool *bpool.BytePool) (int64, error) {
|
||
9 years ago
|
// Offset and length cannot be negative.
|
||
|
if offset < 0 || length < 0 {
|
||
|
return 0, errUnexpected
|
||
|
}
|
||
|
|
||
9 years ago
|
// Can't request more data than what is available.
|
||
|
if offset+length > totalLength {
|
||
|
return 0, errUnexpected
|
||
|
}
|
||
|
|
||
9 years ago
|
// chunkSize is the amount of data that needs to be read from each disk at a time.
|
||
|
chunkSize := getChunkSize(blockSize, dataBlocks)
|
||
|
|
||
9 years ago
|
// bitRotVerify verifies if the file on a particular disk doesn't have bitrot
|
||
9 years ago
|
// by verifying the hash of the contents of the file.
|
||
9 years ago
|
bitRotVerify := func() func(diskIndex int) bool {
|
||
9 years ago
|
verified := make([]bool, len(disks))
|
||
9 years ago
|
// Return closure so that we have reference to []verified and
|
||
9 years ago
|
// not recalculate the hash on it every time the function is
|
||
9 years ago
|
// called for the same disk.
|
||
9 years ago
|
return func(diskIndex int) bool {
|
||
|
if verified[diskIndex] {
|
||
9 years ago
|
// Already validated.
|
||
9 years ago
|
return true
|
||
|
}
|
||
9 years ago
|
// Is this a valid block?
|
||
9 years ago
|
isValid := isValidBlock(disks[diskIndex], volume, path, checkSums[diskIndex], algo)
|
||
9 years ago
|
verified[diskIndex] = isValid
|
||
|
return isValid
|
||
|
}
|
||
|
}()
|
||
|
|
||
|
// Total bytes written to writer
|
||
|
bytesWritten := int64(0)
|
||
|
|
||
9 years ago
|
startBlock := offset / blockSize
|
||
|
endBlock := (offset + length) / blockSize
|
||
|
|
||
|
// curChunkSize = chunk size for the current block in the for loop below.
|
||
|
// curBlockSize = block size for the current block in the for loop below.
|
||
|
// curChunkSize and curBlockSize can change for the last block if totalLength%blockSize != 0
|
||
|
curChunkSize := chunkSize
|
||
|
curBlockSize := blockSize
|
||
9 years ago
|
|
||
|
// For each block, read chunk from each disk. If we are able to read all the data disks then we don't
|
||
|
// need to read parity disks. If one of the data disk is missing we need to read DataBlocks+1 number
|
||
|
// of disks. Once read, we Reconstruct() missing data if needed and write it to the given writer.
|
||
9 years ago
|
for block := startBlock; block <= endBlock; block++ {
|
||
9 years ago
|
// Mark all buffers as unused at the start of the loop so that the buffers
|
||
|
// can be reused.
|
||
|
pool.Reset()
|
||
|
|
||
9 years ago
|
// Each element of enBlocks holds curChunkSize'd amount of data read from its corresponding disk.
|
||
9 years ago
|
enBlocks := make([][]byte, len(disks))
|
||
9 years ago
|
|
||
9 years ago
|
if ((offset + bytesWritten) / blockSize) == (totalLength / blockSize) {
|
||
|
// This is the last block for which curBlockSize and curChunkSize can change.
|
||
|
// For ex. if totalLength is 15M and blockSize is 10MB, curBlockSize for
|
||
|
// the last block should be 5MB.
|
||
|
curBlockSize = totalLength % blockSize
|
||
|
curChunkSize = getChunkSize(curBlockSize, dataBlocks)
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
9 years ago
|
// NOTE: That for the offset calculation we have to use chunkSize and
|
||
|
// not curChunkSize. If we use curChunkSize for offset calculation
|
||
|
// then it can result in wrong offset for the last block.
|
||
|
blockOffset := block * chunkSize
|
||
9 years ago
|
|
||
9 years ago
|
// nextIndex - index from which next set of parallel reads
|
||
|
// should happen.
|
||
|
nextIndex := 0
|
||
|
|
||
9 years ago
|
for {
|
||
9 years ago
|
// readDisks - disks from which we need to read in parallel.
|
||
|
var readDisks []StorageAPI
|
||
|
var err error
|
||
9 years ago
|
// get readable disks slice from which we can read parallelly.
|
||
|
readDisks, nextIndex, err = getReadDisks(disks, nextIndex, dataBlocks)
|
||
9 years ago
|
if err != nil {
|
||
9 years ago
|
return bytesWritten, err
|
||
9 years ago
|
}
|
||
9 years ago
|
// Issue a parallel read across the disks specified in readDisks.
|
||
9 years ago
|
parallelRead(volume, path, readDisks, disks, enBlocks, blockOffset, curChunkSize, bitRotVerify, pool)
|
||
9 years ago
|
if isSuccessDecodeBlocks(enBlocks, dataBlocks) {
|
||
9 years ago
|
// If enough blocks are available to do rs.Reconstruct()
|
||
|
break
|
||
|
}
|
||
9 years ago
|
if nextIndex == len(disks) {
|
||
9 years ago
|
// No more disks to read from.
|
||
|
return bytesWritten, errXLReadQuorum
|
||
9 years ago
|
}
|
||
9 years ago
|
// We do not have enough enough data blocks to reconstruct the data
|
||
|
// hence continue the for-loop till we have enough data blocks.
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
9 years ago
|
// If we have all the data blocks no need to decode, continue to write.
|
||
9 years ago
|
if !isSuccessDataBlocks(enBlocks, dataBlocks) {
|
||
9 years ago
|
// Reconstruct the missing data blocks.
|
||
9 years ago
|
if err := decodeData(enBlocks, dataBlocks, parityBlocks); err != nil {
|
||
9 years ago
|
return bytesWritten, err
|
||
9 years ago
|
}
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
9 years ago
|
// Offset in enBlocks from where data should be read from.
|
||
|
enBlocksOffset := int64(0)
|
||
9 years ago
|
|
||
9 years ago
|
// Total data to be read from enBlocks.
|
||
|
enBlocksLength := curBlockSize
|
||
9 years ago
|
|
||
9 years ago
|
// If this is the start block then enBlocksOffset might not be 0.
|
||
9 years ago
|
if block == startBlock {
|
||
9 years ago
|
enBlocksOffset = offset % blockSize
|
||
|
enBlocksLength -= enBlocksOffset
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
9 years ago
|
remaining := length - bytesWritten
|
||
|
if remaining < enBlocksLength {
|
||
9 years ago
|
// We should not send more data than what was requested.
|
||
9 years ago
|
enBlocksLength = remaining
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
9 years ago
|
// Write data blocks.
|
||
9 years ago
|
n, err := writeDataBlocks(writer, enBlocks, dataBlocks, enBlocksOffset, enBlocksLength)
|
||
9 years ago
|
if err != nil {
|
||
|
return bytesWritten, err
|
||
|
}
|
||
9 years ago
|
|
||
|
// Update total bytes written.
|
||
9 years ago
|
bytesWritten += n
|
||
9 years ago
|
|
||
|
if bytesWritten == length {
|
||
|
// Done writing all the requested data.
|
||
|
break
|
||
|
}
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
9 years ago
|
// Success.
|
||
9 years ago
|
return bytesWritten, nil
|
||
9 years ago
|
}
|
||
9 years ago
|
|
||
|
// isValidBlock - calculates the checksum hash for the block and
|
||
|
// validates if its correct returns true for valid cases, false otherwise.
|
||
9 years ago
|
func isValidBlock(disk StorageAPI, volume, path, checkSum, checkSumAlgo string) (ok bool) {
|
||
9 years ago
|
// Disk is not available, not a valid block.
|
||
9 years ago
|
if disk == nil {
|
||
|
return false
|
||
9 years ago
|
}
|
||
9 years ago
|
// Checksum not available, not a valid block.
|
||
|
if checkSum == "" {
|
||
|
return false
|
||
|
}
|
||
9 years ago
|
// Read everything for a given block and calculate hash.
|
||
9 years ago
|
hashWriter := newHash(checkSumAlgo)
|
||
9 years ago
|
hashBytes, err := hashSum(disk, volume, path, hashWriter)
|
||
9 years ago
|
if err != nil {
|
||
9 years ago
|
errorIf(err, "Unable to calculate checksum %s/%s", volume, path)
|
||
|
return false
|
||
9 years ago
|
}
|
||
9 years ago
|
return hex.EncodeToString(hashBytes) == checkSum
|
||
9 years ago
|
}
|
||
|
|
||
|
// decodeData - decode encoded blocks.
|
||
|
func decodeData(enBlocks [][]byte, dataBlocks, parityBlocks int) error {
|
||
9 years ago
|
// Initialized reedsolomon.
|
||
9 years ago
|
rs, err := reedsolomon.New(dataBlocks, parityBlocks)
|
||
|
if err != nil {
|
||
|
return err
|
||
|
}
|
||
9 years ago
|
|
||
|
// Reconstruct encoded blocks.
|
||
9 years ago
|
err = rs.Reconstruct(enBlocks)
|
||
|
if err != nil {
|
||
|
return err
|
||
|
}
|
||
9 years ago
|
|
||
9 years ago
|
// Verify reconstructed blocks (parity).
|
||
|
ok, err := rs.Verify(enBlocks)
|
||
|
if err != nil {
|
||
|
return err
|
||
|
}
|
||
|
if !ok {
|
||
|
// Blocks cannot be reconstructed, corrupted data.
|
||
|
err = errors.New("Verification failed after reconstruction, data likely corrupted.")
|
||
|
return err
|
||
|
}
|
||
9 years ago
|
|
||
|
// Success.
|
||
9 years ago
|
return nil
|
||
|
}
|