From 67871b8e89f58c67f71e54ecdf296a366645c9ef Mon Sep 17 00:00:00 2001 From: Tomas Dvorak Date: Fri, 18 Sep 2026 22:43:27 +0200 Subject: [PATCH] perf(archive): parallel deflate workers for archive creation (#3316) Archive creation compressed every entry through a single zip.Writer, pinning the task to one CPU core. Entries are now enumerated first, then deflated on a worker pool (adaptive: min(GOMAXPROCS, 8), tunable via fs.WithArchiveWorkers): each worker spools its deflate stream to a temp file with precomputed CRC32/sizes, and a single writer drains results via zip.CreateRaw. Store/dry-run paths stay sequential. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- pkg/filemanager/fs/fs.go | 10 + pkg/filemanager/manager/archive.go | 230 +++++++++++++++--- .../manager/archive_parallel_test.go | 65 +++++ 3 files changed, 271 insertions(+), 34 deletions(-) create mode 100644 pkg/filemanager/manager/archive_parallel_test.go diff --git a/pkg/filemanager/fs/fs.go b/pkg/filemanager/fs/fs.go index c27a340e..0753ec72 100644 --- a/pkg/filemanager/fs/fs.go +++ b/pkg/filemanager/fs/fs.go @@ -327,6 +327,7 @@ type ( SysSkipSoftDelete bool Metadata map[string]string ArchiveCompression bool + ArchiveWorkers int ProgressFunc MaxArchiveSize int64 DryRun CreateArchiveDryRunFunc @@ -586,6 +587,15 @@ func WithArchiveCompression(b bool) Option { }) } +// WithArchiveWorkers sets the number of parallel compression workers for +// archive creation. <=0 selects an adaptive default, 1 keeps the +// sequential single-writer path. +func WithArchiveWorkers(n int) Option { + return OptionFunc(func(o *FsOption) { + o.ArchiveWorkers = n + }) +} + // WithMaxArchiveSize sets maximum size of to be archived file or to-be decompressed // size, 0 for unlimited. func WithMaxArchiveSize(s int64) Option { diff --git a/pkg/filemanager/manager/archive.go b/pkg/filemanager/manager/archive.go index efb4ff4b..d20403fb 100644 --- a/pkg/filemanager/manager/archive.go +++ b/pkg/filemanager/manager/archive.go @@ -2,13 +2,18 @@ package manager import ( "archive/zip" + "compress/flate" "context" "encoding/gob" "fmt" + "hash/crc32" "io" + "os" "path" "path/filepath" + "runtime" "strings" + "sync" "time" "github.com/bodgit/sevenzip" @@ -172,48 +177,48 @@ func (m *manager) CreateArchive(ctx context.Context, uris []*fs.URI, writer io.W zipWriter := zip.NewWriter(writer) defer zipWriter.Close() - var compressed int64 + // Enumerate all archive entries up front; compression itself then runs + // either sequentially or fanned out to workers (#3316). + jobs := make([]createArchiveJob, 0, len(files)) for _, file := range files { if file.Type() == types.FileTypeFile { - if err := m.compressFileToArchive(ctx, "/", file, zipWriter, o.ArchiveCompression, o.DryRun); err != nil { - failed++ - m.l.Warning("Failed to compress file %s: %s, skipping it...", file.Uri(false), err) - } - - compressed += file.Size() - if o.ProgressFunc != nil { - o.ProgressFunc(compressed, file.Size(), 0) - } + jobs = append(jobs, createArchiveJob{parent: "/", file: file}) + continue + } - if o.MaxArchiveSize > 0 && compressed > o.MaxArchiveSize { - return 0, fs.ErrArchiveSrcSizeTooBig + if err := m.Walk(ctx, file.Uri(false), intsets.MaxInt, func(f fs.File, level int) error { + if f.Type() == types.FileTypeFolder || f.IsSymbolic() { + return nil } + jobs = append(jobs, createArchiveJob{ + parent: strings.TrimPrefix(f.Uri(false).Dir(), file.Uri(false).Dir()), + file: f, + }) + return nil + }); err != nil { + m.l.Warning("Failed to walk folder %s: %s, skipping it...", file.Uri(false), err) + failed++ + } + } - } else { - if err := m.Walk(ctx, file.Uri(false), intsets.MaxInt, func(f fs.File, level int) error { - if f.Type() == types.FileTypeFolder || f.IsSymbolic() { - return nil - } - if err := m.compressFileToArchive(ctx, strings.TrimPrefix(f.Uri(false).Dir(), - file.Uri(false).Dir()), f, zipWriter, o.ArchiveCompression, o.DryRun); err != nil { - failed++ - m.l.Warning("Failed to compress file %s: %s, skipping it...", f.Uri(false), err) - } + if o.DryRun == nil && o.ArchiveCompression && o.ArchiveWorkers != 1 && len(jobs) > 1 { + return m.createArchiveParallel(ctx, jobs, zipWriter, o) + } - compressed += f.Size() - if o.ProgressFunc != nil { - o.ProgressFunc(compressed, f.Size(), 0) - } + var compressed int64 + for _, job := range jobs { + if err := m.compressFileToArchive(ctx, job.parent, job.file, zipWriter, o.ArchiveCompression, o.DryRun); err != nil { + failed++ + m.l.Warning("Failed to compress file %s: %s, skipping it...", job.file.Uri(false), err) + } - if o.MaxArchiveSize > 0 && compressed > o.MaxArchiveSize { - return fs.ErrArchiveSrcSizeTooBig - } + compressed += job.file.Size() + if o.ProgressFunc != nil { + o.ProgressFunc(compressed, job.file.Size(), 0) + } - return nil - }); err != nil { - m.l.Warning("Failed to walk folder %s: %s, skipping it...", file.Uri(false), err) - failed++ - } + if o.MaxArchiveSize > 0 && compressed > o.MaxArchiveSize { + return 0, fs.ErrArchiveSrcSizeTooBig } } @@ -257,6 +262,163 @@ func (m *manager) compressFileToArchive(ctx context.Context, parent string, file } +// maxArchiveWorkers caps the adaptive default so disk spooling stays sane. +const maxArchiveWorkers = 8 + +type ( + // createArchiveJob is a single file scheduled into an archive. + createArchiveJob struct { + parent string + file fs.File + } + // archiveEntryResult carries a pre-compressed entry ready for a raw + // zip write, along with the temp file holding its deflate stream. + archiveEntryResult struct { + header *zip.FileHeader + tmpPath string + file fs.File + err error + } +) + +// countingWriter tracks bytes written for zip header sizes. +type countingWriter struct { + w io.Writer + n int64 +} + +func (c *countingWriter) Write(p []byte) (int, error) { + n, err := c.w.Write(p) + c.n += int64(n) + return n, err +} + +// createArchiveParallel deflates archive entries on a worker pool while a +// single writer drains results in completion order via CreateRaw (#3316). +// Compressed data is spooled to temp files so zip's sequential layout is +// preserved without holding entries in memory. +func (m *manager) createArchiveParallel(ctx context.Context, jobs []createArchiveJob, zipWriter *zip.Writer, o *fs.FsOption) (int, error) { + workers := o.ArchiveWorkers + if workers <= 0 { + workers = min(runtime.GOMAXPROCS(0), maxArchiveWorkers) + } + + jobCh := make(chan createArchiveJob) + resCh := make(chan archiveEntryResult, workers) + var wg sync.WaitGroup + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for job := range jobCh { + resCh <- m.spoolEntry(ctx, job) + } + }() + } + go func() { + for _, job := range jobs { + jobCh <- job + } + close(jobCh) + wg.Wait() + close(resCh) + }() + + failed := 0 + var compressed int64 + for res := range resCh { + if res.err != nil { + failed++ + m.l.Warning("Failed to compress file %s: %s, skipping it...", res.file.Uri(false), res.err) + } else { + if err := writeRawEntry(zipWriter, res); err != nil { + failed++ + m.l.Warning("Failed to write archive entry %s: %s, skipping it...", res.file.Uri(false), err) + } + } + + compressed += res.file.Size() + if o.ProgressFunc != nil { + o.ProgressFunc(compressed, res.file.Size(), 0) + } + if o.MaxArchiveSize > 0 && compressed > o.MaxArchiveSize { + return 0, fs.ErrArchiveSrcSizeTooBig + } + } + + return failed, nil +} + +// spoolEntry deflates one file into a temp file and returns a header with +// precomputed CRC32 and sizes for zip.Writer.CreateRaw. +func (m *manager) spoolEntry(ctx context.Context, job createArchiveJob) archiveEntryResult { + res := archiveEntryResult{file: job.file} + es, err := m.GetEntitySource(ctx, job.file.PrimaryEntityID()) + if err != nil { + res.err = fmt.Errorf("failed to get entity source: %w", err) + return res + } + defer es.Close() + + tmp, err := os.CreateTemp("", "cloudreve-zip-*") + if err != nil { + res.err = err + return res + } + + cw := &countingWriter{w: tmp} + crc := crc32.NewIEEE() + fw, err := flate.NewWriter(cw, flate.DefaultCompression) + if err != nil { + tmp.Close() + os.Remove(tmp.Name()) + res.err = err + return res + } + es.Apply(entitysource.WithContext(ctx)) + _, copyErr := io.Copy(fw, io.TeeReader(es, crc)) + closeErr := fw.Close() + tmp.Close() + if copyErr != nil { + os.Remove(tmp.Name()) + res.err = copyErr + return res + } + if closeErr != nil { + os.Remove(tmp.Name()) + res.err = closeErr + return res + } + + res.header = &zip.FileHeader{ + Name: filepath.FromSlash(path.Join(job.parent, job.file.DisplayName())), + Method: zip.Deflate, + Modified: job.file.UpdatedAt(), + CRC32: crc.Sum32(), + CompressedSize64: uint64(cw.n), + UncompressedSize64: uint64(job.file.Size()), + } + res.tmpPath = tmp.Name() + return res +} + +// writeRawEntry writes a pre-compressed entry's local header and streams +// its spooled deflate data, then removes the temp file. +func writeRawEntry(zipWriter *zip.Writer, res archiveEntryResult) error { + defer os.Remove(res.tmpPath) + raw, err := zipWriter.CreateRaw(res.header) + if err != nil { + return err + } + tmp, err := os.Open(res.tmpPath) + if err != nil { + return err + } + defer tmp.Close() + _, err = io.Copy(raw, tmp) + return err +} + func getZipFileList(ctx context.Context, file io.ReaderAt, size int64, textEncoding encoding.Encoding) ([]ArchivedFile, error) { zr, err := zip.NewReader(file, size) if err != nil { diff --git a/pkg/filemanager/manager/archive_parallel_test.go b/pkg/filemanager/manager/archive_parallel_test.go new file mode 100644 index 00000000..e0d6bfbc --- /dev/null +++ b/pkg/filemanager/manager/archive_parallel_test.go @@ -0,0 +1,65 @@ +package manager + +import ( + "archive/zip" + "bytes" + "compress/flate" + "hash/crc32" + "io" + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" +) + +// TestWriteRawEntryRoundTrip verifies that a deflate stream spooled to a +// temp file with precomputed CRC32/sizes produces a valid zip entry via +// CreateRaw — the invariant the parallel path relies on (#3316). +func TestWriteRawEntryRoundTrip(t *testing.T) { + payload := []byte("hello cloudreve parallel zip payload, repeat repeat repeat") + + var compBuf bytes.Buffer + cw := &countingWriter{w: &compBuf} + crc := crc32.NewIEEE() + fw, err := flate.NewWriter(cw, flate.DefaultCompression) + require.NoError(t, err) + _, err = io.Copy(fw, io.TeeReader(bytes.NewReader(payload), crc)) + require.NoError(t, err) + require.NoError(t, fw.Close()) + + tmp := filepath.Join(t.TempDir(), "entry.deflate") + require.NoError(t, os.WriteFile(tmp, compBuf.Bytes(), 0o600)) + + var zipBuf bytes.Buffer + zw := zip.NewWriter(&zipBuf) + err = writeRawEntry(zw, archiveEntryResult{ + header: &zip.FileHeader{ + Name: "dir/a.txt", + Method: zip.Deflate, + CRC32: crc.Sum32(), + CompressedSize64: uint64(cw.n), + UncompressedSize64: uint64(len(payload)), + }, + tmpPath: tmp, + }) + require.NoError(t, err) + require.NoError(t, zw.Close()) + + // Temp file is consumed and removed. + _, statErr := os.Stat(tmp) + require.True(t, os.IsNotExist(statErr)) + + zr, err := zip.NewReader(bytes.NewReader(zipBuf.Bytes()), int64(zipBuf.Len())) + require.NoError(t, err) + require.Len(t, zr.File, 1) + require.Equal(t, "dir/a.txt", zr.File[0].Name) + require.Equal(t, uint64(len(payload)), zr.File[0].UncompressedSize64) + + rc, err := zr.File[0].Open() + require.NoError(t, err) + got, err := io.ReadAll(rc) + require.NoError(t, err) + require.NoError(t, rc.Close()) + require.Equal(t, payload, got) +}