Skip to content

Commit d472855

Browse files
committed
libindex: wire in httpreader.Reader
This handles the "easy" case of simply proxying reads for uncompressed tar archives to range requests. Future improvements would move the "spooling" out of this package and into the `fs.FS` implementation. Signed-off-by: Hank Donnay <hdonnay@redhat.com> Change-Id: I9d200dd841954054df0b187b9fd160a56a6a6964
1 parent ca9be1d commit d472855

3 files changed

Lines changed: 357 additions & 133 deletions

File tree

internal/httpreader/reader.go

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -88,8 +88,14 @@ var (
8888

8989
// ReadAt implements [io.ReaderAt].
9090
func (r *Reader) ReadAt(b []byte, off int64) (int, error) {
91-
if len(b) == 0 {
91+
switch {
92+
case len(b) == 0:
93+
// Nothing to read.
9294
return 0, nil
95+
case off >= r.size:
96+
// Impossible to read anything, so copy [os.File] behavior.
97+
// Use ≥ because it's comparing an index and a length.
98+
return 0, io.EOF
9399
}
94100

95101
req, err := http.NewRequestWithContext(r.ctx, http.MethodGet, r.res, nil)
@@ -105,6 +111,8 @@ func (r *Reader) ReadAt(b []byte, off int64) (int, error) {
105111
switch res.StatusCode {
106112
case http.StatusPartialContent: // OK
107113
case http.StatusOK:
114+
// NOTE(hank) I've seen servers in the wild that just send the file if
115+
// the range starts at 0, even if they support range requests.
108116
if off != 0 {
109117
return 0, errWithStatus(srvBotch(req.Header.Get(`range`)), res.StatusCode)
110118
}

libindex/fetcher.go

Lines changed: 163 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3,26 +3,31 @@ package libindex
33
import (
44
"bufio"
55
"bytes"
6+
"cmp"
67
"context"
78
"encoding/hex"
89
"errors"
910
"fmt"
1011
"io"
1112
"log/slog"
13+
"mime"
1214
"net/http"
1315
"net/url"
1416
"os"
17+
"regexp"
1518
"runtime"
1619
"strings"
1720

1821
"go.opentelemetry.io/otel/attribute"
1922
"go.opentelemetry.io/otel/codes"
23+
semconv "go.opentelemetry.io/otel/semconv/v1.41.0"
2024
"go.opentelemetry.io/otel/trace"
2125
"golang.org/x/sync/errgroup"
2226

2327
"github.com/quay/claircore"
2428
"github.com/quay/claircore/indexer"
2529
"github.com/quay/claircore/internal/cache"
30+
"github.com/quay/claircore/internal/httpreader"
2631
"github.com/quay/claircore/internal/httputil"
2732
"github.com/quay/claircore/internal/wart"
2833
"github.com/quay/claircore/internal/zreader"
@@ -125,6 +130,43 @@ func (a *RemoteFetchArena) fetchInto(ctx context.Context, l *claircore.Layer, cl
125130
span.End()
126131
}()
127132

133+
var d details
134+
d, err = a.inspect(ctx, desc)
135+
if err != nil {
136+
return err
137+
}
138+
139+
HTTPReader:
140+
switch { // Switch for a single case to be able to break.
141+
case d.Uncompressed && d.RangeOK:
142+
var opts []httpreader.Option
143+
if d.ContentLength > 0 {
144+
opts = append(opts, httpreader.WithSize(d.ContentLength))
145+
}
146+
if len(desc.Headers) != 0 {
147+
opts = append(opts, httpreader.WithHeaders(desc.Headers))
148+
}
149+
var rd *httpreader.Reader
150+
rd, err = httpreader.New(ctx, a.wc, desc.URI, opts...)
151+
switch {
152+
case err == nil:
153+
if err = l.Init(ctx, desc, rd); err != nil {
154+
return errors.Join(err, rd.Close())
155+
}
156+
*cl = closeFunc(func() (err error) {
157+
err = errors.Join(l.Close(), rd.Close())
158+
return err
159+
})
160+
a.logger(desc).DebugContext(ctx, "using httpreader")
161+
return nil
162+
case errors.Is(err, errors.ErrUnsupported):
163+
err = nil
164+
break HTTPReader
165+
default:
166+
return err
167+
}
168+
}
169+
128170
// NB This is not closed on purpose. The [io.Closer] populated by this
129171
// function holds the pointer until that function is cleaned up. Once
130172
// nothing has a copy of this [*os.File], the runtime will run all the
@@ -135,19 +177,20 @@ func (a *RemoteFetchArena) fetchInto(ctx context.Context, l *claircore.Layer, cl
135177
var spool *os.File
136178
spool, err = a.files.Get(ctx, key, func(ctx context.Context, _ string) (*os.File, error) {
137179
cacheHit = false
138-
return a.fetchFileForCache(ctx, desc)
180+
return a.fetchFileForCache(ctx, desc, d)
139181
})
140182
if err != nil {
141183
return err
142184
}
185+
var f *os.File
143186
// This is an owned, independent descriptor for the passed [*os.File].
144-
f, err := reopen(a.root, spool)
187+
f, err = reopen(a.root, spool)
145188
if err != nil {
146189
return err
147190
}
148191

149192
// If this succeeds, "f" is now owned by "l"
150-
if err := l.Init(ctx, desc, f); err != nil {
193+
if err = l.Init(ctx, desc, f); err != nil {
151194
return errors.Join(err, f.Close())
152195
}
153196
*cl = closeFunc(func() (err error) {
@@ -173,12 +216,110 @@ func (f closeFunc) Close() error {
173216
return f()
174217
}
175218

219+
// Inspect makes a request to the layer's URI and examines the response for
220+
// useful information.
221+
func (a *RemoteFetchArena) inspect(ctx context.Context, desc *claircore.LayerDescription) (d details, err error) {
222+
ctx, span := tracer.Start(ctx, "RemoteFetchArena.inspect")
223+
defer func() {
224+
a.logger(desc).DebugContext(ctx, "inspected resource", "ok", err == nil, "details", &d)
225+
span.RecordError(err)
226+
span.End()
227+
}()
228+
span.SetStatus(codes.Error, "")
229+
230+
var req *http.Request
231+
var res *http.Response
232+
req, err = http.NewRequestWithContext(ctx, http.MethodGet, desc.URI, nil)
233+
if err != nil {
234+
return d, fmt.Errorf("fetcher: failed to construct request: %w", err)
235+
}
236+
req.Header = http.Header(desc.Headers).Clone()
237+
if req.Header == nil {
238+
req.Header = make(http.Header)
239+
}
240+
req.Header.Set(`claircore-reason`, `inspect`)
241+
req.Header.Set(`range`, `bytes=0-15`)
242+
res, err = a.wc.Do(req)
243+
if err != nil {
244+
return d, fmt.Errorf("fetcher: request failed: %w", err)
245+
}
246+
err = httputil.CheckResponse(res, http.StatusOK, http.StatusPartialContent)
247+
if err != nil {
248+
return d, fmt.Errorf("fetcher: %w", err)
249+
}
250+
head := make([]byte, 16)
251+
_, err = io.ReadFull(res.Body, head)
252+
_ = res.Body.Close()
253+
if err != nil {
254+
return d, fmt.Errorf("fetcher: unexpected read: %w", err)
255+
}
256+
257+
const (
258+
ctKey = `content-type`
259+
arKey = `accept-ranges`
260+
)
261+
ctVal := cmp.Or(res.Header.Get(ctKey), `application/octet-stream`)
262+
ct, _, err := mime.ParseMediaType(ctVal)
263+
if err != nil {
264+
return d, fmt.Errorf("fetcher: %w", err)
265+
}
266+
span.SetAttributes(
267+
semconv.HTTPResponseStatusCode(res.StatusCode),
268+
semconv.HTTPResponseBodySize(int(res.ContentLength)),
269+
semconv.HTTPResponseHeader(ctKey, ct),
270+
semconv.HTTPResponseHeader(arKey, res.Header.Get(arKey)),
271+
)
272+
if isOctetStream(ct) && zreader.DetectCompression(head) == zreader.KindNone {
273+
ct = `application/x-tar`
274+
}
275+
switch res.StatusCode {
276+
case http.StatusOK:
277+
d.ContentLength = res.ContentLength
278+
case http.StatusPartialContent:
279+
var cr httpreader.ContentRange
280+
if err := cr.Parse(res.Header.Get(`content-range`)); err == nil {
281+
d.ContentLength = cr.Length
282+
}
283+
}
284+
d.Uncompressed = ct == "application/x-tar" || strings.HasSuffix(ct, ".tar")
285+
d.RangeOK = res.Header.Get(arKey) == `bytes` || res.StatusCode == http.StatusPartialContent
286+
287+
span.SetStatus(codes.Ok, "")
288+
return d, nil
289+
}
290+
291+
var _ slog.LogValuer = (*details)(nil)
292+
293+
// Details is the useful information reported by [RemoteFetchArena.inspect].
294+
type details struct {
295+
// The content length, if known.
296+
ContentLength int64
297+
// The resource reports to be known non-compressed contents.
298+
Uncompressed bool
299+
// The server reports supporting "bytes" via the "Accept-Ranges" header.
300+
RangeOK bool
301+
}
302+
303+
// LogValue implements [slog.LogValuer].
304+
func (d *details) LogValue() slog.Value {
305+
return slog.GroupValue(
306+
slog.Int64("content-length", d.ContentLength),
307+
slog.Bool("uncompressed", d.Uncompressed),
308+
slog.Bool("range-ok", d.RangeOK),
309+
)
310+
}
311+
312+
// Logger returns a [*slog.Logger] with predefined attributes.
313+
func (a *RemoteFetchArena) logger(desc *claircore.LayerDescription) *slog.Logger {
314+
return slog.With("arena", a.root.Name(), "layer", desc.Digest, "uri", desc.URI)
315+
}
316+
176317
// FetchFileForCache is the inner function used inside the [cache.Live].
177318
//
178319
// Because we know we're the only concurrent call that's dealing with this blob,
179320
// we can be a bit more lax.
180-
func (a *RemoteFetchArena) fetchFileForCache(ctx context.Context, desc *claircore.LayerDescription) (*os.File, error) {
181-
log := slog.With("arena", a.root.Name(), "layer", desc.Digest, "uri", desc.URI)
321+
func (a *RemoteFetchArena) fetchFileForCache(ctx context.Context, desc *claircore.LayerDescription, _ details) (*os.File, error) {
322+
log := a.logger(desc)
182323
ctx, span := tracer.Start(ctx, "RemoteFetchArena.fetchFileForCache")
183324
defer span.End()
184325
span.SetStatus(codes.Error, "")
@@ -210,6 +351,7 @@ func (a *RemoteFetchArena) fetchFileForCache(ctx context.Context, desc *claircor
210351
URL: url,
211352
Header: http.Header(desc.Headers).Clone(),
212353
}).WithContext(ctx)
354+
req.Header.Set(`claircore-reason`, `fetch`)
213355
resp, err := a.wc.Do(req)
214356
if err != nil {
215357
return nil, fmt.Errorf("fetcher: request failed: %w", err)
@@ -235,10 +377,13 @@ func (a *RemoteFetchArena) fetchFileForCache(ctx context.Context, desc *claircor
235377
}
236378
defer zr.Close()
237379
// Look at the content-type and optionally fix it up.
238-
ct := resp.Header.Get("content-type")
380+
ct, _, err := mime.ParseMediaType(resp.Header.Get("content-type"))
381+
if err != nil {
382+
return nil, fmt.Errorf("fetcher: %w", err)
383+
}
239384
log.DebugContext(ctx, "reported content-type", "content-type", ct)
240385
span.SetAttributes(payloadType(ct), payloadCompression(kind))
241-
if ct == "" || ct == "text/plain" || ct == "binary/octet-stream" || ct == "application/octet-stream" {
386+
if ct == "" || ct == "text/plain" || isOctetStream(ct) {
242387
switch kind {
243388
case zreader.KindGzip:
244389
ct = "application/gzip"
@@ -420,3 +565,14 @@ func (p *FetchProxy) Close() error {
420565
}
421566
return errors.Join(errs...)
422567
}
568+
569+
// IsOctetStream tests if the incoming media type string is an "octet-stream"
570+
// type.
571+
//
572+
// Reports false for strings that are not of the form "<type>/<subtype>". For
573+
// interoperability purposes, the nonstandard "binary" type is accepted.
574+
func isOctetStream(ct string) bool {
575+
return reOctetStream.MatchString(ct)
576+
}
577+
578+
var reOctetStream = regexp.MustCompile(`^(application|binary)/octet-stream$`)

0 commit comments

Comments
 (0)