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
21 changes: 20 additions & 1 deletion lib/images/credentials_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,11 +113,30 @@ func TestBorrowedCredentialsExpireWhileQueued(t *testing.T) {
return m.inflightPulls[digest].credentials == nil
}, time.Second, time.Millisecond)

credentials, _, expired := m.borrowedAuth(digest)
credentials, _, expired, stale := m.borrowedAuth(digest, inflight)
assert.True(t, expired)
assert.False(t, stale)
assert.Nil(t, credentials)
}

func TestBorrowedAuthRejectsReplacedInflightPull(t *testing.T) {
m := &manager{inflightPulls: make(map[string]*inflightImagePull)}
const digest = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
first := m.registerInflightPull(digest, &authn.AuthConfig{Username: "first"})
second := m.registerInflightPull(digest, &authn.AuthConfig{Username: "second"})
defer m.releaseInflightPull(digest, second)()

credentials, _, expired, stale := m.borrowedAuth(digest, first)
assert.Nil(t, credentials)
assert.False(t, expired)
assert.True(t, stale)

credentials, _, expired, stale = m.borrowedAuth(digest, second)
assert.Equal(t, "second", credentials.Username)
assert.False(t, expired)
assert.False(t, stale)
}

func TestRecoverInterruptedCredentialedConversionFromCache(t *testing.T) {
p := paths.New(t.TempDir())
img := createTestDockerImage(t)
Expand Down
90 changes: 52 additions & 38 deletions lib/images/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,15 @@ import (
"time"

"github.com/google/go-containerregistry/pkg/authn"
"github.com/google/uuid"
"github.com/kernel/hypeman/lib/paths"
"github.com/kernel/hypeman/lib/queue"
"github.com/kernel/hypeman/lib/tags"
"go.opentelemetry.io/otel/metric"
)

var errStaleBuild = errors.New("stale image build")

const (
StatusPending = "pending"
StatusPulling = "pulling"
Expand Down Expand Up @@ -297,6 +300,9 @@ func (m *manager) registerInflightPull(digest string, credentials *authn.AuthCon
if m.inflightPulls == nil {
m.inflightPulls = make(map[string]*inflightImagePull)
}
if previous := m.inflightPulls[digest]; previous != nil && previous.timer != nil {
previous.timer.Stop()
}
inflight := &inflightImagePull{
fingerprint: credentialFingerprint(credentials),
credentials: credentials,
Expand Down Expand Up @@ -334,17 +340,20 @@ func (m *manager) releaseInflightPull(digest string, inflight *inflightImagePull
}
}

func (m *manager) borrowedAuth(digest string) (*authn.AuthConfig, time.Time, bool) {
func (m *manager) borrowedAuth(digest string, expected *inflightImagePull) (*authn.AuthConfig, time.Time, bool, bool) {
m.createMu.Lock()
defer m.createMu.Unlock()
inflight := m.inflightPulls[digest]
if expected != nil && inflight != expected {
return nil, time.Time{}, false, true
}
if inflight == nil || inflight.credentialsExpireAt.IsZero() {
return nil, time.Time{}, false
return nil, time.Time{}, false, false
}
if inflight.credentials == nil || time.Now().After(inflight.credentialsExpireAt) {
return nil, inflight.credentialsExpireAt, true
return nil, inflight.credentialsExpireAt, true, false
}
return inflight.credentials, inflight.credentialsExpireAt, false
return inflight.credentials, inflight.credentialsExpireAt, false, false
}

func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest, requestedPlatform Platform) (*Image, error) {
Expand All @@ -365,6 +374,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest,
Status: StatusPending,
Request: &storedReq,
BorrowedAuth: req.Credentials != nil,
BuildID: uuid.New().String(),
Tags: tags.Clone(req.Tags),
CreatedAt: time.Now(),
}
Expand All @@ -377,10 +387,14 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest,
// Keep borrowed credentials outside the queued closure so their lifetime is
// bounded even when this job waits behind another pull.
inflight := m.registerInflightPull(ref.Digest(), req.Credentials)
queuePos := m.queue.Enqueue(ref.Digest(), func() {
credentials, deadline, expired := m.borrowedAuth(ref.Digest())
buildID := meta.BuildID
queuePos := m.queue.EnqueueSuccessor(ref.Digest(), func() {
credentials, deadline, expired, stale := m.borrowedAuth(ref.Digest(), inflight)
if stale {
return
}
if expired {
m.updateStatusByDigest(ref, StatusFailed, ErrBorrowedCredentialsExpired)
m.updateStatusByDigest(ref, StatusFailed, ErrBorrowedCredentialsExpired, buildID)
return
}
ctx := context.Background()
Expand All @@ -389,7 +403,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest,
ctx, cancel = context.WithDeadline(ctx, deadline)
defer cancel()
}
m.buildImage(ctx, ref, credentials)
m.buildImage(ctx, ref, credentials, buildID)
Comment thread
cursor[bot] marked this conversation as resolved.
}, m.releaseInflightPull(ref.Digest(), inflight))

img := meta.toImage()
Expand All @@ -399,7 +413,7 @@ func (m *manager) createAndQueueImage(ref *ResolvedRef, req CreateImageRequest,
return img, nil
}

func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials *authn.AuthConfig) {
func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials *authn.AuthConfig, buildID string) {
buildStart := time.Now()
buildStatus := "failed"
buildDir := m.paths.SystemBuild(ref.String())
Expand All @@ -409,7 +423,7 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials
}()

if err := os.MkdirAll(buildDir, 0755); err != nil {
m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("create build dir: %w", err))
m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("create build dir: %w", err), buildID)
return
}

Expand All @@ -420,7 +434,7 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials
m.recordImageBuildPhase(ctx, ref.Digest(), "cleanup", time.Since(start), phaseStatus(err), "not_applicable")
}()

