From 6079ff76c1cc13cedd4848b5b9060ae1a38cad0c Mon Sep 17 00:00:00 2001 From: Chris J Arges Date: Mon, 3 Aug 2026 21:54:02 -0500 Subject: [PATCH] s3: batch published pool cleanup Use S3 multi-object deletion for orphaned published pool files, while retaining per-file deletion for other storage backends. Signed-off-by: Chris J Arges --- AUTHORS | 1 + aptly/interfaces.go | 6 ++ deb/publish.go | 14 ++++- deb/publish_test.go | 49 +++++++++++++++ s3/public.go | 78 ++++++++++++++++++++++- s3/public_test.go | 150 ++++++++++++++++++++++++++++++++++++++++++++ s3/server_test.go | 67 +++++++++++++++++++- 7 files changed, 357 insertions(+), 8 deletions(-) diff --git a/AUTHORS b/AUTHORS index ea69b2987..6b20a4537 100644 --- a/AUTHORS +++ b/AUTHORS @@ -84,3 +84,4 @@ List of contributors, in chronological order: * Zhang Xiao (https://github.com/xzhang1) * Tom Nguyen (https://github.com/lecafard) * Philip Cramer (https://github.com/PhilipCramer) +* Chris J Arges (https://github.com/arges) diff --git a/aptly/interfaces.go b/aptly/interfaces.go index 1412c719b..278099827 100644 --- a/aptly/interfaces.go +++ b/aptly/interfaces.go @@ -87,6 +87,12 @@ type PublishedStorage interface { ReadLink(path string) (string, error) } +// PublishedStorageBulkRemover optionally removes multiple files efficiently. +type PublishedStorageBulkRemover interface { + // RemoveFiles removes multiple files under public path. + RemoveFiles(paths []string) error +} + // FileSystemPublishedStorage is published storage on filesystem type FileSystemPublishedStorage interface { // PublicPath returns root of public part diff --git a/deb/publish.go b/deb/publish.go index 8ae71df94..258788dd7 100644 --- a/deb/publish.go +++ b/deb/publish.go @@ -1769,12 +1769,20 @@ func (collection *PublishedRepoCollection) CleanupPrefixComponentFiles(published sort.Strings(existingFiles) orphanedFiles := utils.StrSlicesSubstract(existingFiles, referencedFiles[component]) + for i := range orphanedFiles { + orphanedFiles[i] = filepath.Join(path, orphanedFiles[i]) + } - for _, file := range orphanedFiles { - err = publishedStorage.Remove(filepath.Join(path, file)) - if err != nil { + if bulkRemover, ok := publishedStorage.(aptly.PublishedStorageBulkRemover); ok { + if err = bulkRemover.RemoveFiles(orphanedFiles); err != nil { return err } + } else { + for _, file := range orphanedFiles { + if err = publishedStorage.Remove(file); err != nil { + return err + } + } } } diff --git a/deb/publish_test.go b/deb/publish_test.go index 20f78c4f9..25c039b53 100644 --- a/deb/publish_test.go +++ b/deb/publish_test.go @@ -70,6 +70,26 @@ func (p *FakeStorageProvider) GetPublishedStorage(name string) (aptly.PublishedS return storage, nil } +type recordingBulkPublishedStorage struct { + aptly.PublishedStorage + removed []string + err error +} + +func (s *recordingBulkPublishedStorage) RemoveFiles(paths []string) error { + s.removed = append(s.removed, paths...) + return s.err +} + +type failingRemovePublishedStorage struct { + aptly.PublishedStorage + err error +} + +func (s *failingRemovePublishedStorage) Remove(string) error { + return s.err +} + type PublishedRepoSuite struct { PackageListMixinSuite repo, repo2, repo3, repo4, repo5 *PublishedRepo @@ -1028,6 +1048,35 @@ func (s *PublishedRepoRemoveSuite) TestRemoveFilesWithPrefixRoot(c *C) { c.Check(filepath.Join(s.publishedStorage2.PublicPath(), "ppa/pool/contrib"), PathExists) } +func (s *PublishedRepoRemoveSuite) TestCleanupPrefixComponentFilesRemovalRouting(c *C) { + orphan := filepath.Join("ppa", "pool", "main", "orphan.deb") + orphanPath := filepath.Join(s.publishedStorage.PublicPath(), orphan) + c.Assert(os.WriteFile(orphanPath, []byte("orphan"), 0644), IsNil) + + err := s.collection.CleanupPrefixComponentFiles(s.provider, s.repo1, []string{"main"}, s.factory, nil) + c.Check(err, IsNil) + c.Check(orphanPath, Not(PathExists)) + + c.Assert(os.WriteFile(orphanPath, []byte("orphan"), 0644), IsNil) + bulkStorage := &recordingBulkPublishedStorage{PublishedStorage: s.publishedStorage} + s.provider.storages[""] = bulkStorage + + err = s.collection.CleanupPrefixComponentFiles(s.provider, s.repo1, []string{"main"}, s.factory, nil) + c.Check(err, IsNil) + c.Check(bulkStorage.removed, DeepEquals, []string{orphan}) + + bulkStorage.err = errors.New("bulk removal failed") + err = s.collection.CleanupPrefixComponentFiles(s.provider, s.repo1, []string{"main"}, s.factory, nil) + c.Check(err, ErrorMatches, "bulk removal failed") + + s.provider.storages[""] = &failingRemovePublishedStorage{ + PublishedStorage: s.publishedStorage, + err: errors.New("single removal failed"), + } + err = s.collection.CleanupPrefixComponentFiles(s.provider, s.repo1, []string{"main"}, s.factory, nil) + c.Check(err, ErrorMatches, "single removal failed") +} + func (s *PublishedRepoRemoveSuite) TestRemoveRepo1and2(c *C) { err := s.collection.Remove(s.provider, "", "ppa", "anaconda", s.factory, nil, false, false) c.Check(err, IsNil) diff --git a/s3/public.go b/s3/public.go index c35de4250..c9c27f331 100644 --- a/s3/public.go +++ b/s3/public.go @@ -58,9 +58,13 @@ type PublishedStorage struct { encryptByDefault bool } +// Amazon S3 accepts at most 1,000 keys in a DeleteObjects request. +const maxDeleteObjectsPerRequest = 1000 + // Check interface var ( - _ aptly.PublishedStorage = (*PublishedStorage)(nil) + _ aptly.PublishedStorage = (*PublishedStorage)(nil) + _ aptly.PublishedStorageBulkRemover = (*PublishedStorage)(nil) ) // NewPublishedStorageRaw creates published storage from raw aws credentials @@ -262,7 +266,7 @@ func (storage *PublishedStorage) Remove(path string) error { // RemoveDirs removes directory structure under public path func (storage *PublishedStorage) RemoveDirs(path string, _ aptly.Progress) error { - const page = 1000 + const page = maxDeleteObjectsPerRequest filelist, _, err := storage.internalFilelist(path, false) if err != nil { @@ -330,6 +334,76 @@ func (storage *PublishedStorage) RemoveDirs(path string, _ aptly.Progress) error return nil } +// RemoveFiles removes multiple files under public path. +func (storage *PublishedStorage) RemoveFiles(paths []string) error { + if storage.disableMultiDel { + for _, path := range paths { + if err := storage.Remove(path); err != nil { + return err + } + } + return nil + } + + if storage.plusWorkaround { + expanded := make([]string, 0, len(paths)) + for _, path := range paths { + expanded = append(expanded, path) + if strings.Contains(path, "+") { + expanded = append(expanded, strings.ReplaceAll(path, "+", " ")) + } + } + paths = expanded + } + + var failures []string + for offset := 0; offset < len(paths); offset += maxDeleteObjectsPerRequest { + part := paths[offset:min(offset+maxDeleteObjectsPerRequest, len(paths))] + objects := make([]types.ObjectIdentifier, len(part)) + for i, path := range part { + objects[i] = types.ObjectIdentifier{ + Key: aws.String(filepath.Join(storage.prefix, path)), + } + } + + quiet := true + output, err := storage.s3.DeleteObjects(context.TODO(), &s3.DeleteObjectsInput{ + Bucket: aws.String(storage.bucket), + Delete: &types.Delete{ + Objects: objects, + Quiet: &quiet, + }, + }) + if err != nil { + var notFoundErr *smithy.GenericAPIError + if errors.As(err, ¬FoundErr) && notFoundErr.Code == "NoSuchBucket" { + return nil + } + return fmt.Errorf("error deleting multiple paths from %s: %s", storage, err) + } + failed := make(map[string]struct{}, len(output.Errors)) + for _, failure := range output.Errors { + key := aws.ToString(failure.Key) + failed[key] = struct{}{} + failures = append(failures, fmt.Sprintf("%s: %s: %s", key, + aws.ToString(failure.Code), aws.ToString(failure.Message))) + } + storage.pathCacheMutex.Lock() + for _, path := range part { + if _, failed := failed[filepath.Join(storage.prefix, path)]; !failed { + delete(storage.pathCache, path) + } + } + storage.pathCacheMutex.Unlock() + + } + if len(failures) != 0 { + return fmt.Errorf("errors deleting multiple paths from %s: %s", storage, strings.Join(failures, "; ")) + } + + return nil +} + // LinkFromPool links package file from pool to dist's pool location // // publishedPrefix is desired prefix for the location in the pool. diff --git a/s3/public_test.go b/s3/public_test.go index d857e64dd..33cb9dd4b 100644 --- a/s3/public_test.go +++ b/s3/public_test.go @@ -3,6 +3,7 @@ package s3 import ( "bytes" "context" + "fmt" "io" "os" "path/filepath" @@ -250,6 +251,155 @@ func (s *PublishedStorageSuite) TestRemoveDirsNoSuchBucket(c *C) { c.Check(err, ErrorMatches, ".*StatusCode: 404.*") } +func (s *PublishedStorageSuite) countRequests(method, uriSubstring string) int { + count := 0 + for _, r := range s.srv.Requests { + if r.Method == method && strings.Contains(r.RequestURI, uriSubstring) { + count++ + } + } + + return count +} + +func (s *PublishedStorageSuite) TestRemoveFilesPrefixed(c *C) { + s.prefixedStorage.disableMultiDel = false + + s.PutFile(c, "lala/xyz", []byte("test")) + s.PutFile(c, "lala/abc", []byte("test")) + + err := s.prefixedStorage.RemoveFiles([]string{"xyz"}) + c.Check(err, IsNil) + + s.AssertNoFile(c, "lala/xyz") + + list, err := s.storage.Filelist("") + c.Check(err, IsNil) + c.Check(list, DeepEquals, []string{"lala/abc"}) + c.Check(s.countRequests("POST", "delete"), Equals, 1) + c.Check(s.countRequests("DELETE", ""), Equals, 0) +} + +func (s *PublishedStorageSuite) TestRemoveFilesBatchesAtThousand(c *C) { + s.storage.disableMultiDel = false + + paths := make([]string, maxDeleteObjectsPerRequest+1) + for i := range paths { + paths[i] = fmt.Sprintf("pool/main/p/pkg/file-%d.deb", i) + } + + err := s.storage.RemoveFiles(paths) + c.Check(err, IsNil) + + c.Check(s.countRequests("POST", "delete"), Equals, 2) + c.Check(s.countRequests("DELETE", ""), Equals, 0) +} + +func (s *PublishedStorageSuite) TestRemoveFilesEmpty(c *C) { + s.storage.disableMultiDel = false + + err := s.storage.RemoveFiles(nil) + c.Check(err, IsNil) + c.Check(s.countRequests("POST", "delete"), Equals, 0) +} + +func (s *PublishedStorageSuite) TestRemoveFilesNoSuchBucket(c *C) { + s.noSuchBucketStorage.disableMultiDel = false + + err := s.noSuchBucketStorage.RemoveFiles([]string{"a"}) + c.Check(err, IsNil) +} + +func (s *PublishedStorageSuite) TestRemoveFilesRequestError(c *C) { + s.storage.disableMultiDel = false + s.srv.config.DeleteObjectsError = "multi-delete failed" + + err := s.storage.RemoveFiles([]string{"a"}) + c.Check(err, ErrorMatches, "error deleting multiple paths.*AccessDenied.*multi-delete failed.*") +} + +func (s *PublishedStorageSuite) TestRemoveFilesReportsAllFailuresAndInvalidatesSuccesses(c *C) { + s.prefixedStorage.disableMultiDel = false + s.srv.config.DeleteErrors = map[string]string{ + "lala/a": "a failed", + "lala/c": "c failed", + } + s.prefixedStorage.pathCache = map[string]string{"a": "", "b": "", "c": ""} + for _, path := range []string{"a", "b", "c"} { + s.PutFile(c, filepath.Join("lala", path), []byte("test")) + } + + err := s.prefixedStorage.RemoveFiles([]string{"a", "b", "c"}) + c.Check(err, ErrorMatches, "errors deleting multiple paths.*lala/a: AccessDenied: a failed; lala/c: AccessDenied: c failed") + + s.AssertNoFile(c, "lala/b") + c.Check(s.prefixedStorage.pathCache, DeepEquals, map[string]string{"a": "", "c": ""}) +} + +func (s *PublishedStorageSuite) TestRemoveFilesReportsFailuresAcrossBatches(c *C) { + s.storage.disableMultiDel = false + paths := make([]string, maxDeleteObjectsPerRequest+1) + for i := range paths { + paths[i] = fmt.Sprintf("file-%d", i) + } + s.srv.config.DeleteErrors = map[string]string{ + paths[0]: "first batch failed", + paths[maxDeleteObjectsPerRequest]: "second batch failed", + } + + err := s.storage.RemoveFiles(paths) + c.Check(err, ErrorMatches, "errors deleting multiple paths.*file-0: AccessDenied: first batch failed; file-1000: AccessDenied: second batch failed") + c.Check(s.countRequests("POST", "delete"), Equals, 2) +} + +func (s *PublishedStorageSuite) TestRemoveFilesDisableMultiDel(c *C) { + s.storage.disableMultiDel = true + + paths := []string{"a", "b", "c"} + for _, path := range paths { + s.PutFile(c, path, []byte("test")) + } + + err := s.storage.RemoveFiles([]string{"a", "b"}) + c.Check(err, IsNil) + + list, err := s.storage.Filelist("") + c.Check(err, IsNil) + c.Check(list, DeepEquals, []string{"c"}) + + // Endpoints that cannot multi-delete fall back to one request per file. + c.Check(s.countRequests("POST", "delete"), Equals, 0) + c.Check(s.countRequests("DELETE", ""), Equals, 2) +} + +func (s *PublishedStorageSuite) TestRemoveFilesDisableMultiDelError(c *C) { + s.storage.disableMultiDel = true + s.srv.config.DeleteObjectErrors = map[string]string{"a": "delete failed"} + + err := s.storage.RemoveFiles([]string{"a"}) + c.Check(err, ErrorMatches, "error deleting a.*AccessDenied.*delete failed.*") +} + +func (s *PublishedStorageSuite) TestRemoveFilesPlusWorkaround(c *C) { + s.storage.disableMultiDel = false + s.storage.plusWorkaround = true + + s.PutFile(c, "a/b+c", []byte("test")) + s.PutFile(c, "a/b", []byte("test")) + + // Filelist hides the space-substituted duplicate, so RemoveFiles has to + // expand it the way Remove does or it would be orphaned forever. + err := s.storage.RemoveFiles([]string{"a/b+c"}) + c.Check(err, IsNil) + + s.AssertNoFile(c, "a/b+c") + s.AssertNoFile(c, "a/b c") + + list, err := s.storage.Filelist("") + c.Check(err, IsNil) + c.Check(list, DeepEquals, []string{"a/b"}) +} + func (s *PublishedStorageSuite) TestRenameFile(c *C) { c.Skip("copy not available in s3test") } diff --git a/s3/server_test.go b/s3/server_test.go index 4a2e7111b..9e77d61a7 100644 --- a/s3/server_test.go +++ b/s3/server_test.go @@ -49,6 +49,12 @@ type Config struct { // all other regions. // http://docs.amazonwebservices.com/AmazonS3/latest/API/ErrorResponses.html Send409Conflict bool + // DeleteErrors maps object keys to errors returned by multi-object deletion. + DeleteErrors map[string]string + // DeleteObjectErrors maps object keys to errors returned by single-object deletion. + DeleteObjectErrors map[string]string + // DeleteObjectsError is returned for the entire multi-object deletion request. + DeleteObjectsError string } func (c *Config) send409Conflict() bool { @@ -516,9 +522,59 @@ func (r bucketResource) put(a *action) interface{} { return nil } -func (bucketResource) post(a *action) interface{} { - fatalError(400, "Method", "bucket POST method not available") - return nil +// POST on a bucket with ?delete deletes multiple objects at once. +// https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteObjects.html +func (r bucketResource) post(a *action) interface{} { + if _, multiDelete := a.req.URL.Query()["delete"]; !multiDelete { + fatalError(400, "Method", "bucket POST method not available") + return nil + } + + b := r.bucket + if b == nil { + fatalError(404, "NoSuchBucket", "The specified bucket does not exist") + } + if a.srv.config != nil && a.srv.config.DeleteObjectsError != "" { + fatalError(403, "AccessDenied", "%s", a.srv.config.DeleteObjectsError) + } + + var req struct { + Objects []struct { + Key string `xml:"Key"` + } `xml:"Object"` + } + if err := xml.NewDecoder(a.req.Body).Decode(&req); err != nil { + fatalError(400, "MalformedXML", "cannot parse delete request: %v", err) + } + if len(req.Objects) == 0 { + fatalError(400, "MalformedXML", "delete request contains no objects") + } + if len(req.Objects) > maxDeleteObjectsPerRequest { + fatalError(400, "MalformedXML", "delete request contains more than 1000 objects") + } + + type deleteError struct { + Key string `xml:"Key"` + Code string `xml:"Code"` + Message string `xml:"Message"` + } + var errors []deleteError + for _, obj := range req.Objects { + message, failed := "", false + if a.srv.config != nil { + message, failed = a.srv.config.DeleteErrors[obj.Key] + } + if failed { + errors = append(errors, deleteError{Key: obj.Key, Code: "AccessDenied", Message: message}) + } else { + delete(b.objects, obj.Key) + } + } + + return &struct { + XMLName xml.Name `xml:"DeleteResult"` + Errors []deleteError `xml:"Error"` + }{Errors: errors} } // validBucketName returns whether name is a valid bucket name. @@ -681,6 +737,11 @@ func (objr objectResource) put(a *action) interface{} { } func (objr objectResource) delete(a *action) interface{} { + if a.srv.config != nil { + if message, failed := a.srv.config.DeleteObjectErrors[objr.name]; failed { + fatalError(403, "AccessDenied", "%s", message) + } + } delete(objr.bucket.objects, objr.name) return nil }