diff --git a/internal/submit/blobstore.go b/internal/submit/blobstore.go index bc21936..06bebee 100644 --- a/internal/submit/blobstore.go +++ b/internal/submit/blobstore.go @@ -2,6 +2,7 @@ package submit import ( "context" + "errors" "fmt" "io" "os" @@ -47,6 +48,10 @@ type LocalContextStore struct { // MinFree is the share of Base's filesystem an upload must leave free; 0 uses // DefaultUploadsMinFree. MinFree float64 + + // sync flushes a file or directory to disk; nil is (*os.File).Sync. Tests + // replace it to watch or fail the flushes. + sync func(*os.File) error } // DefaultUploadsMinFree is the share of the uploads filesystem an upload must @@ -98,7 +103,9 @@ func (s *LocalContextStore) dir(id string) (string, error) { // fully successful copy. So a failed, truncated, or oversize upload never // replaces a good context and never leaves a half-written blob for Kaniko to // read; a re-upload while the submission is still pending simply supersedes the -// previous one. It returns the number of bytes stored. +// previous one. It returns the number of bytes stored, and once it does the blob +// survives a crash: the bytes are flushed before the rename and the directory +// entries after it. func (s *LocalContextStore) Put(_ context.Context, id string, r io.Reader) (int64, error) { dir, err := s.dir(id) if err != nil { @@ -113,6 +120,12 @@ func (s *LocalContextStore) Put(_ context.Context, id string, r io.Reader) (int6 } tmpName := tmp.Name() n, err := io.Copy(tmp, r) + if err == nil { + // The bytes reach the disk before the rename can publish them, so a crash + // right after it finds the whole blob behind the digest the database + // already holds. + err = s.flush(tmp) + } if err != nil { tmp.Close() os.Remove(tmpName) @@ -126,9 +139,37 @@ func (s *LocalContextStore) Put(_ context.Context, id string, r io.Reader) (int6 os.Remove(tmpName) return 0, fmt.Errorf("submit: commit context blob: %w", err) } + // The rename, and a first upload's new directory, are entries in their parent + // directories: flush those too, or a crash can undo them. + for _, d := range []string{dir, s.Base} { + if err := s.syncDir(d); err != nil { + return 0, fmt.Errorf("submit: sync context dir: %w", err) + } + } return n, nil } +func (s *LocalContextStore) flush(f *os.File) error { + if s.sync != nil { + return s.sync(f) + } + return f.Sync() +} + +// syncDir flushes dir's entries to disk. A filesystem that cannot sync a +// directory has no stronger promise to give. +func (s *LocalContextStore) syncDir(dir string) error { + d, err := os.Open(dir) + if err != nil { + return err + } + defer d.Close() + if err := s.flush(d); err != nil && !errors.Is(err, syscall.EINVAL) && !errors.Is(err, syscall.ENOTSUP) { + return err + } + return nil +} + // Exists reports whether a context blob has been stored for id. Approve consults // it so a submission whose context was never uploaded is refused BEFORE the CAS, // instead of being approved into a build Kaniko cannot pull. diff --git a/internal/submit/blobstore_test.go b/internal/submit/blobstore_test.go index 69d4e54..2f395ca 100644 --- a/internal/submit/blobstore_test.go +++ b/internal/submit/blobstore_test.go @@ -7,6 +7,7 @@ import ( "os" "path/filepath" "strings" + "syscall" "testing" ) @@ -109,6 +110,101 @@ func TestLocalContextStorePutOverwrites(t *testing.T) { } } +// Put flushes the new bytes to disk before the rename publishes them, then the +// directory entries the rename and a first upload's mkdir wrote. A crash at any +// point leaves the previous blob or the whole new one, never an empty file behind +// a digest the database already recorded. +func TestLocalContextStorePutFlushesAroundTheRename(t *testing.T) { + base := t.TempDir() + s := &LocalContextStore{Base: base} + blob := filepath.Join(base, "sub-abc", contextBlobName) + var synced []string + s.sync = func(f *os.File) error { + name := f.Name() + if strings.HasSuffix(name, ".tmp") { + name = "temp" + } + published := "none" + if b, err := os.ReadFile(blob); err == nil { + published = string(b) + } + synced = append(synced, name+" published="+published) + return f.Sync() + } + + if _, err := s.Put(context.Background(), "sub-abc", strings.NewReader("\x1f\x8bbytes")); err != nil { + t.Fatalf("Put: %v", err) + } + want := []string{ + "temp published=none", + filepath.Join(base, "sub-abc") + " published=\x1f\x8bbytes", + base + " published=\x1f\x8bbytes", + } + if strings.Join(synced, "\n") != strings.Join(want, "\n") { + t.Fatalf("flushes =\n%q\nwant\n%q", synced, want) + } +} + +// A flush that fails refuses the upload. Before the rename the previous context +// stays in place and no temp file is left; after it the new bytes are in place +// but not promised, so the error still reaches the caller. A filesystem that +// cannot sync a directory at all has no stronger promise to give. +func TestLocalContextStorePutFlushFailures(t *testing.T) { + boom := errors.New("disk said no") + for _, tc := range []struct { + name string + fail func(f *os.File) error + wantErr error + want string + }{ + {"file", func(f *os.File) error { + if strings.HasSuffix(f.Name(), ".tmp") { + return boom + } + return nil + }, boom, "\x1f\x8bold"}, + {"directory", func(f *os.File) error { + if strings.HasSuffix(f.Name(), ".tmp") { + return nil + } + return boom + }, boom, "\x1f\x8bnew"}, + {"directory sync invalid", func(f *os.File) error { + if strings.HasSuffix(f.Name(), ".tmp") { + return nil + } + return &os.PathError{Op: "sync", Path: f.Name(), Err: syscall.EINVAL} + }, nil, "\x1f\x8bnew"}, + {"directory sync not supported", func(f *os.File) error { + if strings.HasSuffix(f.Name(), ".tmp") { + return nil + } + return &os.PathError{Op: "sync", Path: f.Name(), Err: syscall.ENOTSUP} + }, nil, "\x1f\x8bnew"}, + } { + t.Run(tc.name, func(t *testing.T) { + base := t.TempDir() + s := &LocalContextStore{Base: base} + ctx := context.Background() + if _, err := s.Put(ctx, "sub-abc", strings.NewReader("\x1f\x8bold")); err != nil { + t.Fatalf("first Put: %v", err) + } + s.sync = tc.fail + if _, err := s.Put(ctx, "sub-abc", strings.NewReader("\x1f\x8bnew")); !errors.Is(err, tc.wantErr) { + t.Fatalf("Put = %v, want %v", err, tc.wantErr) + } + got, err := os.ReadFile(filepath.Join(base, "sub-abc", contextBlobName)) + if err != nil || string(got) != tc.want { + t.Fatalf("stored = %q, %v; want %q", got, err, tc.want) + } + entries, _ := os.ReadDir(filepath.Join(base, "sub-abc")) + if len(entries) != 1 { + t.Fatalf("dir holds %d entries, want only %s", len(entries), contextBlobName) + } + }) + } +} + // Delete removes the blob and its id-namespaced directory, and is idempotent — // the retry-safety the withdraw/delete cleanup depends on. func TestLocalContextStoreDelete(t *testing.T) {