Skip to content
Open
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
74 changes: 41 additions & 33 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 @@ -365,6 +371,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 +384,11 @@ 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() {
buildID := meta.BuildID
queuePos := m.queue.EnqueueSuccessor(ref.Digest(), func() {
credentials, deadline, expired := m.borrowedAuth(ref.Digest())
if expired {
m.updateStatusByDigest(ref, StatusFailed, ErrBorrowedCredentialsExpired)
m.updateStatusByDigest(ref, StatusFailed, ErrBorrowedCredentialsExpired, buildID)
return
}
ctx := context.Background()
Expand All @@ -389,7 +397,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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Recreate leaves pending builds stuck

High Severity

createAndQueueImage always writes a fresh BuildID and captures it in the queued closure, but EnqueueSuccessor still deduplicates against an existing pending job and keeps that older closure. After delete/recreate while the digest is already pending—or a second recreate behind an active predecessor—the kept job is rejected as stale and the new metadata never runs, so the image stays pending indefinitely. A later CreateImage also returns that pending entry without re-enqueueing.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 2a06de7. Configure here.

}, m.releaseInflightPull(ref.Digest(), inflight))

img := meta.toImage()
Expand All @@ -399,7 +407,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 +417,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 +428,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 +439,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 +457,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 +567,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 +584,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 +615,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