Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions internal/dataplane/http_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -870,7 +870,11 @@ func TestSourceMD5_SizeDerivedDeadlineOutlivesOrdinaryHeaderTimeout(t *testing.T
const (
ordinaryTimeout = 30 * time.Millisecond
hashDelay = 90 * time.Millisecond
wantMD5 = "8c7dd922ad47494fc02c388e12c00eac"
// hashTimeout only has to outlive hashDelay. The wide margin absorbs
// the multi-second stalls a loaded machine shows while the whole
// module's tests run in parallel.
hashTimeout = 10 * time.Second
wantMD5 = "8c7dd922ad47494fc02c388e12c00eac"
)

srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
Expand All @@ -884,13 +888,13 @@ func TestSourceMD5_SizeDerivedDeadlineOutlivesOrdinaryHeaderTimeout(t *testing.T
ordinaryTransport.ResponseHeaderTimeout = ordinaryTimeout

hashTransport := srv.Client().Transport.(*http.Transport).Clone()
hashTransport.ResponseHeaderTimeout = time.Second
hashTransport.ResponseHeaderTimeout = hashTimeout

f := NewFetcher(
&http.Client{Transport: ordinaryTransport},
WithSourceHashDoer(&http.Client{Transport: hashTransport}),
)
f.sourceHashTimeout = func(int64) time.Duration { return time.Second }
f.sourceHashTimeout = func(int64) time.Duration { return hashTimeout }

body, ordinaryErr := f.GetFile(context.Background(), srv.URL)
if ordinaryErr == nil {
Expand Down
7 changes: 5 additions & 2 deletions internal/lockfile/lockfile.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ func Acquire(path string, staleAfter time.Duration, onReclaim func(age time.Dura
// reclaim removes the orphaned lock described by stale.
//
// - Renames path to a unique scratch path; exactly one concurrent reclaimer wins.
// - Checks file identity: if the moved file is not stale, a concurrent acquirer
// - Checks file identity (inode and mtime): if the moved file is not stale, a concurrent acquirer
// created a fresh lock after our staleness check - restore it and return ErrLocked.
// - Otherwise removes the scratch file and optionally calls onReclaim.
func reclaim(path string, stale os.FileInfo, onReclaim func(age time.Duration)) error {
Expand All @@ -84,7 +84,10 @@ func reclaim(path string, stale os.FileInfo, onReclaim func(age time.Duration))

moved, err := os.Stat(reclaimScratch)
// Rename may have moved a fresh lock created after our Stat in Acquire.
grabbedFresh := err == nil && !os.SameFile(stale, moved)
// The inode alone cannot tell: ext4 hands a freed inode number straight to
// the next file, so the fresh lock may reuse the stale one's. Rename keeps
// mtime, and a re-created lock cannot carry the stale lock's old mtime.
grabbedFresh := err == nil && (!os.SameFile(stale, moved) || !moved.ModTime().Equal(stale.ModTime()))

if grabbedFresh {
if err := os.Rename(reclaimScratch, path); err != nil {
Expand Down
8 changes: 4 additions & 4 deletions internal/mirror/dist/cli.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,15 +114,15 @@ func (svc *Service) PullCLI(ctx context.Context) error {
// newest published stable one.
//
// A registry that simply does not offer the binary is not a failure: it yields
// an empty tag and a short skipReason for the summary, the way an absent
// an empty tag and a short skip reason for the summary, the way an absent
// platform plugin yields a warning. The registry's own error text goes to the
// debug log rather than the summary - a wrapped NAME_UNKNOWN chain names the
// repository the label already names, and buries the one fact the operator
// needs behind four levels of transport detail.
//
// err is reserved for a pinned version that cannot be resolved: the user asked
// for it by name, so it stops the pull.
func (svc *Service) resolveCLITag(ctx context.Context) (tag versionTag, skipReason string, err error) {
// The error is reserved for a pinned version that cannot be resolved: the user
// asked for it by name, so it stops the pull.
func (svc *Service) resolveCLITag(ctx context.Context) (versionTag, string, error) {
if pinned := svc.options.CLITag; pinned != "" {
// An explicit version is checked before the pull so a typo fails
// immediately instead of after the transfer retries are exhausted.
Expand Down
5 changes: 4 additions & 1 deletion internal/snapshot/aggapi/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1027,7 +1027,10 @@ func TestAggregatedAPICalls_RejectOversizedChunkedResponses(t *testing.T) {

for _, tc := range aggregatedAPICallCases() {
t.Run(tc.name, func(t *testing.T) {
client := newBoundedTestClient(t, handler, 5*time.Second, maxResponseBytes)
// The timeout only guards against a hang; the call answers at once,
// but a loaded machine running the whole module can stall it for
// seconds.
client := newBoundedTestClient(t, handler, 30*time.Second, maxResponseBytes)

err := tc.call(context.Background(), client)
if !errors.Is(err, ErrResponseTooLarge) {
Expand Down
2 changes: 1 addition & 1 deletion internal/snapshot/snapimport/import_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1409,7 +1409,7 @@ func TestLinuxMountedRunEscapeHelper(t *testing.T) {
root := buildTwoLevelArchive(t)
sourcePath, targetPath := matchingOutsideMountFixture(t, root, strings.HasSuffix(scenario, "regular-file"))

if err := bindMountForTest(sourcePath, targetPath); err != nil {
if err := bindMountForTest(t, sourcePath, targetPath); err != nil {
fmt.Printf("mount namespace unavailable: %v\n", err)

return
Expand Down
15 changes: 13 additions & 2 deletions internal/snapshot/snapimport/plan_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1722,7 +1722,7 @@ func TestLinuxMountedPlanEscapeHelper(t *testing.T) {
sourcePath, targetPath := matchingOutsideMountFixture(t, root, strings.HasSuffix(scenario, "regular-file"))

mount := func() error {
return bindMountForTest(sourcePath, targetPath)
return bindMountForTest(t, sourcePath, targetPath)
}

var err error
Expand Down Expand Up @@ -1809,12 +1809,23 @@ func matchingOutsideMountFixture(t *testing.T, root string, regularFile bool) (s
return source, target
}

func bindMountForTest(source, target string) error {
func bindMountForTest(t *testing.T, source, target string) error {
t.Helper()

output, err := exec.Command("mount", "--bind", source, target).CombinedOutput()
if err != nil {
return fmt.Errorf("mount --bind: %w: %s", err, output)
}

// The TempDir cleanup cannot unlink a mount point (EBUSY), so detach the
// bind mount first: cleanups run last-in first-out, and the TempDir one
// was registered before this mount.
t.Cleanup(func() {
if output, err := exec.Command("umount", target).CombinedOutput(); err != nil {
t.Errorf("umount %s: %v: %s", target, err, output)
}
})

return nil
}

Expand Down
13 changes: 11 additions & 2 deletions internal/snapshot/transport/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -315,16 +315,24 @@ func (c *Client) NewPersistentHTTPClient() (*PersistentHTTPClient, error) {
var (
clonedHTTPTransport *http.Transport
enableHTTP2 bool
protocols *http.Protocols
)

if transport, ok := rt.(*http.Transport); ok {
cloned := transport.Clone()
enableHTTP2 = cloned.ForceAttemptHTTP2 || cloned.TLSNextProto["h2"] != nil
enableHTTP2 = cloned.ForceAttemptHTTP2 || cloned.TLSNextProto["h2"] != nil ||
(cloned.Protocols != nil && cloned.Protocols.HTTP2())
protocols = cloned.Protocols
// Clone copies client-go's x/net/http2 TLSNextProto closures. An
// empty map prevents intermediate caller wrappers from enabling or
// copying HTTP/2 while they clone this transport.
// copying HTTP/2 while they clone this transport. Since Go 1.27
// x/net/http2 enables HTTP/2 through Protocols instead, and each
// clone that keeps it enabled gets a TLSNextProto["h2"] stub that
// the next clone copies but cannot serve h2 with, so Protocols is
// dropped here as well and restored on the final transport.
cloned.TLSNextProto = make(map[string]func(string, *tls.Conn) http.RoundTripper)
cloned.ForceAttemptHTTP2 = false
cloned.Protocols = nil
clonedHTTPTransport = cloned
rt = cloned
}
Expand All @@ -343,6 +351,7 @@ func (c *Client) NewPersistentHTTPClient() (*PersistentHTTPClient, error) {
if ownedHTTPTransport != nil && enableHTTP2 {
ownedHTTPTransport.TLSNextProto = nil
ownedHTTPTransport.ForceAttemptHTTP2 = true
ownedHTTPTransport.Protocols = protocols
utilnet.SetTransportDefaults(ownedHTTPTransport)
}

Expand Down
20 changes: 15 additions & 5 deletions internal/snapshot/transport/http_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -673,6 +673,16 @@ func TestPersistentHTTPClient_IsolatesHTTP2PoolsAndCleanup(t *testing.T) {
func TestPersistentHTTPClient_IsolatesHTTP2ResponseHeaderTimeouts(t *testing.T) {
t.Parallel()

// The short request must end on its own timeout, long before the long
// timeout could end it, so leaking either timeout into the other client
// fails the test. The two stay far apart so that a stall on a loaded
// machine cannot push the short request past maxShortElapsed.
const (
shortTimeout = 50 * time.Millisecond
longTimeout = 10 * time.Second
maxShortElapsed = 5 * time.Second
)

shortStarted := make(chan struct{})
longStarted := make(chan struct{})
releaseLong := make(chan struct{}, 1)
Expand Down Expand Up @@ -701,13 +711,13 @@ func TestPersistentHTTPClient_IsolatesHTTP2ResponseHeaderTimeouts(t *testing.T)
})

sc := newSharedHTTP2Client(t, srv)
shortClient, err := newPersistentTestClient(sc, srv, 50*time.Millisecond)
shortClient, err := newPersistentTestClient(sc, srv, shortTimeout)
if err != nil {
t.Fatalf("build short-timeout client: %v", err)
}
t.Cleanup(shortClient.CloseIdleConnections)

longClient, err := newPersistentTestClient(sc, srv, time.Second)
longClient, err := newPersistentTestClient(sc, srv, longTimeout)
if err != nil {
t.Fatalf("build long-timeout client: %v", err)
}
Expand Down Expand Up @@ -737,12 +747,12 @@ func TestPersistentHTTPClient_IsolatesHTTP2ResponseHeaderTimeouts(t *testing.T)
t.Fatalf("short HTTP/2 request error = %v, want response-header timeout", err)
}

if elapsed < 25*time.Millisecond {
if elapsed < shortTimeout/2 {
t.Fatalf("short HTTP/2 response-header timeout took only %v", elapsed)
}

if elapsed > 500*time.Millisecond {
t.Fatalf("short HTTP/2 response-header timeout took %v, want under 500ms", elapsed)
if elapsed > maxShortElapsed {
t.Fatalf("short HTTP/2 response-header timeout took %v, want under %v", elapsed, maxShortElapsed)
}

<-shortStarted
Expand Down
Loading