mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-21 14:46:58 +00:00
* mount: re-resolve volume locations after a failed chunk read NewChunkGroup passed nil as the ReaderCache's CacheInvalidator, so retryFetchAfterCacheInvalidation was dead code on the FUSE read path. A mount that cached a volume's locations while one server was down kept retrying that server after it died, then returned EIO, even though the master and filer both resolved the live replica. The S3 gateway already passes its filerClient; do the same for the mount. * test: FUSE integration tests for volume server failover One mount appends while a second tails, and a volume server is killed, started or restarted mid-stream against a 001-replicated cluster of three volume servers. Automates the scenario matrix reported for Docker Swarm mounts, including the large-file variant and a no-chaos control. * test: report the filer's own view when append content mismatches A mismatch between what the writer wrote and what the reader sees can come from either side's cache. Read the file back through the filer's HTTP handler as well, and let the mount verbosity be raised from the environment, so a failing run says which layer lost the data. * test: wait for the reader mount to converge before comparing A mount caches metadata for about a second, so reading the file the instant the writer's last close returned can legitimately come back short. Poll the reader until it matches or the timeout expires; content that is wrong rather than merely late never converges and still fails, now with the writer's mount and the filer's own view alongside it. * test: detect a failover cluster child that exited at startup Signal(0) succeeds for a zombie and nothing reaped these children until shutdown, so a process that died on startup looked alive until the readiness timeout expired. Reap each child as it is started and consult the result. * test: read a file the killed volume server actually holds Placement decides which two of three servers back each volume, so killing volume N and reading readfile-N could pass without the victim ever holding a replica of it. Resolve each file's volumes through the filer and the master, and pick one the victim backs, preferring a file the reader has not cached. * ci: stop persisting checkout credentials in the failover workflow The job does not use the token after cloning. Also tag the README's command block as bash and match the timeout the workflow actually uses. * test: discard the ignored errors errcheck flags in the failover harness * test: resolve manifests when mapping a file to its volumes A manifest chunk's own fid names the volume holding the manifest, not the volumes holding the data, so a large enough file would point the failover victim at the wrong server. * test: pin the stale-location recovery path with a primed reader Reading a file for the first time after a server dies proves nothing: the lookup is fresh and returns the survivor. Kill one holder and wait for the master to drop it, read a file on that volume so the reader caches the lone survivor, restart the first server, then kill the survivor. The reader's only cached location is now dead while the data is live elsewhere, which is the case the invalidator exists for: EIO without it, recovery with it. * filer: re-look-up a chunk's locations as soon as they all fail A read that fails against every location it was given is far more likely to be holding a stale list than to be hitting a cluster that is briefly slow, but the retry loops spent the whole backoff ladder, about 13 s, before the caller got a chance to invalidate and look the chunk up again. Give the loops a refresh hook and let the reader cache invalidate on the first fully failed pass, so recovery starts in milliseconds. Clients without an invalidator keep the old behavior. The filer's streaming read path has its own fetch loop and is not covered. * filer: refresh locations on the random-read path too readChunkSliceAt bypasses the chunk cacher in random-access mode and fetches the range directly, which left it without the invalidation the cacher does: a random reader parked on a stale location had no way back at all. Hoist the refresh hook onto the reader cache so both paths share it. * filer: compare chunk locations as a set, not in order Lookups shuffle the locations they return, so comparing positionally reads a reshuffle of the very same replicas as a fresh set and spends an immediate retry on locations that just failed. weed/filer already had an order-independent comparison for this; move it next to the retry loops so both callers share one helper.
522 lines
17 KiB
Go
522 lines
17 KiB
Go
package filer
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"math"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"slices"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/security"
|
|
"github.com/seaweedfs/seaweedfs/weed/stats"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
|
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
|
)
|
|
|
|
var getLookupFileIdBackoffSchedule = []time.Duration{
|
|
150 * time.Millisecond,
|
|
600 * time.Millisecond,
|
|
1800 * time.Millisecond,
|
|
}
|
|
|
|
var (
|
|
jwtSigningReadKey security.SigningKey
|
|
jwtSigningReadKeyExpires int
|
|
loadJwtConfigOnce sync.Once
|
|
)
|
|
|
|
func loadJwtConfig() {
|
|
v := util.GetViper()
|
|
jwtSigningReadKey = security.SigningKey(v.GetString("jwt.signing.read.key"))
|
|
jwtSigningReadKeyExpires = v.GetInt("jwt.signing.read.expires_after_seconds")
|
|
if jwtSigningReadKeyExpires == 0 {
|
|
jwtSigningReadKeyExpires = 60
|
|
}
|
|
}
|
|
|
|
// JwtForVolumeServer generates a JWT token for volume server read operations if jwt.signing.read is configured
|
|
func JwtForVolumeServer(fileId string) string {
|
|
loadJwtConfigOnce.Do(loadJwtConfig)
|
|
if len(jwtSigningReadKey) == 0 {
|
|
return ""
|
|
}
|
|
return string(security.GenJwtForVolumeServer(jwtSigningReadKey, jwtSigningReadKeyExpires, fileId))
|
|
}
|
|
|
|
func HasData(entry *filer_pb.Entry) bool {
|
|
|
|
if len(entry.Content) > 0 {
|
|
return true
|
|
}
|
|
|
|
return len(entry.GetChunks()) > 0
|
|
}
|
|
|
|
func IsSameData(a, b *filer_pb.Entry) bool {
|
|
|
|
if len(a.Content) > 0 || len(b.Content) > 0 {
|
|
return bytes.Equal(a.Content, b.Content)
|
|
}
|
|
|
|
return isSameChunks(a.Chunks, b.Chunks)
|
|
}
|
|
|
|
func isSameChunks(a, b []*filer_pb.FileChunk) bool {
|
|
if len(a) != len(b) {
|
|
return false
|
|
}
|
|
slices.SortFunc(a, func(i, j *filer_pb.FileChunk) int {
|
|
return strings.Compare(i.ETag, j.ETag)
|
|
})
|
|
slices.SortFunc(b, func(i, j *filer_pb.FileChunk) int {
|
|
return strings.Compare(i.ETag, j.ETag)
|
|
})
|
|
for i := 0; i < len(a); i++ {
|
|
if a[i].ETag != b[i].ETag {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
func NewFileReader(filerClient filer_pb.FilerClient, entry *filer_pb.Entry) io.Reader {
|
|
if len(entry.Content) > 0 {
|
|
return bytes.NewReader(entry.Content)
|
|
}
|
|
return NewChunkStreamReader(filerClient, entry.GetChunks())
|
|
}
|
|
|
|
type DoStreamContent func(writer io.Writer) error
|
|
|
|
func PrepareStreamContent(masterClient wdclient.HasLookupFileIdFunction, jwtFunc VolumeServerJwtFunction, chunks []*filer_pb.FileChunk, offset int64, size int64) (DoStreamContent, error) {
|
|
return PrepareStreamContentWithThrottler(context.Background(), masterClient, jwtFunc, chunks, offset, size, 0)
|
|
}
|
|
|
|
type VolumeServerJwtFunction func(fileId string) string
|
|
|
|
// retryFetchWithFreshLocations is the shared self-heal for the read paths: when a chunk fetch
|
|
// fails, invalidate the cached volume locations, re-lookup, and call refetch only when the
|
|
// resolved locations actually changed (so we never retry against the same servers). originalErr
|
|
// is returned unchanged when no retry is attempted, so callers surface the real fetch failure.
|
|
func retryFetchWithFreshLocations(ctx context.Context, invalidator CacheInvalidator, lookupFn wdclient.LookupFileIdFunctionType, fileId string, oldUrls []string, originalErr error, refetch func(newUrls []string) error) error {
|
|
if invalidator == nil {
|
|
return originalErr
|
|
}
|
|
|
|
glog.V(0).InfofCtx(ctx, "read chunk %s failed, invalidating cache and retrying: %v", fileId, originalErr)
|
|
invalidator.InvalidateCache(fileId)
|
|
|
|
newUrls, lookupErr := lookupFn(ctx, fileId)
|
|
if lookupErr != nil {
|
|
glog.WarningfCtx(ctx, "failed to re-lookup chunk %s after cache invalidation: %v", fileId, lookupErr)
|
|
return fmt.Errorf("re-lookup chunk %s after cache invalidation: %w", fileId, lookupErr)
|
|
}
|
|
if len(newUrls) == 0 {
|
|
glog.WarningfCtx(ctx, "re-lookup for chunk %s returned no locations, skipping retry", fileId)
|
|
return fmt.Errorf("re-lookup chunk %s returned no locations", fileId)
|
|
}
|
|
if util_http.SameUrls(oldUrls, newUrls) {
|
|
glog.V(0).InfofCtx(ctx, "re-lookup returned same locations for chunk %s, skipping retry", fileId)
|
|
return originalErr
|
|
}
|
|
|
|
glog.V(0).InfofCtx(ctx, "retrying read chunk %s with %d new locations", fileId, len(newUrls))
|
|
return refetch(newUrls)
|
|
}
|
|
|
|
func PrepareStreamContentWithThrottler(ctx context.Context, masterClient wdclient.HasLookupFileIdFunction, jwtFunc VolumeServerJwtFunction, chunks []*filer_pb.FileChunk, offset int64, size int64, downloadMaxBytesPs int64) (DoStreamContent, error) {
|
|
glog.V(4).InfofCtx(ctx, "prepare to stream content for chunks: %d", len(chunks))
|
|
chunkViews := ViewFromChunks(ctx, masterClient.GetLookupFileIdFunction(), chunks, offset, size)
|
|
|
|
fileId2Url := make(map[string][]string)
|
|
|
|
for x := chunkViews.Front(); x != nil; x = x.Next {
|
|
chunkView := x.Value
|
|
var urlStrings []string
|
|
var err error
|
|
for _, backoff := range getLookupFileIdBackoffSchedule {
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
urlStrings, err = masterClient.GetLookupFileIdFunction()(ctx, chunkView.FileId)
|
|
if err == nil && len(urlStrings) > 0 {
|
|
break
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
glog.V(4).InfofCtx(ctx, "waiting for chunk: %s", chunkView.FileId)
|
|
timer := time.NewTimer(backoff)
|
|
select {
|
|
case <-ctx.Done():
|
|
if !timer.Stop() {
|
|
<-timer.C
|
|
}
|
|
return nil, ctx.Err()
|
|
case <-timer.C:
|
|
}
|
|
}
|
|
if err != nil {
|
|
glog.V(1).InfofCtx(ctx, "operation LookupFileId %s failed, err: %v", chunkView.FileId, err)
|
|
return nil, err
|
|
} else if len(urlStrings) == 0 {
|
|
errUrlNotFound := fmt.Errorf("operation LookupFileId %s failed, err: urls not found", chunkView.FileId)
|
|
glog.ErrorCtx(ctx, errUrlNotFound)
|
|
return nil, errUrlNotFound
|
|
}
|
|
fileId2Url[chunkView.FileId] = urlStrings
|
|
}
|
|
|
|
return func(writer io.Writer) error {
|
|
downloadThrottler := util.NewWriteThrottler(downloadMaxBytesPs)
|
|
remaining := size
|
|
for x := chunkViews.Front(); x != nil; x = x.Next {
|
|
chunkView := x.Value
|
|
if offset < chunkView.ViewOffset {
|
|
gap := chunkView.ViewOffset - offset
|
|
remaining -= gap
|
|
glog.V(4).InfofCtx(ctx, "zero [%d,%d)", offset, chunkView.ViewOffset)
|
|
err := writeZero(writer, gap)
|
|
if err != nil {
|
|
return fmt.Errorf("write zero [%d,%d)", offset, chunkView.ViewOffset)
|
|
}
|
|
offset = chunkView.ViewOffset
|
|
}
|
|
urlStrings := fileId2Url[chunkView.FileId]
|
|
start := time.Now()
|
|
jwt := jwtFunc(chunkView.FileId)
|
|
written, err := retriedStreamFetchChunkData(ctx, writer, urlStrings, jwt, chunkView.CipherKey, chunkView.IsGzipped, chunkView.IsFullChunk(), chunkView.OffsetInChunk, int(chunkView.ViewSize))
|
|
|
|
if err != nil && ctx.Err() != nil {
|
|
return ctx.Err()
|
|
}
|
|
|
|
// If read failed, try to invalidate cache and re-lookup
|
|
if err != nil && written == 0 {
|
|
invalidator, _ := masterClient.(CacheInvalidator)
|
|
err = retryFetchWithFreshLocations(ctx, invalidator, masterClient.GetLookupFileIdFunction(), chunkView.FileId, urlStrings, err, func(newUrls []string) error {
|
|
_, refetchErr := retriedStreamFetchChunkData(ctx, writer, newUrls, jwt, chunkView.CipherKey, chunkView.IsGzipped, chunkView.IsFullChunk(), chunkView.OffsetInChunk, int(chunkView.ViewSize))
|
|
if refetchErr == nil {
|
|
// Update the map so subsequent references use fresh URLs
|
|
fileId2Url[chunkView.FileId] = newUrls
|
|
}
|
|
return refetchErr
|
|
})
|
|
}
|
|
|
|
offset += int64(chunkView.ViewSize)
|
|
remaining -= int64(chunkView.ViewSize)
|
|
stats.FilerRequestHistogram.WithLabelValues("chunkDownload").Observe(time.Since(start).Seconds())
|
|
if err != nil {
|
|
stats.FilerHandlerCounter.WithLabelValues("chunkDownloadError").Inc()
|
|
return fmt.Errorf("read chunk: %w", err)
|
|
}
|
|
stats.FilerHandlerCounter.WithLabelValues("chunkDownload").Inc()
|
|
downloadThrottler.MaybeSlowdown(int64(chunkView.ViewSize))
|
|
}
|
|
if remaining > 0 {
|
|
glog.V(4).InfofCtx(ctx, "zero [%d,%d)", offset, offset+remaining)
|
|
err := writeZero(writer, remaining)
|
|
if err != nil {
|
|
return fmt.Errorf("write zero [%d,%d)", offset, offset+remaining)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}, nil
|
|
}
|
|
|
|
// PrepareStreamContentWithPrefetch is like PrepareStreamContentWithThrottler but uses
|
|
// concurrent chunk prefetching to overlap network I/O. When prefetchAhead > 1, fetch
|
|
// goroutines establish HTTP connections to volume servers ahead of time, streaming data
|
|
// through io.Pipe with minimal memory overhead.
|
|
//
|
|
// prefetchAhead controls the number of chunks fetched concurrently:
|
|
// - 0 or 1: falls back to sequential fetching (same as PrepareStreamContentWithThrottler)
|
|
// - 2+: uses pipe-based prefetch pipeline with that many concurrent fetches
|
|
func PrepareStreamContentWithPrefetch(ctx context.Context, masterClient wdclient.HasLookupFileIdFunction, jwtFunc VolumeServerJwtFunction, chunks []*filer_pb.FileChunk, offset int64, size int64, downloadMaxBytesPs int64, prefetchAhead int) (DoStreamContent, error) {
|
|
if prefetchAhead <= 1 {
|
|
return PrepareStreamContentWithThrottler(ctx, masterClient, jwtFunc, chunks, offset, size, downloadMaxBytesPs)
|
|
}
|
|
|
|
glog.V(4).InfofCtx(ctx, "prepare to stream content with prefetch=%d for chunks: %d", prefetchAhead, len(chunks))
|
|
chunkViews := ViewFromChunks(ctx, masterClient.GetLookupFileIdFunction(), chunks, offset, size)
|
|
|
|
fileId2Url := make(map[string][]string)
|
|
|
|
for x := chunkViews.Front(); x != nil; x = x.Next {
|
|
chunkView := x.Value
|
|
var urlStrings []string
|
|
var err error
|
|
for _, backoff := range getLookupFileIdBackoffSchedule {
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
urlStrings, err = masterClient.GetLookupFileIdFunction()(ctx, chunkView.FileId)
|
|
if err == nil && len(urlStrings) > 0 {
|
|
break
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
glog.V(4).InfofCtx(ctx, "waiting for chunk: %s", chunkView.FileId)
|
|
timer := time.NewTimer(backoff)
|
|
select {
|
|
case <-ctx.Done():
|
|
if !timer.Stop() {
|
|
<-timer.C
|
|
}
|
|
return nil, ctx.Err()
|
|
case <-timer.C:
|
|
}
|
|
}
|
|
if err != nil {
|
|
glog.V(1).InfofCtx(ctx, "operation LookupFileId %s failed, err: %v", chunkView.FileId, err)
|
|
return nil, err
|
|
} else if len(urlStrings) == 0 {
|
|
errUrlNotFound := fmt.Errorf("operation LookupFileId %s failed, err: urls not found", chunkView.FileId)
|
|
glog.ErrorCtx(ctx, errUrlNotFound)
|
|
return nil, errUrlNotFound
|
|
}
|
|
fileId2Url[chunkView.FileId] = urlStrings
|
|
}
|
|
|
|
return func(writer io.Writer) error {
|
|
return streamChunksPrefetched(ctx, writer, chunkViews, fileId2Url, jwtFunc, masterClient, offset, size, downloadMaxBytesPs, prefetchAhead)
|
|
}, nil
|
|
}
|
|
|
|
func StreamContent(masterClient wdclient.HasLookupFileIdFunction, writer io.Writer, chunks []*filer_pb.FileChunk, offset int64, size int64) error {
|
|
streamFn, err := PrepareStreamContent(masterClient, JwtForVolumeServer, chunks, offset, size)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return streamFn(writer)
|
|
}
|
|
|
|
// ---------------- ReadAllReader ----------------------------------
|
|
|
|
func writeZero(w io.Writer, size int64) (err error) {
|
|
zeroPadding := make([]byte, 1024)
|
|
var written int
|
|
for size > 0 {
|
|
if size > 1024 {
|
|
written, err = w.Write(zeroPadding)
|
|
} else {
|
|
written, err = w.Write(zeroPadding[:size])
|
|
}
|
|
size -= int64(written)
|
|
if err != nil {
|
|
return
|
|
}
|
|
}
|
|
return
|
|
}
|
|
|
|
// ---------------- ChunkStreamReader ----------------------------------
|
|
type ChunkStreamReader struct {
|
|
head *Interval[*ChunkView]
|
|
chunkView *Interval[*ChunkView]
|
|
totalSize int64
|
|
logicOffset int64
|
|
buffer []byte
|
|
bufferOffset int64
|
|
bufferLock sync.Mutex
|
|
chunk string
|
|
lookupFileId wdclient.LookupFileIdFunctionType
|
|
}
|
|
|
|
var _ = io.ReadSeeker(&ChunkStreamReader{})
|
|
var _ = io.ReaderAt(&ChunkStreamReader{})
|
|
var _ = io.Closer(&ChunkStreamReader{})
|
|
|
|
func doNewChunkStreamReader(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunks []*filer_pb.FileChunk) *ChunkStreamReader {
|
|
|
|
chunkViews := ViewFromChunks(ctx, lookupFileIdFn, chunks, 0, math.MaxInt64)
|
|
|
|
var totalSize int64
|
|
for x := chunkViews.Front(); x != nil; x = x.Next {
|
|
chunk := x.Value
|
|
totalSize += int64(chunk.ViewSize)
|
|
}
|
|
|
|
return &ChunkStreamReader{
|
|
head: chunkViews.Front(),
|
|
chunkView: chunkViews.Front(),
|
|
lookupFileId: lookupFileIdFn,
|
|
totalSize: totalSize,
|
|
}
|
|
}
|
|
|
|
func NewChunkStreamReaderFromFiler(ctx context.Context, masterClient *wdclient.MasterClient, chunks []*filer_pb.FileChunk) *ChunkStreamReader {
|
|
|
|
lookupFileIdFn := func(ctx context.Context, fileId string) (targetUrl []string, err error) {
|
|
return masterClient.LookupFileId(ctx, fileId)
|
|
}
|
|
|
|
return doNewChunkStreamReader(ctx, lookupFileIdFn, chunks)
|
|
}
|
|
|
|
// NewChunkStreamReaderFromLookup creates a ChunkStreamReader from a lookup function.
|
|
// Used by clients that already have a LookupFileIdFunctionType (e.g., from FilerSource).
|
|
func NewChunkStreamReaderFromLookup(ctx context.Context, lookupFn wdclient.LookupFileIdFunctionType, chunks []*filer_pb.FileChunk) *ChunkStreamReader {
|
|
return doNewChunkStreamReader(ctx, lookupFn, chunks)
|
|
}
|
|
|
|
func NewChunkStreamReader(filerClient filer_pb.FilerClient, chunks []*filer_pb.FileChunk) *ChunkStreamReader {
|
|
|
|
lookupFileIdFn := LookupFn(filerClient)
|
|
|
|
return doNewChunkStreamReader(context.Background(), lookupFileIdFn, chunks)
|
|
}
|
|
|
|
func (c *ChunkStreamReader) ReadAt(p []byte, off int64) (n int, err error) {
|
|
c.bufferLock.Lock()
|
|
defer c.bufferLock.Unlock()
|
|
if err = c.prepareBufferFor(off); err != nil {
|
|
return
|
|
}
|
|
c.logicOffset = off
|
|
return c.doRead(p)
|
|
}
|
|
|
|
func (c *ChunkStreamReader) Read(p []byte) (n int, err error) {
|
|
c.bufferLock.Lock()
|
|
defer c.bufferLock.Unlock()
|
|
return c.doRead(p)
|
|
}
|
|
|
|
func (c *ChunkStreamReader) doRead(p []byte) (n int, err error) {
|
|
// fmt.Printf("do read [%d,%d) at %s[%d,%d)\n", c.logicOffset, c.logicOffset+int64(len(p)), c.chunk, c.bufferOffset, c.bufferOffset+int64(len(c.buffer)))
|
|
for n < len(p) {
|
|
// println("read", c.logicOffset)
|
|
if err = c.prepareBufferFor(c.logicOffset); err != nil {
|
|
return
|
|
}
|
|
t := copy(p[n:], c.buffer[c.logicOffset-c.bufferOffset:])
|
|
n += t
|
|
c.logicOffset += int64(t)
|
|
}
|
|
return
|
|
}
|
|
|
|
func (c *ChunkStreamReader) isBufferEmpty() bool {
|
|
return len(c.buffer) <= int(c.logicOffset-c.bufferOffset)
|
|
}
|
|
|
|
func (c *ChunkStreamReader) Seek(offset int64, whence int) (int64, error) {
|
|
c.bufferLock.Lock()
|
|
defer c.bufferLock.Unlock()
|
|
|
|
var err error
|
|
switch whence {
|
|
case io.SeekStart:
|
|
case io.SeekCurrent:
|
|
offset += c.logicOffset
|
|
case io.SeekEnd:
|
|
offset = c.totalSize + offset
|
|
}
|
|
if offset > c.totalSize {
|
|
err = io.ErrUnexpectedEOF
|
|
} else {
|
|
c.logicOffset = offset
|
|
}
|
|
|
|
return offset, err
|
|
|
|
}
|
|
|
|
func insideChunk(offset int64, chunk *ChunkView) bool {
|
|
return chunk.ViewOffset <= offset && offset < chunk.ViewOffset+int64(chunk.ViewSize)
|
|
}
|
|
|
|
func (c *ChunkStreamReader) prepareBufferFor(offset int64) (err error) {
|
|
// stay in the same chunk
|
|
if c.bufferOffset <= offset && offset < c.bufferOffset+int64(len(c.buffer)) {
|
|
return nil
|
|
}
|
|
// glog.V(2).Infof("c.chunkView: %v buffer:[%d,%d) offset:%d totalSize:%d", c.chunkView, c.bufferOffset, c.bufferOffset+int64(len(c.buffer)), offset, c.totalSize)
|
|
|
|
// find a possible chunk view
|
|
p := c.chunkView
|
|
for p != nil {
|
|
chunk := p.Value
|
|
// glog.V(2).Infof("prepareBufferFor check chunk:[%d,%d)", chunk.ViewOffset, chunk.ViewOffset+int64(chunk.ViewSize))
|
|
if insideChunk(offset, chunk) {
|
|
if c.isBufferEmpty() || c.bufferOffset != chunk.ViewOffset {
|
|
c.chunkView = p
|
|
return c.fetchChunkToBuffer(chunk)
|
|
}
|
|
}
|
|
if offset < c.bufferOffset {
|
|
p = p.Prev
|
|
} else {
|
|
p = p.Next
|
|
}
|
|
}
|
|
|
|
return io.EOF
|
|
}
|
|
|
|
func (c *ChunkStreamReader) fetchChunkToBuffer(chunkView *ChunkView) error {
|
|
urlStrings, err := c.lookupFileId(context.Background(), chunkView.FileId)
|
|
if err != nil {
|
|
glog.V(1).Infof("operation LookupFileId %s failed, err: %v", chunkView.FileId, err)
|
|
return err
|
|
}
|
|
var buffer bytes.Buffer
|
|
// pre-size to the known chunk size; avoids bytes.Buffer's doubling regrowth
|
|
buffer.Grow(int(chunkView.ViewSize))
|
|
var shouldRetry bool
|
|
jwt := JwtForVolumeServer(chunkView.FileId)
|
|
for _, urlString := range urlStrings {
|
|
shouldRetry, err = util_http.ReadUrlAsStream(context.Background(), util_http.AppendQueryParameter(urlString, "readDeleted", "true"), jwt, chunkView.CipherKey, chunkView.IsGzipped, chunkView.IsFullChunk(), chunkView.OffsetInChunk, int(chunkView.ViewSize), func(data []byte) {
|
|
buffer.Write(data)
|
|
})
|
|
if !shouldRetry {
|
|
break
|
|
}
|
|
if err != nil {
|
|
glog.V(1).Infof("read %s failed, err: %v", chunkView.FileId, err)
|
|
buffer.Reset()
|
|
} else {
|
|
break
|
|
}
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
c.buffer = buffer.Bytes()
|
|
c.bufferOffset = chunkView.ViewOffset
|
|
c.chunk = chunkView.FileId
|
|
|
|
// glog.V(0).Infof("fetched %s [%d,%d)", chunkView.FileId, chunkView.ViewOffset, chunkView.ViewOffset+int64(chunkView.ViewSize))
|
|
|
|
return nil
|
|
}
|
|
|
|
func (c *ChunkStreamReader) Close() error {
|
|
c.bufferLock.Lock()
|
|
defer c.bufferLock.Unlock()
|
|
c.buffer = nil
|
|
c.head = nil
|
|
c.chunkView = nil
|
|
return nil
|
|
}
|
|
|
|
func VolumeId(fileId string) string {
|
|
lastCommaIndex := strings.LastIndex(fileId, ",")
|
|
if lastCommaIndex > 0 {
|
|
return fileId[:lastCommaIndex]
|
|
}
|
|
return fileId
|
|
}
|