@@ -3,26 +3,31 @@ package libindex
33import (
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