Skip to content
Merged
39 changes: 28 additions & 11 deletions pkg/cgroup/cgroup.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"path/filepath"
"strconv"
"strings"
"syscall"
)

// Root is where a Linux host mounts the cgroup v2 hierarchy.
Expand Down Expand Up @@ -154,25 +155,32 @@ func Populated(dir string) (bool, error) {
return false, nil
}

// Procs lists the processes in a cgroup and in every cgroup under it; a cgroup that is gone answers ErrNotFound.
// Procs lists the processes in a cgroup and in every cgroup under it; a cgroup that is gone answers ErrNotFound, and a child that goes mid-walk is skipped.
func Procs(dir string) ([]int, error) {
var pids []int
err := filepath.WalkDir(dir, func(path string, entry fs.DirEntry, err error) error {
if err != nil && path == dir && errors.Is(err, fs.ErrNotExist) {
return ErrNotFound
}
// A guest may remove its own child cgroups at any time, so their ENOENT never reads as a cgroup v1 host (SHARD-637).
if path != dir && vanished(err) {
return fs.SkipDir
}
if err != nil {
return err
}
if !entry.IsDir() {
return nil
}

raw, err := read(path, "cgroup.procs")
raw, err := os.ReadFile(filepath.Join(path, "cgroup.procs"))
if path != dir && vanished(err) {
return fs.SkipDir
}
if err != nil {
return err
return readFailed(path, "cgroup.procs", err)
}
for field := range strings.FieldsSeq(raw) {
for field := range strings.FieldsSeq(string(raw)) {
pid, err := strconv.Atoi(field)
if err != nil {
return fmt.Errorf("read %s: %q is not a pid", filepath.Join(path, "cgroup.procs"), field)
Expand Down Expand Up @@ -244,19 +252,28 @@ func readBound(dir, file string) (int64, error) {
}

func read(dir, file string) (string, error) {
path := filepath.Join(dir, file)

raw, err := os.ReadFile(path)
if errors.Is(err, fs.ErrNotExist) {
return "", fmt.Errorf("read %s: %w", path, missing(dir))
}
raw, err := os.ReadFile(filepath.Join(dir, file))
if err != nil {
return "", fmt.Errorf("read %s: %w", path, err)
return "", readFailed(dir, file, err)
}

return strings.TrimSpace(string(raw)), nil
}

func readFailed(dir, file string, err error) error {
path := filepath.Join(dir, file)
if errors.Is(err, fs.ErrNotExist) {
return fmt.Errorf("read %s: %w", path, missing(dir))
}

return fmt.Errorf("read %s: %w", path, err)
}

// vanished is how a cgroup removed under a reader fails: ENOENT once it is gone, ENODEV while the kernel tears its files down.
func vanished(err error) bool {
return errors.Is(err, fs.ErrNotExist) || errors.Is(err, syscall.ENODEV)
}

// missing tells the two ways a control file goes missing apart, because they need opposite answers:
// a cgroup that is gone is the caller's problem, and a controller that is gone is the host's.
func missing(dir string) error {
Expand Down
90 changes: 90 additions & 0 deletions pkg/cgroup/procs_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"path/filepath"
"slices"
"strings"
"syscall"
"testing"

"github.com/presmihaylov/shard/pkg/cgroup"
Expand Down Expand Up @@ -65,3 +66,92 @@ func TestProcsRefusesAnEntryThatIsNotAPid(t *testing.T) {
t.Fatalf("Procs = %v, want the entry named", err)
}
}

// removedMidRead makes dir's cgroup.procs a fifo that answers only once remove is gone, so the walk meets the removal at a known point.
func removedMidRead(t *testing.T, dir, contents, remove string) <-chan error {
t.Helper()

if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatalf("make the cgroup: %v", err)
}
fifo := filepath.Join(dir, "cgroup.procs")
if err := syscall.Mkfifo(fifo, 0o600); err != nil {
t.Fatalf("make the fifo: %v", err)
}

done := make(chan error, 1)
go func() {
w, err := os.OpenFile(fifo, os.O_WRONLY, 0)
if err != nil {
done <- err
return
}
if err := os.RemoveAll(remove); err != nil {
done <- errors.Join(err, w.Close())
return
}
_, err = w.WriteString(contents)
done <- errors.Join(err, w.Close())
}()

return done
}

func TestProcsSkipsAChildTheGuestRemovesBeforeItIsRead(t *testing.T) {
dir := t.TempDir()
procs(t, dir, "1101\n")
procs(t, filepath.Join(dir, "c2"), "1103\n")
done := removedMidRead(t, filepath.Join(dir, "c1"), "1102\n", filepath.Join(dir, "c2"))

got, err := cgroup.Procs(dir)
if err != nil {
t.Fatalf("Procs with a child gone mid-walk: %v", err)
}
if err := <-done; err != nil {
t.Fatalf("remove the child: %v", err)
}
if want := []int{1101, 1102}; !slices.Equal(got, want) {
t.Fatalf("Procs = %v, want %v", got, want)
}
}

func TestProcsSkipsAChildTheGuestRemovesAfterItsProcsAreRead(t *testing.T) {
dir := t.TempDir()
procs(t, dir, "1101\n")
child := filepath.Join(dir, "c1")
done := removedMidRead(t, child, "1102\n", child)

got, err := cgroup.Procs(dir)
if err != nil {
t.Fatalf("Procs with a child gone before its own children were listed: %v", err)
}
if err := <-done; err != nil {
t.Fatalf("remove the child: %v", err)
}
if want := []int{1101, 1102}; !slices.Equal(got, want) {
t.Fatalf("Procs = %v, want %v", got, want)
}
}

func TestProcsNeverReadsAChildMadeAgainAsCgroupV1(t *testing.T) {
dir := t.TempDir()
procs(t, dir, "1101\n")
if err := os.Mkdir(filepath.Join(dir, "c1"), 0o755); err != nil {
t.Fatalf("make the child: %v", err)
}

got, err := cgroup.Procs(dir)
if err != nil {
t.Fatalf("Procs with a child that has no cgroup.procs yet: %v", err)
}
if want := []int{1101}; !slices.Equal(got, want) {
t.Fatalf("Procs = %v, want %v", got, want)
}
}

func TestProcsOfACgroupWithoutItsProcsIsNoController(t *testing.T) {
_, err := cgroup.Procs(t.TempDir())
if !errors.Is(err, cgroup.ErrNoController) {
t.Fatalf("Procs on a cgroup with no cgroup.procs = %v, want ErrNoController", err)
}
}
19 changes: 18 additions & 1 deletion services/image/image.go
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,19 @@ func (s *Service) pullLocked(ctx context.Context, ref string) (Image, error) {
return Image{}, err
}

// A held tag is never resolved again, so the format this provider lacks comes from its layers, and a failed build leaves it held.
held, err := s.store.Get(ref)
if err == nil {
if err := s.unpack(ctx, held); err != nil {
return Image{}, err
}

return s.finish(progress, held)
}
if !errors.Is(err, registry.ErrNotCached) {
return Image{}, err
}

pulled, err := s.store.Pull(ctx, ref, pullReport{p: progress})
if err != nil {
return Image{}, errors.Join(err, s.reclaim())
Expand All @@ -183,7 +196,11 @@ func (s *Service) pullLocked(ctx context.Context, ref string) (Image, error) {
return Image{}, errors.Join(err, s.store.Remove(ref))
}

img, err = s.describe(pulled)
return s.finish(progress, pulled)
}

func (s *Service) finish(progress *Progress, unpacked registry.Image) (Image, error) {
img, err := s.describe(unpacked)
if err != nil {
return Image{}, err
}
Expand Down
4 changes: 2 additions & 2 deletions services/image/image_disk_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@ func TestPullWithDisksBuildsOneDiskPerDigest(t *testing.T) {
t.Fatalf("second Pull: %v", err)
}
again.Close()
// The tree is already there, so the rebuild says building alone and never unpacking (SHARD-385).
want := []string{image.StatusPulling, image.StatusLayer, image.StatusBuilding, image.StatusPulled}
// The tree and the tag are already held, so the rebuild never unpacks (SHARD-385) and never asks the registry.
want := []string{image.StatusBuilding, image.StatusPulled}
if got := statuses(events(t, again)); !slices.Equal(got, want) {
t.Errorf("the rebuild of the disk said %v, want %v", got, want)
}
Expand Down
20 changes: 13 additions & 7 deletions services/image/image_erofs_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,16 @@ func fakeMkfsErofs(t *testing.T) string {
return log
}

// failingMkfsErofs puts a mkfs.erofs stand-in first on PATH that fails the way a full disk does.
func failingMkfsErofs(t *testing.T) {
t.Helper()
dir := t.TempDir()
if err := os.WriteFile(filepath.Join(dir, "mkfs.erofs"), []byte("#!/bin/sh\necho 'no space' >&2; exit 1\n"), 0o755); err != nil {
t.Fatal(err)
}
t.Setenv("PATH", dir+string(os.PathListSeparator)+os.Getenv("PATH"))
}

func TestPullWithErofsBuildsOneImagePerDigestFromTheTree(t *testing.T) {
log := fakeMkfsErofs(t)
server, ref := servedImage(t, "app:1.0", map[string]string{"etc/hostname": "box"})
Expand Down Expand Up @@ -70,8 +80,8 @@ func TestPullWithErofsBuildsOneImagePerDigestFromTheTree(t *testing.T) {
t.Fatalf("second Pull: %v", err)
}
again.Close()
// The tree is already there, so the rebuild says building alone and never unpacking (SHARD-385).
want := []string{image.StatusPulling, image.StatusLayer, image.StatusBuilding, image.StatusPulled}
// The tree and the tag are already held, so the rebuild never unpacks (SHARD-385) and never asks the registry.
want := []string{image.StatusBuilding, image.StatusPulled}
if got := statuses(events(t, again)); !slices.Equal(got, want) {
t.Errorf("the rebuild of the image said %v, want %v", got, want)
}
Expand All @@ -81,11 +91,7 @@ func TestPullWithErofsBuildsOneImagePerDigestFromTheTree(t *testing.T) {
}

func TestPullWithErofsFailsWhenTheToolDoes(t *testing.T) {
dir := t.TempDir()
if err := os.WriteFile(filepath.Join(dir, "mkfs.erofs"), []byte("#!/bin/sh\necho 'no space' >&2; exit 1\n"), 0o755); err != nil {
t.Fatal(err)
}
t.Setenv("PATH", dir+string(os.PathListSeparator)+os.Getenv("PATH"))
failingMkfsErofs(t)
server, ref := servedImage(t, "app:1.0", map[string]string{"etc/hostname": "box"})
root := t.TempDir()
svc := newServiceAt(t, root, server, image.WithErofs())
Expand Down
108 changes: 108 additions & 0 deletions services/image/image_format_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
package image_test

import (
"os"
"path/filepath"
"slices"
"testing"

"github.com/presmihaylov/shard/services/image"
)

// formats are the files a VM provider boots from, which a container provider never built for the tags it pulled.
var formats = []struct {
name string
opt image.Option
path func(image.Image) string
}{
{name: "disk", opt: image.WithDisks(), path: func(img image.Image) string { return img.Disk }},
{name: "erofs", opt: image.WithErofs(), path: func(img image.Image) string { return img.Erofs }},
}

func TestANewFormatKeepsTheDigestTheTagHeld(t *testing.T) {
for _, format := range formats {
t.Run(format.name, func(t *testing.T) {
fakeMkfsErofs(t)
server, ref := servedImage(t, "app:1.0", map[string]string{"etc/hostname": "first"})
root := t.TempDir()
held, err := newServiceAt(t, root, server).Pull(t.Context(), ref)
if err != nil {
t.Fatalf("Pull: %v", err)
}
pushImage(t, server, "app:1.0", map[string]string{"etc/hostname": "second"})

progress := image.NewProgress()
img, err := newServiceAt(t, root, server, format.opt).Pull(image.WithProgress(t.Context(), progress), ref)
if err != nil {
t.Fatalf("Pull with the %s: %v", format.name, err)
}
progress.Close()
if img.Digest != held.Digest {
t.Errorf("the %s pull resolved the tag again: digest %s, want the held %s", format.name, img.Digest, held.Digest)
}
hostname, err := os.ReadFile(filepath.Join(img.RootFS, "etc", "hostname"))
if err != nil || string(hostname) != "first" {
t.Errorf("the tree holds %q (%v), want the held first", hostname, err)
}
if _, err := os.Stat(format.path(img)); err != nil {
t.Errorf("the %s is not there: %v", format.name, err)
}
want := []string{image.StatusBuilding, image.StatusPulled}
if got := statuses(events(t, progress)); !slices.Equal(got, want) {
t.Errorf("the %s pull said %v, want %v", format.name, got, want)
}
})
}
}

func TestANewFormatBuildsWithTheRegistryGone(t *testing.T) {
for _, format := range formats {
t.Run(format.name, func(t *testing.T) {
fakeMkfsErofs(t)
server, ref := servedImage(t, "app:1.0", map[string]string{"etc/hostname": "box"})
root := t.TempDir()
held, err := newServiceAt(t, root, server).Pull(t.Context(), ref)
if err != nil {
t.Fatalf("Pull: %v", err)
}
server.Close()

img, err := newServiceAt(t, root, server, format.opt).Pull(t.Context(), ref)
if err != nil {
t.Fatalf("Pull with the %s and no registry: %v", format.name, err)
}
if img.Digest != held.Digest {
t.Errorf("digest %s, want the held %s", img.Digest, held.Digest)
}
if _, err := os.Stat(format.path(img)); err != nil {
t.Errorf("the %s is not there: %v", format.name, err)
}
})
}
}

func TestAFailedBuildOfANewFormatKeepsTheHeldTag(t *testing.T) {
server, ref := servedImage(t, "app:1.0", map[string]string{"etc/hostname": "box"})
root := t.TempDir()
held, err := newServiceAt(t, root, server).Pull(t.Context(), ref)
if err != nil {
t.Fatalf("Pull: %v", err)
}
failingMkfsErofs(t)

svc := newServiceAt(t, root, server, image.WithErofs())
if _, err := svc.Pull(t.Context(), ref); err == nil {
t.Fatal("Pull with a failing mkfs.erofs returned no error")
}

images, err := svc.List()
if err != nil {
t.Fatalf("List: %v", err)
}
if len(images) != 1 || images[0].Digest != held.Digest {
t.Errorf("the store lists %+v, want the held %s", images, held.Digest)
}
if _, err := os.Stat(held.RootFS); err != nil {
t.Errorf("the held tree went: %v", err)
}
}
4 changes: 2 additions & 2 deletions services/provider/firecracker/checkpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -405,7 +405,7 @@ func (p *Provider) restore(ctx context.Context, id, stateDir string, r record, d
return nil, err
}
if err := m.reseed(ctx); err != nil {
return nil, errors.Join(err, p.end(ctx, m))
return nil, errors.Join(err, p.endAnyway(ctx, m))
}

// Only a running sandbox is ever paused, so what a checkpoint brings back is running and Status says so.
Expand Down Expand Up @@ -507,7 +507,7 @@ func (p *Provider) forkCheckpoint(ctx context.Context, dir string, spec models.S
}
// The restored guest still answers to the source's address and MAC, which the readdress replaces in place.
if err := m.readdress(ctx, r); err != nil {
return errors.Join(err, p.end(ctx, m), os.Remove(filepath.Join(spec.StateDir, recordFile)))
return errors.Join(err, p.endAnyway(ctx, m), os.Remove(filepath.Join(spec.StateDir, recordFile)))
}

return nil
Expand Down
Loading
Loading