Skip to content

Commit 81efb9e

Browse files
committed
review: whole-repo round — unify backend capability gates, extend fsync-skip to streamed paths
vm debug shared none of createVM's capability gate list and had drifted (missing fc+windows and windows-requires-cloudimg); both now call the same validators. ExtractTar's per-file fsync goes away: clone/restore streams follow the #160 no-sync contract and snapshot Create/Import rely on the SyncTree barrier — Import previously never made dir entries durable before the record committed, so it gains that barrier. Multiprocess storm tests share one spawn/wait harness in contracttest instead of four copies.
1 parent e9502f6 commit 81efb9e

8 files changed

Lines changed: 103 additions & 117 deletions

File tree

cmd/vm/debug.go

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -44,11 +44,8 @@ func (h Handler) Debug(cmd *cobra.Command, args []string) error {
4444
if err != nil {
4545
return err
4646
}
47-
if conf.UseFirecracker && vmCfg.SharedMemory {
48-
return fmt.Errorf("--fc and --shared-memory are mutually exclusive: Firecracker does not support vhost-user-fs hot-plug")
49-
}
50-
if conf.UseFirecracker && vmCfg.HugePages {
51-
return fmt.Errorf("--fc and --hugepages are mutually exclusive: Firecracker cannot restore hugetlbfs-backed snapshots")
47+
if err = validateBackendFlags(conf, vmCfg); err != nil {
48+
return err
5249
}
5350
if len(vmCfg.DataDisks) > 0 {
5451
fmt.Fprintln(os.Stderr, "warning: --data-disk is ignored in debug mode (debug only prints the hypervisor launch command; data disks need PrepareDataDisks to materialize)")
@@ -58,11 +55,11 @@ func (h Handler) Debug(cmd *cobra.Command, args []string) error {
5855
if err != nil {
5956
return err
6057
}
58+
if err = validateBootCompat(conf, vmCfg, boot); err != nil {
59+
return err
60+
}
6161

6262
if conf.UseFirecracker {
63-
if boot.KernelPath == "" {
64-
return fmt.Errorf("--fc requires OCI images (direct kernel boot)")
65-
}
6663
// FC requires uncompressed ELF kernel — resolve vmlinux path for debug output.
6764
if err := firecracker.EnsureVmlinuxBoot(boot); err != nil {
6865
return err

cmd/vm/run.go

Lines changed: 31 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -419,14 +419,8 @@ func (h Handler) createVM(cmd *cobra.Command, image string) (context.Context, *t
419419
return nil, nil, nil, err
420420
}
421421

422-
if conf.UseFirecracker && vmCfg.Windows {
423-
return nil, nil, nil, fmt.Errorf("--fc and --windows are mutually exclusive: Firecracker does not support Windows guests")
424-
}
425-
if conf.UseFirecracker && vmCfg.SharedMemory {
426-
return nil, nil, nil, fmt.Errorf("--fc and --shared-memory are mutually exclusive: Firecracker does not support vhost-user-fs hot-plug")
427-
}
428-
if conf.UseFirecracker && vmCfg.HugePages {
429-
return nil, nil, nil, fmt.Errorf("--fc and --hugepages are mutually exclusive: Firecracker cannot restore hugetlbfs-backed snapshots")
422+
if err = validateBackendFlags(conf, vmCfg); err != nil {
423+
return nil, nil, nil, err
430424
}
431425
bridgeDev, _ := cmd.Flags().GetString("bridge")
432426
if bridgeDev != "" && vmCfg.Network != "" {
@@ -446,11 +440,8 @@ func (h Handler) createVM(cmd *cobra.Command, image string) (context.Context, *t
446440
if err != nil {
447441
return nil, nil, nil, err
448442
}
449-
if vmCfg.Windows && bootCfg.KernelPath != "" {
450-
return nil, nil, nil, fmt.Errorf("--windows requires cloudimg (UEFI boot), got OCI direct boot image")
451-
}
452-
if conf.UseFirecracker && bootCfg.KernelPath == "" {
453-
return nil, nil, nil, fmt.Errorf("--fc requires OCI images (direct kernel boot): Firecracker does not support UEFI/cloudimg boot")
443+
if err = validateBootCompat(conf, vmCfg, bootCfg); err != nil {
444+
return nil, nil, nil, err
454445
}
455446
cmdcore.EnsureFirmwarePath(conf, bootCfg)
456447

@@ -485,6 +476,33 @@ func (h Handler) createVM(cmd *cobra.Command, image string) (context.Context, *t
485476
return ctx, info, hyper, nil
486477
}
487478

479+
// validateBackendFlags fast-fails flag combinations the selected backend can never launch; boot-mode-dependent checks live in validateBootCompat. Shared by create and debug so the capability gate list cannot drift.
480+
func validateBackendFlags(conf *config.Config, vmCfg *types.VMConfig) error {
481+
if !conf.UseFirecracker {
482+
return nil
483+
}
484+
switch {
485+
case vmCfg.Windows:
486+
return fmt.Errorf("--fc and --windows are mutually exclusive: Firecracker does not support Windows guests")
487+
case vmCfg.SharedMemory:
488+
return fmt.Errorf("--fc and --shared-memory are mutually exclusive: Firecracker does not support vhost-user-fs hot-plug")
489+
case vmCfg.HugePages:
490+
return fmt.Errorf("--fc and --hugepages are mutually exclusive: Firecracker cannot restore hugetlbfs-backed snapshots")
491+
}
492+
return nil
493+
}
494+
495+
func validateBootCompat(conf *config.Config, vmCfg *types.VMConfig, bootCfg *types.BootConfig) error {
496+
directBoot := hypervisor.IsDirectBoot(bootCfg)
497+
if vmCfg.Windows && directBoot {
498+
return fmt.Errorf("--windows requires cloudimg (UEFI boot), got OCI direct boot image")
499+
}
500+
if conf.UseFirecracker && !directBoot {
501+
return fmt.Errorf("--fc requires OCI images (direct kernel boot): Firecracker does not support UEFI/cloudimg boot")
502+
}
503+
return nil
504+
}
505+
488506
// pinResolvedBlobs holds the resolved image's digest locks until the reserve
489507
// commits; the empty set (bridge/dataless) pins nothing.
490508
func pinResolvedBlobs(ctx context.Context, backends []imagebackend.Images, ref string, blobIDs map[string]struct{}) (func(), error) {

hypervisor/cloudhypervisor/utils.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -181,8 +181,7 @@ func restoreVM(ctx context.Context, hc *http.Client, sourceDir, restoreMode stri
181181
case restoreModeMmap:
182182
cfg.MemoryRestoreMode = chMemoryRestoreMmap
183183
default:
184-
// Fail loud rather than silently falling back to eager copy — a
185-
// misspelled or wrong-case mode is a ~6x restore-latency regression.
184+
// Fail loud: a misspelled mode silently falling back to eager copy is a ~6x restore-latency regression.
186185
return fmt.Errorf("unknown restore mode %q", restoreMode)
187186
}
188187
body, err := json.Marshal(cfg)

meta/contracttest/workers.go

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
package contracttest
2+
3+
import (
4+
"bytes"
5+
"os"
6+
"os/exec"
7+
"testing"
8+
)
9+
10+
// SpawnWorkers re-execs the test binary n times against the helper test testName; env yields each worker's extra environment.
11+
func SpawnWorkers(t *testing.T, n int, testName string, env func(w int) []string) []*exec.Cmd {
12+
t.Helper()
13+
cmds := make([]*exec.Cmd, 0, n)
14+
for w := range n {
15+
cmd := exec.Command(os.Args[0], "-test.run="+testName+"$", "-test.count=1") //nolint:gosec
16+
cmd.Env = append(os.Environ(), env(w)...)
17+
out := &bytes.Buffer{}
18+
cmd.Stdout, cmd.Stderr = out, out
19+
if err := cmd.Start(); err != nil {
20+
t.Fatalf("start worker %d: %v", w, err)
21+
}
22+
cmds = append(cmds, cmd)
23+
}
24+
return cmds
25+
}
26+
27+
// WaitWorkers fails on the first worker that exits non-zero, echoing its combined output; what labels the failure.
28+
func WaitWorkers(t *testing.T, cmds []*exec.Cmd, what string) {
29+
t.Helper()
30+
for w, cmd := range cmds {
31+
if err := cmd.Wait(); err != nil {
32+
t.Fatalf("worker %d %s: %v\n%s", w, what, err, cmd.Stdout)
33+
}
34+
}
35+
}

meta/json/multiprocess_test.go

Lines changed: 9 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ package json
22

33
import (
44
"bufio"
5-
"bytes"
65
"fmt"
76
"os"
87
"os/exec"
@@ -12,6 +11,7 @@ import (
1211
"time"
1312

1413
"github.com/cocoonstack/cocoon/meta"
14+
"github.com/cocoonstack/cocoon/meta/contracttest"
1515
)
1616

1717
const (
@@ -39,27 +39,10 @@ func TestMultiProcessCorrectness(t *testing.T) {
3939
}
4040
dir := t.TempDir()
4141
workers := stormWorkers()
42-
cmds := make([]*exec.Cmd, 0, workers)
43-
for w := range workers {
44-
cmd := exec.Command(os.Args[0], "-test.run=TestMultiProcessWorker$", "-test.count=1") //nolint:gosec
45-
cmd.Env = append(
46-
os.Environ(),
47-
"META_MP_DIR="+dir,
48-
"META_MP_WORKER="+strconv.Itoa(w),
49-
fmt.Sprintf("META_MP_OPS=%d", mpOps),
50-
)
51-
out := &bytes.Buffer{}
52-
cmd.Stdout, cmd.Stderr = out, out
53-
if err := cmd.Start(); err != nil {
54-
t.Fatalf("start worker %d: %v", w, err)
55-
}
56-
cmds = append(cmds, cmd)
57-
}
58-
for w, cmd := range cmds {
59-
if err := cmd.Wait(); err != nil {
60-
t.Fatalf("worker %d failed: %v\n%s", w, err, cmd.Stdout)
61-
}
62-
}
42+
cmds := contracttest.SpawnWorkers(t, workers, "TestMultiProcessWorker", func(w int) []string {
43+
return []string{"META_MP_DIR=" + dir, "META_MP_WORKER=" + strconv.Itoa(w), fmt.Sprintf("META_MP_OPS=%d", mpOps)}
44+
})
45+
contracttest.WaitWorkers(t, cmds, "failed")
6346

6447
s := newStore(t, dir, "alpha")
6548
c := meta.NewCollection[map[string]int](s, "alpha", "records")
@@ -115,27 +98,10 @@ func TestInverseScopeNoDeadlock(t *testing.T) {
11598
}
11699
dir := t.TempDir()
117100
scopes := []string{"alpha:beta", "beta:alpha"}
118-
cmds := make([]*exec.Cmd, 0, 2*inverseWorkers)
119-
for w := range 2 * inverseWorkers {
120-
cmd := exec.Command(os.Args[0], "-test.run=TestInverseScopeWorker$", "-test.count=1") //nolint:gosec
121-
cmd.Env = append(
122-
os.Environ(),
123-
"META_MP_DIR="+dir,
124-
"META_MP_WORKER="+strconv.Itoa(w),
125-
"META_MP_SCOPE="+scopes[w%2],
126-
)
127-
out := &bytes.Buffer{}
128-
cmd.Stdout, cmd.Stderr = out, out
129-
if err := cmd.Start(); err != nil {
130-
t.Fatalf("start worker %d: %v", w, err)
131-
}
132-
cmds = append(cmds, cmd)
133-
}
134-
for w, cmd := range cmds {
135-
if err := cmd.Wait(); err != nil {
136-
t.Fatalf("worker %d failed (deadlock or error): %v\n%s", w, err, cmd.Stdout)
137-
}
138-
}
101+
cmds := contracttest.SpawnWorkers(t, 2*inverseWorkers, "TestInverseScopeWorker", func(w int) []string {
102+
return []string{"META_MP_DIR=" + dir, "META_MP_WORKER=" + strconv.Itoa(w), "META_MP_SCOPE=" + scopes[w%2]}
103+
})
104+
contracttest.WaitWorkers(t, cmds, "failed (deadlock or error)")
139105
s := newStore(t, dir, "alpha", "beta")
140106
for _, ns := range []string{"alpha", "beta"} {
141107
total := 0

meta/sqlite/multiprocess_test.go

Lines changed: 9 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@ package sqlite
22

33
import (
44
"bufio"
5-
"bytes"
65
"encoding/json"
76
"errors"
87
"fmt"
@@ -14,6 +13,7 @@ import (
1413
"testing"
1514

1615
"github.com/cocoonstack/cocoon/meta"
16+
"github.com/cocoonstack/cocoon/meta/contracttest"
1717
)
1818

1919
const (
@@ -70,27 +70,10 @@ func TestMultiProcessCorrectness(t *testing.T) {
7070
dir := t.TempDir()
7171
_ = newStore(t, dir, "alpha") // parent initializes once; workers only open
7272
workers := stormWorkers()
73-
cmds := make([]*exec.Cmd, 0, workers)
74-
for w := range workers {
75-
cmd := exec.Command(os.Args[0], "-test.run=TestMultiProcessWorker$", "-test.count=1") //nolint:gosec
76-
cmd.Env = append(
77-
os.Environ(),
78-
"META_MP_DIR="+dir,
79-
"META_MP_WORKER="+strconv.Itoa(w),
80-
fmt.Sprintf("META_MP_OPS=%d", mpOps),
81-
)
82-
out := &bytes.Buffer{}
83-
cmd.Stdout, cmd.Stderr = out, out
84-
if err := cmd.Start(); err != nil {
85-
t.Fatalf("start worker %d: %v", w, err)
86-
}
87-
cmds = append(cmds, cmd)
88-
}
89-
for w, cmd := range cmds {
90-
if err := cmd.Wait(); err != nil {
91-
t.Fatalf("worker %d failed: %v\n%s", w, err, cmd.Stdout)
92-
}
93-
}
73+
cmds := contracttest.SpawnWorkers(t, workers, "TestMultiProcessWorker", func(w int) []string {
74+
return []string{"META_MP_DIR=" + dir, "META_MP_WORKER=" + strconv.Itoa(w), fmt.Sprintf("META_MP_OPS=%d", mpOps)}
75+
})
76+
contracttest.WaitWorkers(t, cmds, "failed")
9477

9578
s := newStore(t, dir, "alpha")
9679
records := map[string]struct{}{}
@@ -157,27 +140,10 @@ func TestInverseScopeNoDeadlock(t *testing.T) {
157140
dir := t.TempDir()
158141
_ = newStore(t, dir, "alpha", "beta")
159142
scopes := []string{"alpha:beta", "beta:alpha"}
160-
cmds := make([]*exec.Cmd, 0, 2*inverseWorkers)
161-
for w := range 2 * inverseWorkers {
162-
cmd := exec.Command(os.Args[0], "-test.run=TestInverseScopeWorker$", "-test.count=1") //nolint:gosec
163-
cmd.Env = append(
164-
os.Environ(),
165-
"META_MP_DIR="+dir,
166-
"META_MP_WORKER="+strconv.Itoa(w),
167-
"META_MP_SCOPE="+scopes[w%2],
168-
)
169-
out := &bytes.Buffer{}
170-
cmd.Stdout, cmd.Stderr = out, out
171-
if err := cmd.Start(); err != nil {
172-
t.Fatalf("start worker %d: %v", w, err)
173-
}
174-
cmds = append(cmds, cmd)
175-
}
176-
for w, cmd := range cmds {
177-
if err := cmd.Wait(); err != nil {
178-
t.Fatalf("worker %d failed (deadlock or error): %v\n%s", w, err, cmd.Stdout)
179-
}
180-
}
143+
cmds := contracttest.SpawnWorkers(t, 2*inverseWorkers, "TestInverseScopeWorker", func(w int) []string {
144+
return []string{"META_MP_DIR=" + dir, "META_MP_WORKER=" + strconv.Itoa(w), "META_MP_SCOPE=" + scopes[w%2]}
145+
})
146+
contracttest.WaitWorkers(t, cmds, "failed (deadlock or error)")
181147
s := newStore(t, dir, "alpha", "beta")
182148
for _, ns := range []string{"alpha", "beta"} {
183149
total := 0

snapshot/localfile/import.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,11 @@ func (lf *LocalFile) Import(ctx context.Context, r io.Reader, name, description
7070
return "", err
7171
}
7272

73+
// Files and dir entries durable before the record exists (ExtractTar never fsyncs).
74+
if err = utils.SyncTree(dataDir); err != nil {
75+
return "", fmt.Errorf("sync snapshot to disk: %w", err)
76+
}
77+
7378
size, sizeErr := utils.DirSize(dataDir)
7479
if sizeErr != nil {
7580
return "", fmt.Errorf("compute data dir size: %w", sizeErr)

utils/tar.go

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ func TarDir(tw *tar.Writer, dir string) error {
5454
return nil
5555
}
5656

57-
// ExtractTar extracts flat tar entries into dir; entries matching any skip predicate are dropped.
57+
// ExtractTar extracts flat tar entries into dir; entries matching any skip predicate are dropped. It never fsyncs — callers needing durability follow with SyncTree.
5858
func ExtractTar(dir string, r io.Reader, skip ...func(name string) bool) error {
5959
tr := tar.NewReader(r)
6060
for {
@@ -112,17 +112,17 @@ func tarFileFrom(tw *tar.Writer, f *os.File, fi os.FileInfo, nameInTar string) e
112112
}
113113

114114
// extractFileSparse restores a sparse file from its segment map.
115-
func extractFileSparse(path string, r io.Reader, perm os.FileMode, realSize int64, mapJSON string) error {
115+
func extractFileSparse(path string, r io.Reader, perm os.FileMode, realSize int64, mapJSON string) (err error) {
116116
var segments []sparseSegment
117-
if err := json.Unmarshal([]byte(mapJSON), &segments); err != nil {
117+
if err = json.Unmarshal([]byte(mapJSON), &segments); err != nil {
118118
return fmt.Errorf("decode sparse map: %w", err)
119119
}
120120

121121
f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, perm) //nolint:gosec
122122
if err != nil {
123123
return err
124124
}
125-
defer f.Close() //nolint:errcheck
125+
defer func() { err = errors.Join(err, f.Close()) }()
126126

127127
if err := f.Truncate(realSize); err != nil {
128128
return err
@@ -137,15 +137,15 @@ func extractFileSparse(path string, r io.Reader, perm os.FileMode, realSize int6
137137
}
138138
}
139139

140-
return f.Sync()
140+
return nil
141141
}
142142

143-
func extractFile(path string, r io.Reader, perm os.FileMode) error {
143+
func extractFile(path string, r io.Reader, perm os.FileMode) (err error) {
144144
f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, perm) //nolint:gosec
145145
if err != nil {
146146
return err
147147
}
148-
defer f.Close() //nolint:errcheck
148+
defer func() { err = errors.Join(err, f.Close()) }()
149149

150150
bp := sparseBlockPool.Get().(*[]byte)
151151
buf := *bp
@@ -178,7 +178,7 @@ func extractFile(path string, r io.Reader, perm os.FileMode) error {
178178
}
179179
}
180180

181-
return f.Sync()
181+
return nil
182182
}
183183

184184
// writeBlockSparse seeks over all-zero chunks instead of writing them.

0 commit comments

Comments
 (0)