From df7a8159cfc9c5ffc8c89abf917ab9c1b908657d Mon Sep 17 00:00:00 2001 From: dongzhenyangofficial-ctrl Date: Mon, 7 Sep 2026 15:53:53 +0800 Subject: [PATCH] fix(cache): bound repository scans and preserve storage failures --- internal/dao/cache_admin_benchmark_test.go | 88 ++++++++++++++++ internal/dao/cache_admin_dao.go | 14 ++- internal/dao/cache_admin_dao_test.go | 20 ++++ internal/dao/file_dao.go | 20 +++- internal/dao/file_dao_storage_error_test.go | 50 +++++++++ internal/dao/meta_dao.go | 22 +++- internal/dao/meta_dao_degraded_test.go | 80 ++++++++++++++ internal/dao/upload_dao.go | 110 ++++++++------------ internal/dao/upload_dao_test.go | 36 +++++++ internal/server/upload_cleanup.go | 4 +- pkg/util/repo_util.go | 85 +++++++++++++++ pkg/util/repo_util_test.go | 62 +++++++++++ 12 files changed, 520 insertions(+), 71 deletions(-) create mode 100644 internal/dao/cache_admin_benchmark_test.go create mode 100644 internal/dao/file_dao_storage_error_test.go create mode 100644 internal/dao/meta_dao_degraded_test.go create mode 100644 pkg/util/repo_util_test.go diff --git a/internal/dao/cache_admin_benchmark_test.go b/internal/dao/cache_admin_benchmark_test.go new file mode 100644 index 0000000..948e665 --- /dev/null +++ b/internal/dao/cache_admin_benchmark_test.go @@ -0,0 +1,88 @@ +package dao + +import ( + "io/fs" + "os" + "path/filepath" + "sort" + "strconv" + "strings" + "testing" + + "dingospeed/internal/data" + "dingospeed/pkg/config" +) + +// BenchmarkListFilesExactRepo compares the current exact-key path with the old +// call chain, which first discovered every repository using revision as the only +// API marker. The fixture is intentionally small: seven repositories and 500 +// empty paths-info entries per repository. +func BenchmarkListFilesExactRepo(b *testing.B) { + repos := b.TempDir() + oldConfig := config.SysConfig + config.SysConfig = &config.Config{ + Server: config.ServerConfig{Repos: repos}, + Upload: config.Upload{Namespace: "dingo-local"}, + } + b.Cleanup(func() { config.SysConfig = oldConfig }) + + baseData := data.NewBaseData() + lockDao := NewLockDao(baseData) + admin := NewCacheAdminDao(NewFileDao(nil, baseData, lockDao)) + + for repo := 0; repo < 7; repo++ { + orgRepo := "remote/repo-" + string(rune('a'+repo)) + if repo == 0 { + orgRepo = "dingo-local/target" + } + if err := os.MkdirAll(filepath.Join(repoApiRoot("models", orgRepo), "revision", "main"), 0o755); err != nil { + b.Fatal(err) + } + for file := 0; file < 500; file++ { + path := filepath.Join(repoApiRoot("models", orgRepo), "paths-info", "commit", "dir", strconv.Itoa(file)) + if err := os.MkdirAll(path, 0o755); err != nil { + b.Fatal(err) + } + } + } + + b.Run("direct", func(b *testing.B) { + for i := 0; i < b.N; i++ { + _ = admin.ListFiles("models", "dingo-local/target") + } + }) + b.Run("legacy-global-discovery", func(b *testing.B) { + for i := 0; i < b.N; i++ { + _ = legacyListFilesForBenchmark("models", "dingo-local/target") + } + }) +} + +func legacyListFilesForBenchmark(repoType, orgRepo string) []*CacheFileRow { + keys := make(map[repoKey]struct{}) + root := filepath.Join(config.SysConfig.Repos(), "api", repoType) + _ = filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { + if err != nil || !entry.IsDir() || path == root { + return nil + } + if entry.Name() != "revision" { + return nil + } + rel, relErr := filepath.Rel(root, filepath.Dir(path)) + if relErr == nil && rel != "." && !strings.HasPrefix(rel, "..") { + keys[repoKey{RepoType: repoType, OrgRepo: filepath.ToSlash(rel)}] = struct{}{} + } + return fs.SkipDir + }) + ordered := make([]repoKey, 0, len(keys)) + for key := range keys { + ordered = append(ordered, key) + } + sort.Slice(ordered, func(i, j int) bool { return ordered[i].OrgRepo < ordered[j].OrgRepo }) + for _, key := range ordered { + if key.OrgRepo == orgRepo { + return indexRows(buildRepoIndex(key.RepoType, key.OrgRepo)) + } + } + return []*CacheFileRow{} +} diff --git a/internal/dao/cache_admin_dao.go b/internal/dao/cache_admin_dao.go index 796e664..2570f70 100644 --- a/internal/dao/cache_admin_dao.go +++ b/internal/dao/cache_admin_dao.go @@ -490,7 +490,9 @@ func listRepoKeys() []repoKey { repos := config.SysConfig.Repos() for _, repoType := range []string{"models", "datasets", "spaces"} { scanRepoRoots(filepath.Join(repos, "files", repoType), repoType, []string{"blobs", "resolve"}, keys) - scanRepoRoots(filepath.Join(repos, "api", repoType), repoType, []string{"revision"}, keys) + // paths-info 和 recycle 都是仓库根下的数据子树。它们本身足以证明仓库 + // 存在;命中后必须停止向下遍历,否则仓库发现会退化为扫描全部文件。 + scanRepoRoots(filepath.Join(repos, "api", repoType), repoType, []string{"revision", "paths-info", recycleDirName}, keys) } result := make([]repoKey, 0, len(keys)) for k := range keys { @@ -561,6 +563,16 @@ func (d *CacheAdminDao) ListRepos() []*CacheRepo { // ListFiles 返回某个仓库的一级列表;orgRepo 为空时返回全部仓库的合集。 func (d *CacheAdminDao) ListFiles(repoType, orgRepo string) []*CacheFileRow { + // 调用方已经给出完整仓库键时直接构建该仓库的索引。不能为了确认它是否在 + // 列表中先枚举所有仓库;一个小仓库的延迟不应受其他仓库大小影响。 + if repoType != "" && orgRepo != "" { + key := repoKey{RepoType: repoType, OrgRepo: orgRepo} + if validateRepoKey(key) != nil { + return []*CacheFileRow{} + } + return indexRows(buildRepoIndex(repoType, orgRepo)) + } + rows := make([]*CacheFileRow, 0) for _, key := range listRepoKeys() { if repoType != "" && key.RepoType != repoType { diff --git a/internal/dao/cache_admin_dao_test.go b/internal/dao/cache_admin_dao_test.go index cefdc67..edf8e41 100644 --- a/internal/dao/cache_admin_dao_test.go +++ b/internal/dao/cache_admin_dao_test.go @@ -115,6 +115,26 @@ func findOrphan(rows []*RecycleRow, sha string) *RecycleRow { return nil } +func TestListRepoKeysStopsAtPathsInfo(t *testing.T) { + _, _, _ = newTestCacheAdminDao(t) + orgRepo := "dingo-local/paths-heavy" + + // revision 是一个合法仓库标记,但这里故意把同名目录放在 paths-info + // 深处。仓库发现如果进入数据子树,会把这个深层路径误判成另一个仓库。 + deepRevision := filepath.Join(repoApiRoot("models", orgRepo), "paths-info", "commit", "nested", "revision") + if err := os.MkdirAll(deepRevision, 0o755); err != nil { + t.Fatalf("create deep paths-info fixture: %v", err) + } + + keys := listRepoKeys() + if len(keys) != 1 { + t.Fatalf("expected one repository without descending paths-info, got %#v", keys) + } + if keys[0] != (repoKey{RepoType: "models", OrgRepo: orgRepo}) { + t.Fatalf("unexpected repository key: %#v", keys[0]) + } +} + func tombstoneExists(repoType, orgRepo, sha string) bool { return util.FileExists(recycleEntryPath(repoType, orgRepo, sha)) } diff --git a/internal/dao/file_dao.go b/internal/dao/file_dao.go index 61fcf83..65e5552 100644 --- a/internal/dao/file_dao.go +++ b/internal/dao/file_dao.go @@ -18,6 +18,7 @@ import ( "context" "crypto/sha256" "encoding/hex" + "errors" "fmt" "io" "net/http" @@ -82,6 +83,11 @@ func (f *FileDao) GetFileCommitSha(repoType, orgRepo, commit, authorization stri if IsLocalOrgRepo(orgRepo) { commitSha, err = f.GetCommitHfOffline(repoType, orgRepo, commit) if err != nil { + var accessErr *util.FileAccessError + if errors.As(err, &accessErr) { + zap.S().Errorw("local metadata storage unavailable", "kind", accessErr.Kind, "path", accessErr.Path, "error", accessErr.Err) + return "", myerr.Wrap("local metadata storage unavailable", err) + } return "", myerr.NewAppendCode(http.StatusNotFound, fmt.Sprintf("%s is not found", orgRepo)) } return commitSha, nil @@ -91,6 +97,11 @@ func (f *FileDao) GetFileCommitSha(repoType, orgRepo, commit, authorization stri } commitSha, err = f.GetCommitHfOffline(repoType, orgRepo, commit) if err != nil { + var accessErr *util.FileAccessError + if errors.As(err, &accessErr) { + zap.S().Errorw("metadata cache storage unavailable", "kind", accessErr.Kind, "path", accessErr.Path, "error", accessErr.Err) + return "", myerr.Wrap("metadata cache storage unavailable", err) + } if source == "file" { // 若只是发起文件下载(先在线后离线),将不会校验meta文件是否存在,没有就创建,主要是看文件本身是否存在。 goto remoteRequestMeta @@ -167,9 +178,16 @@ func (f *FileDao) RemoteRequestMeta(method, repoType, orgRepo, revision, authori func (f *FileDao) GetCommitHfOffline(repoType, orgRepo, commit string) (string, error) { apiPath := fmt.Sprintf("%s/api/%s/%s/revision/%s/meta_get.json", config.SysConfig.Repos(), repoType, orgRepo, commit) - if util.FileExists(apiPath) { + exists, err := util.PathExists(apiPath) + if err != nil { + return "", err + } + if exists { cacheContent, err := f.ReadCacheRequest(apiPath) if err != nil { + if accessErr, ok := util.ClassifyFileAccessError(apiPath, err); ok { + return "", accessErr + } return "", err } var sha CommitHfSha diff --git a/internal/dao/file_dao_storage_error_test.go b/internal/dao/file_dao_storage_error_test.go new file mode 100644 index 0000000..ecc6899 --- /dev/null +++ b/internal/dao/file_dao_storage_error_test.go @@ -0,0 +1,50 @@ +package dao + +import ( + "errors" + "net/http" + "strings" + "testing" + + "dingospeed/internal/data" + "dingospeed/pkg/config" + myerr "dingospeed/pkg/error" + "dingospeed/pkg/util" +) + +func TestGetCommitHfOfflinePreservesMissingBehavior(t *testing.T) { + oldConfig := config.SysConfig + config.SysConfig = &config.Config{Server: config.ServerConfig{Repos: t.TempDir()}} + t.Cleanup(func() { config.SysConfig = oldConfig }) + + fileDao := NewFileDao(nil, data.NewBaseData(), nil) + _, err := fileDao.GetFileCommitSha("models", "org/repo", "main", "", "meta") + var appErr myerr.Error + if !errors.As(err, &appErr) || appErr.StatusCode() != http.StatusNotFound { + t.Fatalf("missing metadata: got %v, want HTTP 404 application error", err) + } +} + +func TestGetFileCommitShaPreservesStorageFailure(t *testing.T) { + oldConfig := config.SysConfig + // A NUL byte makes stat fail as an invalid path on every supported OS. It + // exercises the inaccessible-path branch without relying on host mounts. + config.SysConfig = &config.Config{Server: config.ServerConfig{Repos: "invalid\x00repos"}} + t.Cleanup(func() { config.SysConfig = oldConfig }) + + fileDao := NewFileDao(nil, data.NewBaseData(), nil) + _, err := fileDao.GetFileCommitSha("models", "org/repo", "main", "", "meta") + if err == nil { + t.Fatal("expected inaccessible storage error") + } + var accessErr *util.FileAccessError + if !errors.As(err, &accessErr) { + t.Fatalf("got %T %v, want wrapped FileAccessError", err, err) + } + if accessErr.Kind != util.FileAccessUnavailable { + t.Fatalf("got kind %q, want %q", accessErr.Kind, util.FileAccessUnavailable) + } + if strings.Contains(err.Error(), "invalid\x00repos") { + t.Fatalf("public application error exposes storage path: %q", err.Error()) + } +} diff --git a/internal/dao/meta_dao.go b/internal/dao/meta_dao.go index 9f6b874..ef4eec2 100644 --- a/internal/dao/meta_dao.go +++ b/internal/dao/meta_dao.go @@ -15,6 +15,7 @@ package dao import ( + "errors" "fmt" "net/http" "path/filepath" @@ -191,7 +192,7 @@ func (m *MetaDao) requestAndSaveMeta(repoType, orgRepo, revision, commitSha, met if revision == mainVersion { err = m.writeApiMetaFile(repoType, orgRepo, revision, method, resp.StatusCode, extractHeaders, resp.Body) if err != nil { - return nil, err + m.logMetadataCacheWriteFailure(orgRepo, revision, method, err) } } else { apiDir := fmt.Sprintf("%s/api/%s/%s/revision/%s", config.SysConfig.Repos(), repoType, orgRepo, mainVersion) @@ -199,14 +200,14 @@ func (m *MetaDao) requestAndSaveMeta(repoType, orgRepo, revision, commitSha, met if !util.FileExists(apiMetaPath) { err = m.writeApiMetaFile(repoType, orgRepo, mainVersion, method, resp.StatusCode, extractHeaders, resp.Body) // create main dir if err != nil { - return nil, err + m.logMetadataCacheWriteFailure(orgRepo, mainVersion, method, err) } } } err = m.writeApiMetaFile(repoType, orgRepo, commitSha, method, resp.StatusCode, extractHeaders, resp.Body) if err != nil { - return nil, err + m.logMetadataCacheWriteFailure(orgRepo, commitSha, method, err) } return &common.CacheContent{ StatusCode: resp.StatusCode, @@ -215,16 +216,31 @@ func (m *MetaDao) requestAndSaveMeta(repoType, orgRepo, revision, commitSha, met }, nil } +func (m *MetaDao) logMetadataCacheWriteFailure(orgRepo, revision, method string, err error) { + fields := []interface{}{"repo", orgRepo, "revision", revision, "method", method, "error", err} + var accessErr *util.FileAccessError + if errors.As(err, &accessErr) { + fields = append(fields, "kind", accessErr.Kind, "path", accessErr.Path) + } + zap.S().Warnw("serving remote metadata without cache", fields...) +} + func (m *MetaDao) writeApiMetaFile(repoType, orgRepo, commitSha, method string, statusCode int, extractHeaders map[string]string, body []byte) error { apiDir := fmt.Sprintf("%s/api/%s/%s/revision/%s", config.SysConfig.Repos(), repoType, orgRepo, commitSha) apiMetaPath := fmt.Sprintf("%s/%s", apiDir, fmt.Sprintf("meta_%s.json", method)) err := util.MakeDirs(apiMetaPath) if err != nil { zap.S().Errorf("create %s dir err.%v", apiMetaPath, err) + if accessErr, ok := util.ClassifyFileAccessError(apiMetaPath, err); ok { + return accessErr + } return err } if err = m.fileDao.WriteCacheRequest(apiMetaPath, statusCode, extractHeaders, body); err != nil { zap.S().Errorf("writeCacheRequest err.%v", err) + if accessErr, ok := util.ClassifyFileAccessError(apiMetaPath, err); ok { + return accessErr + } return err } return nil diff --git a/internal/dao/meta_dao_degraded_test.go b/internal/dao/meta_dao_degraded_test.go new file mode 100644 index 0000000..b92439f --- /dev/null +++ b/internal/dao/meta_dao_degraded_test.go @@ -0,0 +1,80 @@ +package dao + +import ( + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + + "dingospeed/internal/data" + "dingospeed/pkg/config" + "dingospeed/pkg/consts" +) + +func TestRequestAndSaveMetaServesRemoteResponseWhenCacheUnavailable(t *testing.T) { + const body = `{"sha":"remote-commit","id":"org/repo"}` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(body)) + })) + t.Cleanup(server.Close) + + oldConfig := config.SysConfig + config.SysConfig = &config.Config{ + Server: config.ServerConfig{ + Online: true, + Repos: "invalid\x00repos", + HfScheme: "http", + HfNetLoc: strings.TrimPrefix(server.URL, "http://"), + }, + Retry: config.Retry{Attempts: 1}, + } + t.Cleanup(func() { config.SysConfig = oldConfig }) + + baseData := data.NewBaseData() + fileDao := NewFileDao(nil, baseData, NewLockDao(baseData)) + metaDao := NewMetaDao(fileDao, nil, baseData) + got, err := metaDao.requestAndSaveMeta("models", "org/repo", "main", "remote-commit", consts.RequestTypeGet, "") + if err != nil { + t.Fatalf("remote metadata should survive cache failure: %v", err) + } + if string(got.OriginContent) != body { + t.Fatalf("got body %q, want %q", got.OriginContent, body) + } +} + +func TestRequestAndSaveMetaStillCachesWhenStorageAvailable(t *testing.T) { + const body = `{"sha":"remote-commit","id":"org/repo"}` + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte(body)) + })) + t.Cleanup(server.Close) + + repos := t.TempDir() + oldConfig := config.SysConfig + config.SysConfig = &config.Config{ + Server: config.ServerConfig{ + Online: true, + Repos: repos, + HfScheme: "http", + HfNetLoc: strings.TrimPrefix(server.URL, "http://"), + }, + Retry: config.Retry{Attempts: 1}, + } + t.Cleanup(func() { config.SysConfig = oldConfig }) + + baseData := data.NewBaseData() + fileDao := NewFileDao(nil, baseData, NewLockDao(baseData)) + metaDao := NewMetaDao(fileDao, nil, baseData) + if _, err := metaDao.requestAndSaveMeta("models", "org/repo", "main", "remote-commit", consts.RequestTypeGet, ""); err != nil { + t.Fatal(err) + } + for _, revision := range []string{"main", "remote-commit"} { + path := filepath.Join(repos, "api", "models", "org", "repo", "revision", revision, "meta_get.json") + if _, err := os.Stat(path); err != nil { + t.Fatalf("expected metadata cache %s: %v", path, err) + } + } +} diff --git a/internal/dao/upload_dao.go b/internal/dao/upload_dao.go index d11f211..31cf709 100644 --- a/internal/dao/upload_dao.go +++ b/internal/dao/upload_dao.go @@ -956,56 +956,63 @@ func (u *UploadDao) CleanupExpiredStagedUploads(retention time.Duration) (int, e if retention <= 0 { retention = 7 * 24 * time.Hour } - root := filepath.Join(config.SysConfig.Repos(), "files") cutoff := time.Now().Add(-retention) removed := 0 - err := filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { - if err != nil { - if os.IsNotExist(err) { + namespace := "dingo-local" + if config.SysConfig != nil && config.SysConfig.Upload.Namespace != "" { + namespace = config.SysConfig.Upload.Namespace + } + filesRoot := filepath.Join(config.SysConfig.Repos(), "files") + for _, repoType := range []string{"models", "datasets"} { + root := filepath.Join(filesRoot, repoType, namespace) + err := filepath.WalkDir(root, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if entry.IsDir() || !strings.HasSuffix(entry.Name(), localUploadStageSuffix) { return nil } - return err - } - if entry.IsDir() || !strings.HasSuffix(entry.Name(), localUploadStageSuffix) { - return nil - } - repoType, orgRepo, sha, ok := stagedBlobIdentity(root, path) - if !ok { - return nil - } - info, statErr := entry.Info() - if statErr != nil { - if os.IsNotExist(statErr) { + parsedRepoType, orgRepo, sha, ok := stagedBlobIdentity(filesRoot, path) + if !ok || parsedRepoType != repoType || !IsLocalOrgRepo(orgRepo) { return nil } - return statErr - } - if info.ModTime().After(cutoff) { - return nil - } - blobKey := uploadBlobLockKey(repoType, orgRepo, sha) - uploadBlobLocks.Lock(blobKey) - defer uploadBlobLocks.Unlock(blobKey) - info, statErr = os.Stat(path) - if statErr != nil { - if os.IsNotExist(statErr) { + info, statErr := entry.Info() + if statErr != nil { + if os.IsNotExist(statErr) { + return nil + } + return statErr + } + if info.ModTime().After(cutoff) { return nil } - return statErr - } - if info.ModTime().After(cutoff) { + blobKey := uploadBlobLockKey(repoType, orgRepo, sha) + uploadBlobLocks.Lock(blobKey) + defer uploadBlobLocks.Unlock(blobKey) + info, statErr = os.Stat(path) + if statErr != nil { + if os.IsNotExist(statErr) { + return nil + } + return statErr + } + if info.ModTime().After(cutoff) { + return nil + } + if rmErr := os.Remove(path); rmErr != nil && !os.IsNotExist(rmErr) { + return rmErr + } + removed++ return nil + }) + if err != nil && !os.IsNotExist(err) { + return removed, err } - if rmErr := os.Remove(path); rmErr != nil && !os.IsNotExist(rmErr) { - return rmErr - } - removed++ - return nil - }) - if os.IsNotExist(err) { - return removed, nil } - return removed, err + return removed, nil } // CleanupUnreferencedBlobs 回收“已经完整落盘、但不被任何清单引用”的内容。 @@ -1239,31 +1246,6 @@ func (u *UploadDao) RunStagedUploadCleanup(ctx context.Context) { } else if removed > 0 { zap.S().Infof("cleanup expired staged uploads removed %d file(s)", removed) } - // 待回收快照要排在 blob 回收之前:快照还在盘上时它的清单仍然算引用 - // (referencedShas 扫的是全部快照),里面的内容永远轮不到回收。 - dropped, err := u.CleanupSupersededSnapshots(config.SysConfig.GetUploadSupersededRetention()) - if err != nil { - zap.S().Warnf("drop superseded snapshots failed: %v", err) - } else if dropped > 0 { - zap.S().Infof("dropped %d superseded snapshot(s)", dropped) - } - retention := config.SysConfig.GetUploadOrphanRetention() - reclaimed, err := u.CleanupUnreferencedBlobs(retention) - if err != nil { - zap.S().Warnf("reclaim unreferenced upload content failed: %v", err) - } else if reclaimed > 0 { - zap.S().Infof("reclaimed %d unreferenced upload content file(s)", reclaimed) - } - // 上一趟只覆盖本地命名空间。远端缓存里被一级删除的内容由墓碑驱动, - // 与上传内容共用同一个保留期(orphanRetentionHours,默认 168h)。 - recycled, err := u.CleanupRecycledBlobs(retention) - if err != nil { - zap.S().Warnf("reclaim recycled cache content failed: %v", err) - continue - } - if recycled > 0 { - zap.S().Infof("reclaimed %d recycled cache content file(s)", recycled) - } } } } diff --git a/internal/dao/upload_dao_test.go b/internal/dao/upload_dao_test.go index d60c3ed..f258345 100644 --- a/internal/dao/upload_dao_test.go +++ b/internal/dao/upload_dao_test.go @@ -471,6 +471,42 @@ func TestStagedUploadRetentionAndCleanup(t *testing.T) { } } +func TestStagedUploadCleanupOnlyScansLocalNamespace(t *testing.T) { + u, repos := newTestUploadDao(t) + stale := time.Now().Add(-48 * time.Hour) + + localPath := stagedBlobPath(repos, "models", "dingo-local/local-repo", "local-sha") + remotePath := stagedBlobPath(repos, "models", "remote-org/remote-repo", "remote-sha") + localDatasetPath := stagedBlobPath(repos, "datasets", "dingo-local/local-dataset", "dataset-sha") + for _, path := range []string{localPath, remotePath, localDatasetPath} { + if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, []byte("staged"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Chtimes(path, stale, stale); err != nil { + t.Fatal(err) + } + } + + removed, err := u.CleanupExpiredStagedUploads(24 * time.Hour) + if err != nil { + t.Fatal(err) + } + if removed != 2 { + t.Fatalf("removed=%d, want 2 local staged files", removed) + } + for _, path := range []string{localPath, localDatasetPath} { + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Fatalf("local staged file must be removed: %s (err=%v)", path, err) + } + } + if _, err := os.Stat(remotePath); err != nil { + t.Fatalf("remote namespace must not be scanned or removed: %v", err) + } +} + func TestRepeatedFullUploadInterruptionsReuseOneStagedFile(t *testing.T) { u, repos := newTestUploadDao(t) content := []byte("0123456789abcdef0123456789abcdef-tail") diff --git a/internal/server/upload_cleanup.go b/internal/server/upload_cleanup.go index 10c08b4..8badd95 100644 --- a/internal/server/upload_cleanup.go +++ b/internal/server/upload_cleanup.go @@ -18,8 +18,8 @@ func NewUploadCleanupServer(uploadDao *dao.UploadDao) *UploadCleanupServer { } func (s *UploadCleanupServer) Start(ctx context.Context) error { - zap.S().Infof("[UPLOAD-CLEANUP] staged upload cleanup interval=%s retention=%s", - config.SysConfig.GetUploadStagingCleanupInterval(), config.SysConfig.GetUploadStagingRetention()) + zap.S().Infof("[UPLOAD-CLEANUP] local namespace=%s staged upload cleanup interval=%s retention=%s", + config.SysConfig.Upload.Namespace, config.SysConfig.GetUploadStagingCleanupInterval(), config.SysConfig.GetUploadStagingRetention()) go s.uploadDao.RunStagedUploadCleanup(ctx) return nil } diff --git a/pkg/util/repo_util.go b/pkg/util/repo_util.go index 2054761..65f554a 100644 --- a/pkg/util/repo_util.go +++ b/pkg/util/repo_util.go @@ -21,9 +21,11 @@ import ( "io" "os" "path/filepath" + "runtime" "sort" "strings" "sync" + "syscall" "time" "dingospeed/pkg/common" @@ -36,6 +38,89 @@ var ( symlinkLock sync.Mutex ) +type FileAccessErrorKind string + +const ( + FileAccessPermissionDenied FileAccessErrorKind = "permission_denied" + FileAccessMountDisconnected FileAccessErrorKind = "mount_disconnected" + FileAccessIOFailure FileAccessErrorKind = "io_failure" + FileAccessUnavailable FileAccessErrorKind = "path_unavailable" +) + +// FileAccessError distinguishes an inaccessible storage path from a path that +// simply does not exist. Err retains the operating-system error for diagnosis. +type FileAccessError struct { + Kind FileAccessErrorKind + Path string + Err error +} + +func (e *FileAccessError) Error() string { + return fmt.Sprintf("storage path %s is unavailable (%s): %v", e.Path, e.Kind, e.Err) +} + +func (e *FileAccessError) Unwrap() error { return e.Err } + +// PathExists is the error-preserving alternative to FileExists. A missing path +// is a normal negative result; every other stat failure is returned as a typed +// FileAccessError so callers cannot mistake a storage failure for a cache miss. +func PathExists(filePath string) (bool, error) { + _, err := os.Stat(filePath) + if err == nil { + return true, nil + } + if os.IsNotExist(err) { + return false, nil + } + accessErr, _ := ClassifyFileAccessError(filePath, err) + return false, accessErr +} + +// ClassifyFileAccessError recognizes filesystem access failures even when +// another layer has wrapped the original *os.PathError. Content and parsing +// errors are deliberately left unclassified. +func ClassifyFileAccessError(filePath string, err error) (*FileAccessError, bool) { + if err == nil || os.IsNotExist(err) { + return nil, false + } + var pathErr *os.PathError + if !errors.As(err, &pathErr) { + return nil, false + } + return &FileAccessError{Kind: classifyFileAccessError(err), Path: filePath, Err: err}, true +} + +func classifyFileAccessError(err error) FileAccessErrorKind { + return classifyFileAccessErrorForOS(err, runtime.GOOS) +} + +func classifyFileAccessErrorForOS(err error, goos string) FileAccessErrorKind { + if goos != "linux" && os.IsPermission(err) { + return FileAccessPermissionDenied + } + + // DingoFS runs on Linux. Numeric errno values keep this package buildable + // on development hosts where Linux-specific syscall names are unavailable. + const ( + linuxEIO syscall.Errno = 5 + linuxENOTCONN syscall.Errno = 107 + linuxESTALE syscall.Errno = 116 + ) + var errno syscall.Errno + if errors.As(err, &errno) { + switch errno { + case linuxENOTCONN, linuxESTALE: + return FileAccessMountDisconnected + case linuxEIO: + return FileAccessIOFailure + } + } + if os.IsPermission(err) { + return FileAccessPermissionDenied + } + return FileAccessUnavailable +} + func GetOrgRepo(org, repo string) string { if org == "" { return repo diff --git a/pkg/util/repo_util_test.go b/pkg/util/repo_util_test.go new file mode 100644 index 0000000..c4a7fc7 --- /dev/null +++ b/pkg/util/repo_util_test.go @@ -0,0 +1,62 @@ +package util + +import ( + "errors" + "fmt" + "os" + "path/filepath" + "syscall" + "testing" +) + +func TestPathExistsResults(t *testing.T) { + dir := t.TempDir() + exists, err := PathExists(filepath.Join(dir, "missing")) + if err != nil || exists { + t.Fatalf("missing path: exists=%v err=%v", exists, err) + } + + path := filepath.Join(dir, "present") + if err := os.WriteFile(path, []byte("ok"), 0o600); err != nil { + t.Fatal(err) + } + exists, err = PathExists(path) + if err != nil || !exists { + t.Fatalf("existing path: exists=%v err=%v", exists, err) + } +} + +func TestClassifyWrappedFileAccessError(t *testing.T) { + original := &os.PathError{Op: "read", Path: "/repos/meta.json", Err: syscall.Errno(107)} + accessErr, ok := ClassifyFileAccessError("/repos/meta.json", fmt.Errorf("cache read: %w", original)) + if !ok { + t.Fatal("expected wrapped filesystem error to be classified") + } + if accessErr.Kind != FileAccessMountDisconnected { + t.Fatalf("got kind %q, want %q", accessErr.Kind, FileAccessMountDisconnected) + } + if _, ok = ClassifyFileAccessError("/repos/meta.json", errors.New("invalid metadata")); ok { + t.Fatal("content error must not be classified as a storage access failure") + } +} + +func TestClassifyFileAccessError(t *testing.T) { + tests := []struct { + name string + err error + want FileAccessErrorKind + }{ + {name: "permission", err: os.ErrPermission, want: FileAccessPermissionDenied}, + {name: "disconnected", err: &os.PathError{Op: "stat", Path: "/repos", Err: syscall.Errno(107)}, want: FileAccessMountDisconnected}, + {name: "stale", err: &os.PathError{Op: "stat", Path: "/repos", Err: syscall.Errno(116)}, want: FileAccessMountDisconnected}, + {name: "io", err: &os.PathError{Op: "stat", Path: "/repos", Err: syscall.Errno(5)}, want: FileAccessIOFailure}, + {name: "other", err: errors.New("other"), want: FileAccessUnavailable}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := classifyFileAccessErrorForOS(tc.err, "linux"); got != tc.want { + t.Fatalf("got %q, want %q", got, tc.want) + } + }) + } +}