Loading internal/submit/blobstore.go +42 −1 Changes for internal/submit/blobstore.go: 42 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -2,6 +2,7 @@ package submit import ( "context" "errors" "fmt" "io" "os" Loading Loading @@ -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 Loading Loading @@ -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 { Loading @@ -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) Loading @@ -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. Loading internal/submit/blobstore_test.go +96 −0 Changes for internal/submit/blobstore_test.go: 96 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -7,6 +7,7 @@ import ( "os" "path/filepath" "strings" "syscall" "testing" ) Loading Loading @@ -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) { Loading Loading
internal/submit/blobstore.go +42 −1 Changes for internal/submit/blobstore.go: 42 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -2,6 +2,7 @@ package submit import ( "context" "errors" "fmt" "io" "os" Loading Loading @@ -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 Loading Loading @@ -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 { Loading @@ -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) Loading @@ -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. Loading
internal/submit/blobstore_test.go +96 −0 Changes for internal/submit/blobstore_test.go: 96 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -7,6 +7,7 @@ import ( "os" "path/filepath" "strings" "syscall" "testing" ) Loading Loading @@ -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) { Loading