m.updateStatusByDigest(ref, StatusPulling, nil)
m.updateStatusByDigest(ref, StatusPulling, nil, buildID)

// Pull by the digest-pinned reference, not the tag: a digest ref fetches the
// exact manifest regardless of the platform passed downstream, so a
Expand All @@ -431,7 +445,7 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials
result, err := m.ociClient.pullAndExportWithAuth(ctx, pullRef, ref.Digest(), tempDir, credentials)
m.recordPullResultMetrics(ctx, ref.Digest(), result)
if err != nil {
m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("pull and export: %w", err))
m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("pull and export: %w", err), buildID)
m.recordPullMetrics(ctx, "failed")
return
}
Expand All @@ -449,38 +463,40 @@ func (m *manager) buildImage(ctx context.Context, ref *ResolvedRef, credentials
}
}

m.updateStatusByDigest(ref, StatusConverting, nil)
m.updateStatusByDigest(ref, StatusConverting, nil, buildID)

diskPath := digestPath(m.paths, ref.Repository(), ref.DigestHex())
// Use default image format (erofs on Linux, ext4 on Darwin)
convertStart := time.Now()
diskSize, err := ExportRootfs(tempDir, diskPath, DefaultImageFormat)
m.recordImageBuildPhase(ctx, ref.Digest(), "filesystem_export", time.Since(convertStart), phaseStatus(err), "not_applicable")
if err != nil {
m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("convert to %s: %w", DefaultImageFormat, err))
m.updateStatusByDigest(ref, StatusFailed, fmt.Errorf("convert to %s: %w", DefaultImageFormat, err), buildID)
return
}

finalizeStart := time.Now()
err = m.finalizeImage(ref, result, diskSize)
err = m.finalizeImage(ref, result, diskSize, buildID)
m.recordImageBuildPhase(ctx, ref.Digest(), "finalize", time.Since(finalizeStart), phaseStatus(err), "not_applicable")
if err != nil {
m.updateStatusByDigest(ref, StatusFailed, err)
if errors.Is(err, errStaleBuild) {
return
}
m.updateStatusByDigest(ref, StatusFailed, err, buildID)
return
}

buildStatus = "success"
}

