mirror of
https://github.com/gotenberg/gotenberg.git
synced 2026-10-08 05:23:18 +01:00
perf(pdfengines): process a request's files concurrently behind an opt-in ceiling
This commit is contained in:
172
pkg/modules/pdfengines/concurrency.go
Normal file
172
pkg/modules/pdfengines/concurrency.go
Normal file
@@ -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
|
||||
}
|
||||
341
pkg/modules/pdfengines/concurrency_test.go
Normal file
341
pkg/modules/pdfengines/concurrency_test.go
Normal file
@@ -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])
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user