From 63ff0496470cfad11511e73d76fe3c2c3cf1611a Mon Sep 17 00:00:00 2001 From: dis Date: Fri, 17 Jul 2026 16:58:13 +0530 Subject: [PATCH 1/3] OCPBUGS-97746: verify and sync catalog.jsonl file is available before populating updateStatus --- .../core/clustercatalog_controller.go | 6 ++++ .../core/clustercatalog_controller_test.go | 2 ++ internal/catalogd/storage/localdir.go | 33 +++++++++++++++++++ internal/catalogd/storage/storage.go | 6 ++++ .../testutil/mock/storage/mock_instance.go | 14 ++++++++ 5 files changed, 61 insertions(+) diff --git a/internal/catalogd/controllers/core/clustercatalog_controller.go b/internal/catalogd/controllers/core/clustercatalog_controller.go index fedfe500f0..a8ced78f22 100644 --- a/internal/catalogd/controllers/core/clustercatalog_controller.go +++ b/internal/catalogd/controllers/core/clustercatalog_controller.go @@ -270,6 +270,12 @@ func (r *ClusterCatalogReconciler) reconcile(ctx context.Context, catalog *ocv1. } baseURL := r.Storage.BaseURL(catalog.Name) + if err := r.Storage.VerifyAndSync(catalog.Name); err != nil { + verifyErr := fmt.Errorf("error verifying catalog content before serving: %w", err) + updateStatusProgressing(&catalog.Status, catalog.GetGeneration(), verifyErr) + return ctrl.Result{}, verifyErr + } + updateStatusProgressing(&catalog.Status, catalog.GetGeneration(), nil) updateStatusServing(&catalog.Status, canonicalRef, unpackTime, baseURL, catalog.GetGeneration()) diff --git a/internal/catalogd/controllers/core/clustercatalog_controller_test.go b/internal/catalogd/controllers/core/clustercatalog_controller_test.go index f6cbe46dfb..786fa07696 100644 --- a/internal/catalogd/controllers/core/clustercatalog_controller_test.go +++ b/internal/catalogd/controllers/core/clustercatalog_controller_test.go @@ -31,9 +31,11 @@ func newMockStore(ctrl *gomock.Controller, shouldError bool) *mockstorage.MockIn if shouldError { m.EXPECT().Store(gomock.Any(), gomock.Any(), gomock.Any()).Return(errors.New("mockstore store error")).AnyTimes() m.EXPECT().Delete(gomock.Any()).Return(errors.New("mockstore delete error")).AnyTimes() + m.EXPECT().VerifyAndSync(gomock.Any()).Return(errors.New("mockstore verify error")).AnyTimes() } else { m.EXPECT().Store(gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes() m.EXPECT().Delete(gomock.Any()).Return(nil).AnyTimes() + m.EXPECT().VerifyAndSync(gomock.Any()).Return(nil).AnyTimes() } m.EXPECT().BaseURL(gomock.Any()).Return("URL").AnyTimes() m.EXPECT().ContentExists(gomock.Any()).Return(true).AnyTimes() diff --git a/internal/catalogd/storage/localdir.go b/internal/catalogd/storage/localdir.go index 0cd65933b0..77379ee70a 100644 --- a/internal/catalogd/storage/localdir.go +++ b/internal/catalogd/storage/localdir.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "io/fs" "net/http" "net/url" @@ -251,6 +252,38 @@ func (s *LocalDirV1) ContentExists(catalog string) bool { return true } +// VerifyAndSync verifies that catalog.jsonl exists at RootDir//catalog.jsonl, +// calls fsync to ensure the file is durably written to disk, and confirms the file is +// readable by attempting a small read. It should be called after Store() succeeds and +// before marking the catalog as Serving. +func (s *LocalDirV1) VerifyAndSync(catalog string) error { + s.m.RLock() + defer s.m.RUnlock() + + path := catalogFilePath(s.catalogDir(catalog)) + + if _, err := os.Stat(path); err != nil { + return fmt.Errorf("catalog.jsonl not found at %q: %w", path, err) + } + + f, err := os.Open(path) + if err != nil { + return fmt.Errorf("catalog.jsonl not readable at %q: %w", path, err) + } + defer f.Close() + + if err := f.Sync(); err != nil { + return fmt.Errorf("fsync failed for catalog.jsonl at %q: %w", path, err) + } + + buf := make([]byte, 1) + if _, err := f.Read(buf); err != nil && !errors.Is(err, io.EOF) { + return fmt.Errorf("catalog.jsonl read failed at %q: %w", path, err) + } + + return nil +} + func (s *LocalDirV1) catalogDir(catalog string) string { return filepath.Join(s.RootDir, catalog) } diff --git a/internal/catalogd/storage/storage.go b/internal/catalogd/storage/storage.go index af78a669fc..a16a7f4cad 100644 --- a/internal/catalogd/storage/storage.go +++ b/internal/catalogd/storage/storage.go @@ -15,6 +15,12 @@ type Instance interface { Delete(catalog string) error ContentExists(catalog string) bool + // VerifyAndSync confirms that catalog.jsonl exists on disk for the given + // catalog, flushes it to stable storage via fsync, and verifies the file + // is readable. It must be called after Store and before marking the + // catalog as Serving. + VerifyAndSync(catalog string) error + BaseURL(catalog string) string StorageServerHandler() http.Handler } diff --git a/internal/testutil/mock/storage/mock_instance.go b/internal/testutil/mock/storage/mock_instance.go index 2aa9613e3b..245a803c9d 100644 --- a/internal/testutil/mock/storage/mock_instance.go +++ b/internal/testutil/mock/storage/mock_instance.go @@ -111,3 +111,17 @@ func (mr *MockInstanceMockRecorder) Store(ctx, catalog, fsys any) *gomock.Call { mr.mock.ctrl.T.Helper() return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Store", reflect.TypeOf((*MockInstance)(nil).Store), ctx, catalog, fsys) } + +// VerifyAndSync mocks base method. +func (m *MockInstance) VerifyAndSync(catalog string) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "VerifyAndSync", catalog) + ret0, _ := ret[0].(error) + return ret0 +} + +// VerifyAndSync indicates an expected call of VerifyAndSync. +func (mr *MockInstanceMockRecorder) VerifyAndSync(catalog any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "VerifyAndSync", reflect.TypeOf((*MockInstance)(nil).VerifyAndSync), catalog) +} From 5416bc9855ba84fec0eeac8cacf23d57104ef883 Mon Sep 17 00:00:00 2001 From: dis Date: Tue, 15 Sep 2026 11:12:45 +0530 Subject: [PATCH 2/3] fsync storeCatalogData and storeIndexData after finished writing. --- .../core/clustercatalog_controller.go | 6 -- .../core/clustercatalog_controller_test.go | 2 - internal/catalogd/storage/localdir.go | 56 +++++++------------ internal/catalogd/storage/storage.go | 6 -- .../testutil/mock/storage/mock_instance.go | 14 ----- 5 files changed, 21 insertions(+), 63 deletions(-) diff --git a/internal/catalogd/controllers/core/clustercatalog_controller.go b/internal/catalogd/controllers/core/clustercatalog_controller.go index a8ced78f22..fedfe500f0 100644 --- a/internal/catalogd/controllers/core/clustercatalog_controller.go +++ b/internal/catalogd/controllers/core/clustercatalog_controller.go @@ -270,12 +270,6 @@ func (r *ClusterCatalogReconciler) reconcile(ctx context.Context, catalog *ocv1. } baseURL := r.Storage.BaseURL(catalog.Name) - if err := r.Storage.VerifyAndSync(catalog.Name); err != nil { - verifyErr := fmt.Errorf("error verifying catalog content before serving: %w", err) - updateStatusProgressing(&catalog.Status, catalog.GetGeneration(), verifyErr) - return ctrl.Result{}, verifyErr - } - updateStatusProgressing(&catalog.Status, catalog.GetGeneration(), nil) updateStatusServing(&catalog.Status, canonicalRef, unpackTime, baseURL, catalog.GetGeneration()) diff --git a/internal/catalogd/controllers/core/clustercatalog_controller_test.go b/internal/catalogd/controllers/core/clustercatalog_controller_test.go index 786fa07696..f6cbe46dfb 100644 --- a/internal/catalogd/controllers/core/clustercatalog_controller_test.go +++ b/internal/catalogd/controllers/core/clustercatalog_controller_test.go @@ -31,11 +31,9 @@ func newMockStore(ctrl *gomock.Controller, shouldError bool) *mockstorage.MockIn if shouldError { m.EXPECT().Store(gomock.Any(), gomock.Any(), gomock.Any()).Return(errors.New("mockstore store error")).AnyTimes() m.EXPECT().Delete(gomock.Any()).Return(errors.New("mockstore delete error")).AnyTimes() - m.EXPECT().VerifyAndSync(gomock.Any()).Return(errors.New("mockstore verify error")).AnyTimes() } else { m.EXPECT().Store(gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes() m.EXPECT().Delete(gomock.Any()).Return(nil).AnyTimes() - m.EXPECT().VerifyAndSync(gomock.Any()).Return(nil).AnyTimes() } m.EXPECT().BaseURL(gomock.Any()).Return("URL").AnyTimes() m.EXPECT().ContentExists(gomock.Any()).Return(true).AnyTimes() diff --git a/internal/catalogd/storage/localdir.go b/internal/catalogd/storage/localdir.go index 77379ee70a..777fc5bddf 100644 --- a/internal/catalogd/storage/localdir.go +++ b/internal/catalogd/storage/localdir.go @@ -5,7 +5,6 @@ import ( "encoding/json" "errors" "fmt" - "io" "io/fs" "net/http" "net/url" @@ -180,6 +179,13 @@ func (s *LocalDirV1) storeAtomicSwap(ctx context.Context, catalog string, fsys f return "", err } + if err := syncDir(catalogDir); err != nil { + return "", fmt.Errorf("error syncing catalog directory: %w", err) + } + if err := syncDir(s.RootDir); err != nil { + return "", fmt.Errorf("error syncing storage root directory: %w", err) + } + return catalogDir, nil } @@ -252,38 +258,6 @@ func (s *LocalDirV1) ContentExists(catalog string) bool { return true } -// VerifyAndSync verifies that catalog.jsonl exists at RootDir//catalog.jsonl, -// calls fsync to ensure the file is durably written to disk, and confirms the file is -// readable by attempting a small read. It should be called after Store() succeeds and -// before marking the catalog as Serving. -func (s *LocalDirV1) VerifyAndSync(catalog string) error { - s.m.RLock() - defer s.m.RUnlock() - - path := catalogFilePath(s.catalogDir(catalog)) - - if _, err := os.Stat(path); err != nil { - return fmt.Errorf("catalog.jsonl not found at %q: %w", path, err) - } - - f, err := os.Open(path) - if err != nil { - return fmt.Errorf("catalog.jsonl not readable at %q: %w", path, err) - } - defer f.Close() - - if err := f.Sync(); err != nil { - return fmt.Errorf("fsync failed for catalog.jsonl at %q: %w", path, err) - } - - buf := make([]byte, 1) - if _, err := f.Read(buf); err != nil && !errors.Is(err, io.EOF) { - return fmt.Errorf("catalog.jsonl read failed at %q: %w", path, err) - } - - return nil -} - func (s *LocalDirV1) catalogDir(catalog string) string { return filepath.Join(s.RootDir, catalog) } @@ -314,7 +288,7 @@ func storeCatalogData(catalogDir string, metas <-chan *declcfg.Meta) error { return err } } - return nil + return f.Sync() } func storeIndexData(catalogDir string, metas <-chan *declcfg.Meta) error { @@ -328,7 +302,19 @@ func storeIndexData(catalogDir string, metas <-chan *declcfg.Meta) error { enc := json.NewEncoder(f) enc.SetEscapeHTML(false) - return enc.Encode(idx) + if err := enc.Encode(idx); err != nil { + return err + } + return f.Sync() +} + +func syncDir(dir string) error { + d, err := os.Open(dir) + if err != nil { + return err + } + defer d.Close() + return d.Sync() } func discoverAndStoreSchema(catalogDir string, metas <-chan *declcfg.Meta) error { diff --git a/internal/catalogd/storage/storage.go b/internal/catalogd/storage/storage.go index a16a7f4cad..af78a669fc 100644 --- a/internal/catalogd/storage/storage.go +++ b/internal/catalogd/storage/storage.go @@ -15,12 +15,6 @@ type Instance interface { Delete(catalog string) error ContentExists(catalog string) bool - // VerifyAndSync confirms that catalog.jsonl exists on disk for the given - // catalog, flushes it to stable storage via fsync, and verifies the file - // is readable. It must be called after Store and before marking the - // catalog as Serving. - VerifyAndSync(catalog string) error - BaseURL(catalog string) string StorageServerHandler() http.Handler } diff --git a/internal/testutil/mock/storage/mock_instance.go b/internal/testutil/mock/storage/mock_instance.go index 245a803c9d..2aa9613e3b 100644 --- a/internal/testutil/mock/storage/mock_instance.go +++ b/internal/testutil/mock/storage/mock_instance.go @@ -111,17 +111,3 @@ func (mr *MockInstanceMockRecorder) Store(ctx, catalog, fsys any) *gomock.Call { mr.mock.ctrl.T.Helper() return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Store", reflect.TypeOf((*MockInstance)(nil).Store), ctx, catalog, fsys) } - -// VerifyAndSync mocks base method. -func (m *MockInstance) VerifyAndSync(catalog string) error { - m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "VerifyAndSync", catalog) - ret0, _ := ret[0].(error) - return ret0 -} - -// VerifyAndSync indicates an expected call of VerifyAndSync. -func (mr *MockInstanceMockRecorder) VerifyAndSync(catalog any) *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "VerifyAndSync", reflect.TypeOf((*MockInstance)(nil).VerifyAndSync), catalog) -} From dd0fde3626928dfd85395db966b8bb95cd818590 Mon Sep 17 00:00:00 2001 From: dis Date: Mon, 21 Sep 2026 15:56:12 +0530 Subject: [PATCH 3/3] add comment about syncDir --- internal/catalogd/storage/localdir.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/internal/catalogd/storage/localdir.go b/internal/catalogd/storage/localdir.go index 777fc5bddf..6d2ce2100f 100644 --- a/internal/catalogd/storage/localdir.go +++ b/internal/catalogd/storage/localdir.go @@ -179,6 +179,8 @@ func (s *LocalDirV1) storeAtomicSwap(ctx context.Context, catalog string, fsys f return "", err } + // catalog.jsonl and index.json are fsync'd when written; syncing catalogDir and RootDir + // persists directory metadata for the RemoveAll/Rename swap (file fsync alone does not). if err := syncDir(catalogDir); err != nil { return "", fmt.Errorf("error syncing catalog directory: %w", err) }