func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64) error {
// Read current metadata to preserve request info.
func (m *manager) finalizeImage(ref *ResolvedRef, result *pullResult, diskSize int64, buildID string) error {
m.createMu.Lock()
defer m.createMu.Unlock()

// Read current metadata to preserve request info and reject stale builds.
meta, err := readMetadata(m.paths, ref.Repository(), ref.DigestHex())
if err != nil {
meta = &imageMetadata{
Name: ref.String(),
Digest: ref.Digest(),
CreatedAt: time.Now(),
}
if err != nil || meta.BuildID != buildID {
return errStaleBuild
}

// The pulled image config is the source of truth for the platform.
Expand Down Expand Up @@ -557,19 +573,15 @@ func (m *manager) recordImageBuildPhase(ctx context.Context, digest, phase strin
)
}

func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err error) {
func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err error, buildID string) {
m.createMu.Lock()
defer m.createMu.Unlock()

meta, readErr := readMetadata(m.paths, ref.Repository(), ref.DigestHex())
if readErr != nil {
// Create new metadata if it doesn't exist
meta = &imageMetadata{
Name: ref.String(),
Digest: ref.Digest(),
Status: status,
CreatedAt: time.Now(),
}
} else {
meta.Status = status
if readErr != nil || meta.BuildID != buildID {
return
}
meta.Status = status

if err != nil {
errorMsg := err.Error()
Expand All @@ -578,7 +590,8 @@ func (m *manager) updateStatusByDigest(ref *ResolvedRef, status string, err erro

writeMetadata(m.paths, ref.Repository(), ref.DigestHex(), meta)

// Notify subscribers of terminal status
// Notify while holding createMu so a delete/recreate cannot race the
// metadata write and receive a terminal event for the old build.
if status == StatusReady || status == StatusFailed {
m.notifyReady(ref.DigestHex(), status, err)
}
Expand Down Expand Up @@ -608,11 +621,12 @@ func (m *manager) RecoverInterruptedBuilds() {
}
ref := NewResolvedRef(normalized, meta.Digest)
if meta.BorrowedAuth && (meta.Status != StatusConverting || !m.ociClient.existsInLayout(digestToLayoutTag(meta.Digest))) {
m.updateStatusByDigest(ref, StatusFailed, ErrBorrowedCredentialsExpired)
m.updateStatusByDigest(ref, StatusFailed, ErrBorrowedCredentialsExpired, meta.BuildID)
continue
}
buildID := meta.BuildID
m.queue.Enqueue(meta.Digest, func() {
m.buildImage(context.Background(), ref, nil)
m.buildImage(context.Background(), ref, nil, buildID)
}, nil)
}
}
Expand Down
117 changes: 117 additions & 0 deletions lib/images/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -602,6 +602,123 @@ func TestImportLocalImageFromOCICache(t *testing.T) {
t.Logf("Disk path verified: %s (%d bytes)", diskPath, diskStat.Size())
}

func TestDeleteAndRecreateDuringBuildTail(t *testing.T) {
origFormat := DefaultImageFormat
DefaultImageFormat = FormatCpio
defer func() { DefaultImageFormat = origFormat }()

dataDir := t.TempDir()
p := paths.New(dataDir)
mgr, err := NewManager(p, 1, nil)
require.NoError(t, err)
m := mgr.(*manager)

ctx := context.Background()
const repo = "kernel.local/test/recreate-race"
const tag = "v1"

testImg := createTestDockerImage(t)
imgDigest, err := testImg.Digest()
require.NoError(t, err)
digestStr := imgDigest.String()

cacheDir := p.SystemOCICache()
layoutPath, err := layout.Write(cacheDir, empty.Index)
require.NoError(t, err)
require.NoError(t, layoutPath.AppendImage(testImg, layout.WithAnnotations(map[string]string{
"org.opencontainers.image.ref.name": digestToLayoutTag(digestStr),
})))

digestHex := digestToLayoutTag(digestStr)
events := make(chan StatusEvent, 2)
m.subscribeToReady(digestHex, events)
defer m.unsubscribeFromReady(digestHex, events)

_, err = m.ImportLocalImage(ctx, repo, tag, digestStr)
require.NoError(t, err)
select {
case event := <-events:
require.Equal(t, StatusReady, event.Status)
case <-time.After(30 * time.Second):
t.Fatal("first build did not become ready")
}
firstMeta, err := readMetadata(p, repo, digestHex)
require.NoError(t, err)
require.NotEmpty(t, firstMeta.BuildID)

slotHeld := make(chan struct{})
releaseSlot := make(chan struct{})
m.queue.EnqueueSuccessor(digestStr, func() {
close(slotHeld)
<-releaseSlot
}, nil)
select {
case <-slotHeld:
case <-time.After(5 * time.Second):
t.Fatal("queue slot was not held")
}

// Delete by digest so the test does not depend on the tag symlink being
// created after the ready notification.
require.NoError(t, m.DeleteImage(ctx, repo+"@"+digestStr))
recreated, err := m.ImportLocalImage(ctx, repo, tag, digestStr)
require.NoError(t, err)
require.Equal(t, StatusPending, recreated.Status)
require.NotNil(t, recreated.QueuePosition)
require.Equal(t, 1, *recreated.QueuePosition)

currentMeta, err := readMetadata(p, repo, digestHex)
require.NoError(t, err)
require.NotEqual(t, firstMeta.BuildID, currentMeta.BuildID)
require.NoError(t, m.DeleteImage(ctx, repo+"@"+digestStr))
recreatedAgain, err := m.ImportLocalImage(ctx, repo, tag, digestStr)
require.NoError(t, err)
require.Equal(t, StatusPending, recreatedAgain.Status)
require.NotNil(t, recreatedAgain.QueuePosition)
require.Equal(t, 1, *recreatedAgain.QueuePosition)
latestMeta, err := readMetadata(p, repo, digestHex)
require.NoError(t, err)
require.NotEqual(t, currentMeta.BuildID, latestMeta.BuildID)
currentMeta = latestMeta

waitCtx, cancelWait := context.WithCancel(ctx)
defer cancelWait()
waitResult := make(chan error, 1)
go func() {
waitResult <- m.WaitForReady(waitCtx, repo+"@"+digestStr)
}()
select {
case err := <-waitResult:
t.Fatalf("recreated image completed before successor ran: %v", err)
case <-time.After(100 * time.Millisecond):
}
normalized, err := ParseNormalizedRef(repo + "@" + digestStr)
require.NoError(t, err)
staleRef := NewResolvedRef(normalized, digestStr)
m.updateStatusByDigest(staleRef, StatusFailed, errors.New("stale build"), firstMeta.BuildID)
staleResult, _, _, err := m.ociClient.extractOCIImageDetails(digestHex)
require.NoError(t, err)
require.ErrorIs(t, m.finalizeImage(staleRef, &pullResult{Metadata: staleResult}, 1, firstMeta.BuildID), errStaleBuild)
currentMeta, err = readMetadata(p, repo, digestHex)
require.NoError(t, err)
require.Equal(t, StatusPending, currentMeta.Status)
require.Nil(t, currentMeta.Error)

close(releaseSlot)
select {
case event := <-events:
require.Equal(t, StatusReady, event.Status)
case <-time.After(30 * time.Second):
t.Fatal("recreated image did not become ready")
}
select {
case err := <-waitResult:
require.NoError(t, err)
case <-time.After(30 * time.Second):
t.Fatal("WaitForReady did not observe the successor build")
}
}

// waitForReady waits for an image build to complete
func waitForReady(t *testing.T, mgr Manager, ctx context.Context, imageName string) {
for i := 0; i < 600; i++ {
Expand Down
1 change: 1 addition & 0 deletions lib/images/storage.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ type imageMetadata struct {
WorkingDir string `json:"working_dir,omitempty"`
CreatedAt time.Time `json:"created_at"`
BorrowedAuth bool `json:"borrowed_auth,omitempty"`
BuildID string `json:"build_id,omitempty"`
}

func (m *imageMetadata) toImage() *Image {
Expand Down
Loading
Loading