diff --git a/Makefile b/Makefile index 85dace02..eb38231f 100644 --- a/Makefile +++ b/Makefile @@ -87,6 +87,7 @@ LOG_STD_FORMAT=auto LOG_STD_ENABLE_GCP_FIELDS=false LOG_STD_LEVEL_CASE=lower PDFENGINES_DISABLE_ROUTES=false +PDFENGINES_MAX_CONCURRENCY=1 PDFENGINES_MERGE_ENGINES=qpdf,pdfcpu,pdftk PDFENGINES_SPLIT_ENGINES=pdfcpu,qpdf,pdftk PDFENGINES_FLATTEN_ENGINES=qpdf diff --git a/compose.yaml b/compose.yaml index 1f5b61aa..19512e83 100644 --- a/compose.yaml +++ b/compose.yaml @@ -96,6 +96,7 @@ services: - "--pdfengines-embed-engines=${PDFENGINES_EMBED_ENGINES}" - "--pdfengines-embed-metadata-engines=${PDFENGINES_EMBED_METADATA_ENGINES}" - "--pdfengines-factur-x-engines=${PDFENGINES_FACTUR_X_ENGINES}" + - "--pdfengines-max-concurrency=${PDFENGINES_MAX_CONCURRENCY}" - "--pdfengines-disable-routes=${PDFENGINES_DISABLE_ROUTES}" - "--prometheus-namespace=${PROMETHEUS_NAMESPACE}" - "--prometheus-collect-interval=${PROMETHEUS_COLLECT_INTERVAL}" diff --git a/pkg/modules/pdfengines/concurrency.go b/pkg/modules/pdfengines/concurrency.go new file mode 100644 index 00000000..8dae7888 --- /dev/null +++ b/pkg/modules/pdfengines/concurrency.go @@ -0,0 +1,172 @@ +package pdfengines + +import ( + "sync" + + "github.com/gotenberg/gotenberg/v8/pkg/modules/api" +) + +// defaultMaxConcurrency is the number of PDF files a single stub processes at +// once when --pdfengines-max-concurrency (env PDFENGINES_MAX_CONCURRENCY) is +// not set. Each unit of work forks an external binary (qpdf, pdfcpu, pdftk or +// exiftool), so the ceiling trades wall clock against process count and RSS. +// +// It defaults to one, which processes files exactly as the sequential loops +// this package used to run did. Raising it only ever affects a request that +// carries several files, or one that splits into several outputs: a +// single-file request never reaches the concurrent path at all. Operators who +// send multi-file batches and have the memory headroom opt in. +// +// This never covers LibreOffice. libreoffice-pdfengine implements Convert and +// nothing else, every other [gotenberg.PdfEngine] method on it returns +// [gotenberg.ErrPdfEngineMethodNotSupported], and [ConvertStub] deliberately +// does not use this package's helpers. A soffice instance costs too much +// memory to run several of per container, so LibreOffice throughput is scaled +// by adding Gotenberg containers, not by raising this number. +const defaultMaxConcurrency = 1 + +// maxFileConcurrency is how many files one request may have in flight at once. +// It is replaced during [PdfEngines.Provision]. +var maxFileConcurrency = defaultMaxConcurrency + +// engineExtraSlots bounds the concurrency this package ADDS, across the whole +// process rather than per request, and holds one fewer slot than +// [maxFileConcurrency] because every request already owns one unit of its own. +// +// Bounding the added concurrency rather than the total is what keeps the +// ceiling from becoming a throughput regression. The sequential loops this +// helper replaced had no ceiling at all: X concurrent requests ran X engine +// binaries, one apiece. A pool covering the total would cut those X requests +// down to the ceiling, so an operator raising the flag to speed up a single +// multi-file request would slow the server down under real load. Reserving +// each request the unit it always had makes the worst case "what happened +// before, plus at most maxFileConcurrency-1". +// +// A per-request limit would have the opposite failure: X simultaneous requests +// forking X times the limit, trading the timeouts this exists to prevent for +// memory exhaustion. +var engineExtraSlots = make(chan struct{}, defaultMaxConcurrency-1) + +// acquireEngineSlot waits for the first unit of capacity to become available, +// either this request's reserved unit or a slot from the shared pool, and +// returns the function that gives it back. +// +// Waiting on both at once is the whole point. Committing to one source and +// blocking on it strands the other: a goroutine parked on an exhausted pool +// cannot pick up its own request's reserved unit when the file before it +// finishes, so the reserved units sit idle while every file queues on the +// pool, which is slower than having no pool at all. +func acquireEngineSlot(ctx *api.Context, reserved chan struct{}) (func(), error) { + // An [api.Context] carries a request context in production, but one built + // as a literal, which the unit tests do, embeds a nil [context.Context] + // and would panic on Done. A nil channel never fires, which correctly + // leaves the two capacity sources as the only things to wait on. + var done <-chan struct{} + if ctx != nil && ctx.Context != nil { + done = ctx.Done() + } + + // A select whose cancellation and capacity cases are both ready picks + // between them at random, so an already-dead request would start more + // files on a coin flip. Check first and stop taking on work. + if done != nil { + select { + case <-done: + return nil, ctx.Err() + default: + } + } + + select { + case <-reserved: + return func() { reserved <- struct{}{} }, nil + case engineExtraSlots <- struct{}{}: + return func() { <-engineExtraSlots }, nil + case <-done: + return nil, ctx.Err() + } +} + +// forEachInputPath runs fn against every input path, up to +// --pdfengines-max-concurrency (env PDFENGINES_MAX_CONCURRENCY) at a time. +// +// The stubs mutate each PDF in place, so distinct input paths never touch the +// same file and may run together. Callers that layer operations on one file, +// like [WatermarkStub] applying several watermarks in order, must keep that +// outer sequence and parallelize only the file dimension. +// +// Every path is attempted even after one fails, and the error returned is the +// first in input order rather than the first to arrive. That keeps the failing +// filename in the error message identical to what the sequential form +// reported, which the integration scenarios assert on. +func forEachInputPath(ctx *api.Context, inputPaths []string, fn func(inputPath string) error) error { + return forEachInputPathIndexed(ctx, inputPaths, func(_ int, inputPath string) error { + return fn(inputPath) + }) +} + +// forEachInputPathIndexed is [forEachInputPath] with the input path's index, +// for callers collecting a result per file. Writing into a preallocated slice +// at the given index needs no further synchronization; writing into a shared +// map does and must not be done from fn. +func forEachInputPathIndexed(ctx *api.Context, inputPaths []string, fn func(i int, inputPath string) error) error { + if len(inputPaths) == 0 { + return nil + } + + // The common case is a single file. Skip the goroutine and the slot: the + // caller is already inside whatever bound its own route applies. + if len(inputPaths) == 1 { + return fn(0, inputPaths[0]) + } + + // At the default ceiling of one, run the plain sequential loop this helper + // replaced. Racing goroutines for a single slot would serialize the work + // just the same, but the order files are picked up in would be down to the + // scheduler, and every file would be attempted even once one has failed. + // Taking the old path keeps the default a genuine no-op: same order, same + // early return, no goroutines. + if maxFileConcurrency < 2 { + for i, inputPath := range inputPaths { + err := fn(i, inputPath) + if err != nil { + return err + } + } + + return nil + } + + // The unit this request would have had all to itself before any of this + // existed. Whichever file claims it runs without touching the shared pool, + // so concurrent requests can never throttle each other below the + // one-binary-apiece they already got. See [engineExtraSlots]. + reserved := make(chan struct{}, 1) + reserved <- struct{}{} + + errs := make([]error, len(inputPaths)) + + var wg sync.WaitGroup + for i, inputPath := range inputPaths { + wg.Go(func() { + release, err := acquireEngineSlot(ctx, reserved) + if err != nil { + errs[i] = err + return + } + defer release() + + errs[i] = fn(i, inputPath) + }) + } + + wg.Wait() + + for _, err := range errs { + if err != nil { + return err + } + } + + return nil +} diff --git a/pkg/modules/pdfengines/concurrency_test.go b/pkg/modules/pdfengines/concurrency_test.go new file mode 100644 index 00000000..1a1380e3 --- /dev/null +++ b/pkg/modules/pdfengines/concurrency_test.go @@ -0,0 +1,341 @@ +package pdfengines + +import ( + "context" + "errors" + "fmt" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/gotenberg/gotenberg/v8/pkg/modules/api" +) + +func TestForEachInputPath(t *testing.T) { + for _, tc := range []struct { + scenario string + inputPaths []string + fn func(inputPath string) error + expectErr string + }{ + { + scenario: "no input path", + inputPaths: nil, + fn: func(string) error { return errors.New("must not run") }, + }, + { + scenario: "single input path", + inputPaths: []string{"a.pdf"}, + fn: func(string) error { return nil }, + }, + { + scenario: "single input path with error", + inputPaths: []string{"a.pdf"}, + fn: func(p string) error { return fmt.Errorf("boom %s", p) }, + expectErr: "boom a.pdf", + }, + { + scenario: "many input paths", + inputPaths: []string{"a.pdf", "b.pdf", "c.pdf", "d.pdf", "e.pdf"}, + fn: func(string) error { return nil }, + }, + { + scenario: "error is the first in input order, not the first to arrive", + inputPaths: []string{"a.pdf", "b.pdf", "c.pdf"}, + fn: func(p string) error { + // "c.pdf" fails without delay so that it lands well before + // "b.pdf"; the reported error must still be "b.pdf". + if p == "b.pdf" { + var counter int + for i := range 5_000_000 { + counter += i + } + return fmt.Errorf("slow failure %s (%d)", p, counter%1) + } + if p == "c.pdf" { + return fmt.Errorf("fast failure %s", p) + } + return nil + }, + expectErr: "slow failure b.pdf (0)", + }, + } { + t.Run(tc.scenario, func(t *testing.T) { + err := forEachInputPath(new(api.Context), tc.inputPaths, tc.fn) + + if tc.expectErr == "" { + if err != nil { + t.Fatalf("expected no error but got: %v", err) + } + return + } + + if err == nil { + t.Fatalf("expected error %q but got none", tc.expectErr) + } + + if err.Error() != tc.expectErr { + t.Fatalf("expected error %q but got %q", tc.expectErr, err.Error()) + } + }) + } +} + +func TestForEachInputPathRunsEveryPath(t *testing.T) { + inputPaths := make([]string, 50) + for i := range inputPaths { + inputPaths[i] = fmt.Sprintf("%d.pdf", i) + } + + var ( + mu sync.Mutex + seen = make(map[string]int) + ) + + err := forEachInputPath(new(api.Context), inputPaths, func(inputPath string) error { + mu.Lock() + defer mu.Unlock() + seen[inputPath]++ + return nil + }) + if err != nil { + t.Fatalf("expected no error but got: %v", err) + } + + if len(seen) != len(inputPaths) { + t.Fatalf("expected %d distinct paths but got %d", len(inputPaths), len(seen)) + } + + for path, count := range seen { + if count != 1 { + t.Fatalf("expected '%s' to run once but it ran %d times", path, count) + } + } +} + +func TestForEachInputPathRespectsTheSlotCeiling(t *testing.T) { + previousSlots, previousMax := engineExtraSlots, maxFileConcurrency + defer func() { engineExtraSlots, maxFileConcurrency = previousSlots, previousMax }() + + // One request may run its reserved unit plus ceiling-1 borrowed ones. + const ceiling = 3 + maxFileConcurrency = ceiling + engineExtraSlots = make(chan struct{}, ceiling-1) + + inputPaths := make([]string, 40) + for i := range inputPaths { + inputPaths[i] = fmt.Sprintf("%d.pdf", i) + } + + var inFlight, peak atomic.Int64 + + err := forEachInputPath(new(api.Context), inputPaths, func(string) error { + current := inFlight.Add(1) + defer inFlight.Add(-1) + + for { + observed := peak.Load() + if current <= observed || peak.CompareAndSwap(observed, current) { + break + } + } + + // Hold the slot long enough that the ceiling would be exceeded if it + // were not enforced. + var counter int + for i := range 200_000 { + counter += i + } + _ = counter + + return nil + }) + if err != nil { + t.Fatalf("expected no error but got: %v", err) + } + + if peak.Load() > ceiling { + t.Fatalf("expected at most %d concurrent runs but observed %d", ceiling, peak.Load()) + } +} + +func TestForEachInputPathHonorsCancellation(t *testing.T) { + previousSlots, previousMax := engineExtraSlots, maxFileConcurrency + defer func() { engineExtraSlots, maxFileConcurrency = previousSlots, previousMax }() + + // Concurrent path, with the shared pool exhausted by another request, so + // only this request's reserved unit is available. + maxFileConcurrency = 3 + engineExtraSlots = make(chan struct{}, 2) + engineExtraSlots <- struct{}{} + engineExtraSlots <- struct{}{} + + cancelledCtx, cancel := context.WithCancel(context.Background()) + cancel() + + // Neither source of capacity is available: the pool is exhausted by other + // requests and this request's reserved unit is already in use by one of its + // own files. A waiter must observe the cancelled context rather than block + // forever. Driving acquireEngineSlot directly keeps that deterministic: + // through forEachInputPath the reserved unit is reusable, so whether a + // given file waits at all depends on how fast the file before it finishes. + inUse := make(chan struct{}, 1) + + release, err := acquireEngineSlot(&api.Context{Context: cancelledCtx}, inUse) + if !errors.Is(err, context.Canceled) { + t.Fatalf("expected context.Canceled but got: %v", err) + } + + if release != nil { + t.Fatal("expected no release function when acquisition fails") + } + + // A cancelled request stops taking on work even when capacity is free, + // rather than deciding on the coin flip a ready select would give. + inUse <- struct{}{} + + _, err = acquireEngineSlot(&api.Context{Context: cancelledCtx}, inUse) + if !errors.Is(err, context.Canceled) { + t.Fatalf("expected context.Canceled with the reserved unit free but got: %v", err) + } + + // Live request, free reserved unit: acquired and handed back. + release, err = acquireEngineSlot(&api.Context{Context: context.Background()}, inUse) + if err != nil { + t.Fatalf("expected the reserved unit to be acquired but got: %v", err) + } + + release() + + if len(inUse) != 1 { + t.Fatalf("expected the reserved unit to be returned but the channel holds %d", len(inUse)) + } +} + +func TestForEachInputPathCompletesWithACancelledContext(t *testing.T) { + previousSlots, previousMax := engineExtraSlots, maxFileConcurrency + defer func() { engineExtraSlots, maxFileConcurrency = previousSlots, previousMax }() + + maxFileConcurrency = 3 + engineExtraSlots = make(chan struct{}, 2) + engineExtraSlots <- struct{}{} + engineExtraSlots <- struct{}{} + + cancelledCtx, cancel := context.WithCancel(context.Background()) + cancel() + + done := make(chan struct{}) + + go func() { + defer close(done) + _ = forEachInputPath(&api.Context{Context: cancelledCtx}, []string{"a.pdf", "b.pdf", "c.pdf"}, func(string) error { + return nil + }) + }() + + select { + case <-done: + case <-time.After(10 * time.Second): + t.Fatal("forEachInputPath hung on a cancelled context with the shared pool exhausted") + } +} + +func TestForEachInputPathNeverThrottlesBelowOnePerRequest(t *testing.T) { + previousSlots, previousMax := engineExtraSlots, maxFileConcurrency + defer func() { engineExtraSlots, maxFileConcurrency = previousSlots, previousMax }() + + // A small ceiling against far more concurrent requests than it covers. + // Before the shared pool existed each of these ran a binary of its own, so + // the pool must not drop aggregate concurrency below one per request. + const ( + ceiling = 2 + requests = 8 + ) + + maxFileConcurrency = ceiling + engineExtraSlots = make(chan struct{}, ceiling-1) + + // Every runner announces itself and then blocks, so the count of arrivals + // is the true simultaneous concurrency rather than whatever the scheduler + // happened to overlap. + arrived := make(chan struct{}, requests*3) + release := make(chan struct{}) + + var wg sync.WaitGroup + for range requests { + wg.Go(func() { + _ = forEachInputPath(new(api.Context), []string{"a.pdf", "b.pdf", "c.pdf"}, func(string) error { + arrived <- struct{}{} + <-release + return nil + }) + }) + } + + // One runner per request must be able to start without waiting on the + // shared pool. If the pool governed the total instead of the surplus, only + // `ceiling` runners would ever arrive and this would time out. + for i := range requests { + select { + case <-arrived: + case <-time.After(10 * time.Second): + t.Fatalf("only %d runners started concurrently, expected at least one per request (%d); the shared pool is throttling requests against each other", i, requests) + } + } + + close(release) + wg.Wait() +} + +func TestForEachInputPathIsSequentialAtTheDefaultCeiling(t *testing.T) { + if defaultMaxConcurrency != 1 { + t.Fatalf("this test pins the default as a no-op, but defaultMaxConcurrency is %d", defaultMaxConcurrency) + } + + previousSlots, previousMax := engineExtraSlots, maxFileConcurrency + defer func() { engineExtraSlots, maxFileConcurrency = previousSlots, previousMax }() + + maxFileConcurrency = defaultMaxConcurrency + engineExtraSlots = make(chan struct{}, defaultMaxConcurrency-1) + + // At the default the helper must behave exactly like the sequential loops + // it replaced: files in input order, and no file attempted once one has + // failed. + var order []string + + err := forEachInputPath(new(api.Context), []string{"a.pdf", "b.pdf", "c.pdf", "d.pdf"}, func(inputPath string) error { + order = append(order, inputPath) + if inputPath == "b.pdf" { + return errors.New("boom") + } + return nil + }) + + if err == nil || err.Error() != "boom" { + t.Fatalf("expected error \"boom\" but got: %v", err) + } + + if len(order) != 2 || order[0] != "a.pdf" || order[1] != "b.pdf" { + t.Fatalf("expected the run to stop after b.pdf in input order but got %v", order) + } +} + +func TestForEachInputPathIndexed(t *testing.T) { + inputPaths := []string{"a.pdf", "b.pdf", "c.pdf", "d.pdf"} + collected := make([]string, len(inputPaths)) + + err := forEachInputPathIndexed(new(api.Context), inputPaths, func(i int, inputPath string) error { + collected[i] = inputPath + return nil + }) + if err != nil { + t.Fatalf("expected no error but got: %v", err) + } + + for i, inputPath := range inputPaths { + if collected[i] != inputPath { + t.Fatalf("expected index %d to hold '%s' but got '%s'", i, inputPath, collected[i]) + } + } +} diff --git a/pkg/modules/pdfengines/pdfengines.go b/pkg/modules/pdfengines/pdfengines.go index af455ad5..109e1938 100644 --- a/pkg/modules/pdfengines/pdfengines.go +++ b/pkg/modules/pdfengines/pdfengines.go @@ -45,6 +45,7 @@ type PdfEngines struct { rotateNames []string facturXNames []string engines []gotenberg.PdfEngine + maxConcurrency int disableRoutes bool } @@ -70,6 +71,7 @@ func (mod *PdfEngines) Descriptor() gotenberg.ModuleDescriptor { fs.StringSlice("pdfengines-stamp-engines", []string{"pdfcpu", "pdftk"}, "Set the PDF engines and their order for the stamp feature - empty means all") fs.StringSlice("pdfengines-rotate-engines", []string{"pdfcpu", "pdftk"}, "Set the PDF engines and their order for the rotate feature - empty means all") fs.StringSlice("pdfengines-factur-x-engines", []string{"qpdf"}, "Set the PDF engines and their order for the Factur-X XMP feature - empty means all") + fs.Int("pdfengines-max-concurrency", defaultMaxConcurrency, "Set the maximum number of PDF files a feature processes concurrently, across all requests - bounds how many qpdf, pdfcpu, pdftk and exiftool processes run at once, so raising it trades memory for speed. Does not apply to LibreOffice: scale Gotenberg containers instead") fs.Bool("pdfengines-disable-routes", false, "Disable the routes") // Deprecated flags. @@ -105,8 +107,16 @@ func (mod *PdfEngines) Provision(ctx *gotenberg.Context) error { stampNames := flags.MustStringSlice("pdfengines-stamp-engines") rotateNames := flags.MustStringSlice("pdfengines-rotate-engines") facturXNames := flags.MustStringSlice("pdfengines-factur-x-engines") + mod.maxConcurrency = flags.MustInt("pdfengines-max-concurrency") mod.disableRoutes = flags.MustBool("pdfengines-disable-routes") + if mod.maxConcurrency > 0 { + maxFileConcurrency = mod.maxConcurrency + // One fewer than the ceiling: each request already reserves a unit of + // its own. See [engineExtraSlots]. + engineExtraSlots = make(chan struct{}, mod.maxConcurrency-1) + } + engines, err := ctx.Modules(new(gotenberg.PdfEngine)) if err != nil { return fmt.Errorf("get PDF engines: %w", err) @@ -222,6 +232,10 @@ func (mod *PdfEngines) Validate() error { return errors.New("no PDF engine is available; enable at least one engine module (e.g. qpdf, pdfcpu, pdftk, libreoffice-pdfengine, exiftool)") } + if mod.maxConcurrency < 1 { + return fmt.Errorf("PDF engines max concurrency must be at least 1, got %d; set --pdfengines-max-concurrency (env PDFENGINES_MAX_CONCURRENCY) to a positive value", mod.maxConcurrency) + } + availableEngines := make([]string, len(mod.engines)) for i, engine := range mod.engines { @@ -296,6 +310,7 @@ func (mod *PdfEngines) SystemMessages() []string { fmt.Sprintf("stamp engines - %s", strings.Join(mod.stampNames, " ")), fmt.Sprintf("rotate engines - %s", strings.Join(mod.rotateNames, " ")), fmt.Sprintf("factur-x engines - %s", strings.Join(mod.facturXNames, " ")), + fmt.Sprintf("max concurrency - %d", mod.maxConcurrency), } } diff --git a/pkg/modules/pdfengines/routes.go b/pkg/modules/pdfengines/routes.go index f6460bf7..d019b14d 100644 --- a/pkg/modules/pdfengines/routes.go +++ b/pkg/modules/pdfengines/routes.go @@ -208,14 +208,14 @@ func RotateStub(ctx *api.Context, engine gotenberg.PdfEngine, angle int, pages s return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.Rotate(ctx, ctx.Log(), inputPath, angle, pages) if err != nil { return fmt.Errorf("rotate '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // ValidatePdfFormatsCompat checks for incompatible combinations of PDF formats @@ -334,14 +334,14 @@ func SplitPdfStub(ctx *api.Context, engine gotenberg.PdfEngine, mode gotenberg.S // FlattenStub merges annotation appearances with page content for each given // PDF, effectively deleting the original annotations. func FlattenStub(ctx *api.Context, engine gotenberg.PdfEngine, inputPaths []string) error { - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.Flatten(ctx, ctx.Log(), inputPath) if err != nil { return fmt.Errorf("flatten '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // defaultImageQuality is the JPEG quality applied by the image optimization @@ -393,19 +393,27 @@ func OptimizeStub(ctx *api.Context, engine gotenberg.PdfEngine, optimizeImages b return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.OptimizeImages(ctx, ctx.Log(), imageQuality, inputPath) if err != nil { return fmt.Errorf("optimize images of '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // ConvertStub transforms a given PDF to the specified formats defined in // [gotenberg.PdfFormats]. If no format, it does nothing and returns the input // paths. +// +// This loop stays sequential on purpose. Convert is the one PDF engine method +// LibreOffice implements, and libreoffice-pdfengine is the default and only +// convert engine, so every iteration here drives the single soffice daemon. A +// LibreOffice instance is far too memory-hungry to run several of per +// container: the way to convert more documents at once is to scale Gotenberg +// containers, not to widen this loop. Do not route it through +// [forEachInputPath]. func ConvertStub(ctx *api.Context, engine gotenberg.PdfEngine, formats gotenberg.PdfFormats, inputPaths []string) ([]string, error) { zeroValued := gotenberg.PdfFormats{} if formats == zeroValued { @@ -432,14 +440,14 @@ func WriteMetadataStub(ctx *api.Context, engine gotenberg.PdfEngine, metadata ma return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.WriteMetadata(ctx, ctx.Log(), metadata, inputPath) if err != nil { return fmt.Errorf("write metadata into '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // documentTitle returns the input PDF's Title metadata entry, falling back to @@ -491,28 +499,33 @@ func WriteBookmarksStub(ctx *api.Context, engine gotenberg.PdfEngine, bookmarks return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.WriteBookmarks(ctx, ctx.Log(), inputPath, b) if err != nil { return fmt.Errorf("write bookmarks into '%s': %w", inputPath, err) } - } + + return nil + }) case map[string][]gotenberg.Bookmark: - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { filename := ctx.OriginalFilename(inputPath) - if specificBookmarks, ok := b[filename]; ok { - err := engine.WriteBookmarks(ctx, ctx.Log(), inputPath, specificBookmarks) - if err != nil { - return fmt.Errorf("write bookmarks into '%s': %w", inputPath, err) - } + specificBookmarks, ok := b[filename] + if !ok { + return nil } - } + + err := engine.WriteBookmarks(ctx, ctx.Log(), inputPath, specificBookmarks) + if err != nil { + return fmt.Errorf("write bookmarks into '%s': %w", inputPath, err) + } + + return nil + }) default: // Should not happen. return fmt.Errorf("bookmarks type '%T' not supported", bookmarks) } - - return nil } // FormDataPdfEmbeds extracts embedded file paths from form data. @@ -537,14 +550,14 @@ func EmbedFilesMetadataStub(ctx *api.Context, engine gotenberg.PdfEngine, metada return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.EmbedFilesMetadata(ctx, ctx.Log(), metadata, inputPath) if err != nil { return fmt.Errorf("set embeds metadata on PDF '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // FormDataPdfFacturX extracts the Factur-X parameters and the invoice XML path @@ -737,14 +750,14 @@ func InjectFacturXXMPStub(ctx *api.Context, engine gotenberg.PdfEngine, facturX return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.InjectFacturXXMP(ctx, ctx.Log(), facturX, inputPath) if err != nil { return fmt.Errorf("inject Factur-X XMP into PDF '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // FormDataPdfEncrypt extracts the encryption parameters and permissions from @@ -783,14 +796,14 @@ func EncryptPdfStub(ctx *api.Context, engine gotenberg.PdfEngine, opts gotenberg return nil } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.Encrypt(ctx, ctx.Log(), inputPath, opts) if err != nil { return fmt.Errorf("encrypt PDF '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // EmbedFilesStub embeds files into PDF files. @@ -819,14 +832,14 @@ func EmbedFilesStub(ctx *api.Context, engine gotenberg.PdfEngine, embedPaths []s resolvedPaths[i] = resolvedPath } - for _, inputPath := range inputPaths { + return forEachInputPath(ctx, inputPaths, func(inputPath string) error { err := engine.EmbedFiles(ctx, ctx.Log(), resolvedPaths, inputPath) if err != nil { return fmt.Errorf("embed files into PDF '%s': %w", inputPath, err) } - } - return nil + return nil + }) } // FormDataPdfStamps builds the ordered list of stamps from the repeated stamp @@ -943,16 +956,23 @@ func bindStampOrWatermarkFiles(stamps []gotenberg.Stamp, files []string, kind st // WatermarkStub applies each watermark to a list of PDF files, in order. // Entries with no source are skipped, so an empty list does nothing. func WatermarkStub(ctx *api.Context, engine gotenberg.PdfEngine, watermarks []gotenberg.Stamp, inputPaths []string) error { + // Watermarks stack on the same file, so the outer loop stays sequential; + // only the file dimension is parallel. for _, watermark := range watermarks { if watermark.Source == "" { continue } - for _, inputPath := range inputPaths { - err := engine.Watermark(ctx, ctx.Log(), inputPath, watermark) - if err != nil { - return fmt.Errorf("watermark '%s': %w", inputPath, err) + err := forEachInputPath(ctx, inputPaths, func(inputPath string) error { + errWatermark := engine.Watermark(ctx, ctx.Log(), inputPath, watermark) + if errWatermark != nil { + return fmt.Errorf("watermark '%s': %w", inputPath, errWatermark) } + + return nil + }) + if err != nil { + return err } } @@ -962,16 +982,23 @@ func WatermarkStub(ctx *api.Context, engine gotenberg.PdfEngine, watermarks []go // StampStub applies each stamp to a list of PDF files, in order. Entries with // no source are skipped, so an empty list does nothing. func StampStub(ctx *api.Context, engine gotenberg.PdfEngine, stamps []gotenberg.Stamp, inputPaths []string) error { + // Stamps stack on the same file, so the outer loop stays sequential; only + // the file dimension is parallel. for _, stamp := range stamps { if stamp.Source == "" { continue } - for _, inputPath := range inputPaths { - err := engine.Stamp(ctx, ctx.Log(), inputPath, stamp) - if err != nil { - return fmt.Errorf("stamp '%s': %w", inputPath, err) + err := forEachInputPath(ctx, inputPaths, func(inputPath string) error { + errStamp := engine.Stamp(ctx, ctx.Log(), inputPath, stamp) + if errStamp != nil { + return fmt.Errorf("stamp '%s': %w", inputPath, errStamp) } + + return nil + }) + if err != nil { + return err } } @@ -1476,14 +1503,25 @@ func readMetadataRoute(engine gotenberg.PdfEngine) api.Route { return fmt.Errorf("validate form data: %w", err) } - res := make(map[string]map[string]any, len(inputPaths)) - for _, inputPath := range inputPaths { - metadata, err := engine.ReadMetadata(ctx, ctx.Log(), inputPath) - if err != nil { - return fmt.Errorf("read metadata: %w", err) + // Collected per index, then folded into the map on this + // goroutine: a shared map cannot be written concurrently. + collected := make([]map[string]any, len(inputPaths)) + err = forEachInputPathIndexed(ctx, inputPaths, func(i int, inputPath string) error { + metadata, errRead := engine.ReadMetadata(ctx, ctx.Log(), inputPath) + if errRead != nil { + return fmt.Errorf("read metadata: %w", errRead) } - res[ctx.OriginalFilename(inputPath)] = metadata + collected[i] = metadata + return nil + }) + if err != nil { + return err + } + + res := make(map[string]map[string]any, len(inputPaths)) + for i, inputPath := range inputPaths { + res[ctx.OriginalFilename(inputPath)] = collected[i] } err = c.JSON(http.StatusOK, res) @@ -1554,14 +1592,25 @@ func readBookmarksRoute(engine gotenberg.PdfEngine) api.Route { return fmt.Errorf("validate form data: %w", err) } - res := make(map[string][]gotenberg.Bookmark, len(inputPaths)) - for _, inputPath := range inputPaths { - bookmarks, err := engine.ReadBookmarks(ctx, ctx.Log(), inputPath) - if err != nil { - return fmt.Errorf("read bookmarks: %w", err) + // Collected per index, then folded into the map on this + // goroutine: a shared map cannot be written concurrently. + collected := make([][]gotenberg.Bookmark, len(inputPaths)) + err = forEachInputPathIndexed(ctx, inputPaths, func(i int, inputPath string) error { + bookmarks, errRead := engine.ReadBookmarks(ctx, ctx.Log(), inputPath) + if errRead != nil { + return fmt.Errorf("read bookmarks: %w", errRead) } - res[ctx.OriginalFilename(inputPath)] = bookmarks + collected[i] = bookmarks + return nil + }) + if err != nil { + return err + } + + res := make(map[string][]gotenberg.Bookmark, len(inputPaths)) + for i, inputPath := range inputPaths { + res[ctx.OriginalFilename(inputPath)] = collected[i] } err = c.JSON(http.StatusOK, res)