/
githubmirror
/
nbs
Обзор
Документация
Войти
/
githubmirror
/
nbs
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
cloud/disk_manager/internal/pkg/dataplane/url/common/reader.go
241 строка
4 KB
Sergei
issue-557: [Disk Manager] Add metrics for url source (#6444)
15 июл 2026, 17:37
Не верифицирован
15 июл 2026, 17:37
77c4ac1
Код
Авторство
О чём код?
package common import ( "bytes" "context" "encoding/binary" "io" net_url "net/url" "time" "github.com/ydb-platform/nbs/cloud/disk_manager/internal/pkg/dataplane/url/common/cache" url_metrics "github.com/ydb-platform/nbs/cloud/disk_manager/internal/pkg/dataplane/url/metrics" "github.com/ydb-platform/nbs/cloud/tasks/errors" ) //////////////////////////////////////////////////////////////////////////////// type Reader interface { EnableCache() // Returns total amount of bytes that can be read. Size() uint64 ETag() string // Returns the number of bytes read. Read(ctx context.Context, start uint64, data []byte) (uint64, error) ReadBinary( ctx context.Context, start uint64, size uint64, byteOrder binary.ByteOrder, data interface{}, ) error CacheMissedRequestsCount() uint64 } //////////////////////////////////////////////////////////////////////////////// type urlReader struct { httpClient httpClientInterface url string etag string size uint64 cache *cache.Cache metrics url_metrics.Metrics } func NewURLReader( ctx context.Context, httpClientTimeout time.Duration, httpClientMinRetryTimeout time.Duration, httpClientMaxRetryTimeout time.Duration, httpClientMaxRetries uint32, url string, metrics url_metrics.Metrics, ) (_ Reader, err error) { parsed, err := net_url.Parse(url) if err != nil { return nil, NewSourceInvalidError("parse url: %w", err) } if parsed.Scheme != "https" && parsed.Scheme != "http" { return nil, NewSourceInvalidError( "invalid protocol scheme %q", parsed.Scheme, ) } httpClient := newHTTPClient( ctx, httpClientTimeout, httpClientMinRetryTimeout, httpClientMaxRetryTimeout, httpClientMaxRetries, url, metrics, ) defer metrics.StatRequest("head")(&err) resp, err := httpClient.Head(ctx) if err != nil { return nil, err } etag := resp.Header.Get("Etag") if len(etag) == 0 { return nil, NewSourceInvalidError( "missing Etag header in response", ) } if resp.ContentLength == 0 { return nil, NewSourceInvalidError( "url ContentLength should not be zero", ) } return &urlReader{ httpClient: httpClient, url: url, etag: etag, size: uint64(resp.ContentLength), metrics: metrics, }, nil } //////////////////////////////////////////////////////////////////////////////// func (r *urlReader) EnableCache() { if r.cache == nil { r.cache = cache.NewCache(r.read, r.metrics.OnCacheHit) } } func (r *urlReader) Size() uint64 { return r.size } func (r *urlReader) ETag() string { return r.etag } func (r *urlReader) CacheMissedRequestsCount() uint64 { return r.httpClient.RequestsCount() } func (r *urlReader) validateRange( start, end uint64, // Half-open interval [start:end). ) error { if start >= end || end > r.size { return NewSourceInvalidError("range [%v:%v) is invalid", start, end) } return nil } func (r *urlReader) read( ctx context.Context, start uint64, data []byte, ) (err error) { defer r.metrics.StatRequest("get")(&err) end := start + uint64(len(data)) if end > r.size { end = r.size if r.size <= start { return NewSourceInvalidError( "size %v should be greater than start %v", r.size, start, ) } // Cut the tail. data = data[:r.size-start] } err = r.validateRange(start, end) if err != nil { return err } reader, err := r.httpClient.Body(ctx, start, end, r.etag) if err != nil { return err } defer reader.Close() _, err = io.ReadFull(reader, data) if err != nil { // NBS-3324: interpret all errors as retriable. return errors.NewRetriableError(err) } return nil } func (r *urlReader) Read( ctx context.Context, start uint64, data []byte, ) (uint64, error) { size := uint64(len(data)) if r.cache == nil { err := r.read(ctx, start, data) if err != nil { return 0, err } return size, nil } err := r.cache.Read(ctx, start, data) if err != nil { return 0, err } return size, nil } func (r *urlReader) ReadBinary( ctx context.Context, start uint64, size uint64, byteOrder binary.ByteOrder, data interface{}, ) error { end := start + size err := r.validateRange(start, end) if err != nil { return err } byteData := make([]byte, size) _, err = r.Read(ctx, start, byteData) if err != nil { return err } err = binary.Read(bytes.NewReader(byteData), byteOrder, data) // NBS-3324: interpret all errors as retriable. if err != nil { return errors.NewRetriableError(err) } return nil }