Compare commits

...

6 Commits

19 changed files with 933 additions and 323 deletions

View File

@@ -83,12 +83,16 @@ runs:
INPUT_PLATFORM: ${{ inputs.platform }}
INPUT_ALTERNATE_REPOSITORY: ${{ inputs.alternate_repository }}
INPUT_DRY_RUN: ${{ inputs.dry_run }}
# Exporting the build cache needs a registry login. Forks run without
# credentials, so they import the cache but never export it.
INPUT_CACHE_WRITABLE: ${{ inputs.docker_hub_username != '' }}
run: |
.github/actions/build-test-push/build.sh \
--version "$INPUT_VERSION" \
--platform "$INPUT_PLATFORM" \
--alternate-repository "$INPUT_ALTERNATE_REPOSITORY" \
--dry-run "$INPUT_DRY_RUN"
--dry-run "$INPUT_DRY_RUN" \
--cache-writable "$INPUT_CACHE_WRITABLE"
- name: Run integration tests
if: inputs.skip_integrations_tests != 'true'

View File

@@ -12,6 +12,7 @@ version=""
platform=""
alternate_repository=""
dry_run=""
cache_writable=""
while [[ $# -gt 0 ]]; do
case $1 in
@@ -31,6 +32,10 @@ while [[ $# -gt 0 ]]; do
dry_run="$2"
shift 2
;;
--cache-writable)
cache_writable="$2"
shift 2
;;
*)
echo "Unknown option $1"
exit 1
@@ -44,11 +49,41 @@ echo
echo "Gotenberg version: $version"
echo "Target platform: $platform"
# The build cache lives under the canonical repository, captured before the
# alternate-repository override below. Pull requests build into "snapshot", so
# deriving the cache ref after the override would give them a cache namespace
# of their own and they would never import what main published, which is the
# population that benefits most.
cache_image="$DOCKER_REGISTRY/$DOCKER_REPOSITORY"
# Layers are per-architecture, so each platform keeps its own cache manifest.
cache_platform="${platform//\//-}"
# Layers running "apt-get upgrade" install whatever versions are current at
# build time, and the packages are deliberately not pinned. A persistent cache
# would turn those into hits and freeze security patches into a published
# image until debian:13-slim itself changes digest. Keying them on the ISO week
# bounds that staleness to seven days while leaving every build within a week
# free to reuse the cache.
apt_snapshot="$(date -u +%G-W%V)"
# Only a build that is not redirected to an alternate repository writes the
# cache, so a pull request cannot make its own state the baseline for main.
# Reading stays enabled everywhere, including forks, since the cache ref is
# public and needs no credentials.
cache_to_enabled="false"
if [ "$cache_writable" = "true" ] && [ -z "$alternate_repository" ]; then
cache_to_enabled="true"
fi
if [ -n "$alternate_repository" ]; then
DOCKER_REPOSITORY=$alternate_repository
echo "⚠️ Using $alternate_repository for DOCKER_REPOSITORY"
fi
echo "Build cache: $cache_image:buildcache-<target>-$cache_platform (write: $cache_to_enabled)"
echo "APT snapshot: $apt_snapshot"
if [ "$dry_run" = "true" ]; then
echo "🚧 Dry run"
fi
@@ -189,12 +224,36 @@ join() {
echo "$*"
}
# cache_flags echoes the buildx cache arguments for a build target. Each target
# keeps its own manifest so that the Chromium and LibreOffice variants, which
# branch from common-stage rather than from each other, do not overwrite one
# another's entry.
#
# mode=max exports intermediate stages too, not just the final layers, which is
# what makes the expensive apt and jlink stages reusable. type=registry, not
# type=gha: the GitHub Actions cache is capped at 10 GB per repository and is
# already carrying the Go and golangci-lint caches that the lint and test jobs
# depend on. Multi-GB image layers across five platforms would evict them.
cache_flags() {
local target="$1"
local ref="$cache_image:buildcache-$target-$cache_platform"
local flags="--cache-from type=registry,ref=$ref"
if [ "$cache_to_enabled" = "true" ]; then
flags="$flags --cache-to type=registry,ref=$ref,mode=max"
fi
echo "$flags"
}
no_arch_tag="$DOCKER_REGISTRY/$DOCKER_REPOSITORY:$version"
# Full variant.
cmd="docker buildx build \
--target gotenberg \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg) \
--platform $platform \
--load \
${tags_flags[*]} \
@@ -207,6 +266,8 @@ run_cmd "$cmd"
cmd="docker buildx build \
--target gotenberg-chromium \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-chromium) \
--platform $platform \
--load \
${tags_chromium_flags[*]} \
@@ -218,6 +279,8 @@ run_cmd "$cmd"
cmd="docker buildx build \
--target gotenberg-libreoffice \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-libreoffice) \
--platform $platform \
--load \
${tags_libreoffice_flags[*]} \
@@ -230,6 +293,8 @@ if [ "$platform" = "linux/amd64" ]; then
cmd="docker buildx build \
--target gotenberg-cloudrun \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-cloudrun) \
--platform $platform \
--load \
${tags_cloud_run_flags[*]} \
@@ -240,6 +305,8 @@ if [ "$platform" = "linux/amd64" ]; then
cmd="docker buildx build \
--target gotenberg-cloudrun-chromium \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-cloudrun-chromium) \
--platform $platform \
--load \
${tags_cloud_run_chromium_flags[*]} \
@@ -250,6 +317,8 @@ if [ "$platform" = "linux/amd64" ]; then
cmd="docker buildx build \
--target gotenberg-cloudrun-libreoffice \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-cloudrun-libreoffice) \
--platform $platform \
--load \
${tags_cloud_run_libreoffice_flags[*]} \
@@ -263,6 +332,8 @@ if [ "$platform" = "linux/amd64" ] || [ "$platform" = "linux/arm64" ]; then
cmd="docker buildx build \
--target gotenberg-aws-lambda \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-aws-lambda) \
--platform $platform \
--load \
${tags_aws_lambda_flags[*]} \
@@ -273,6 +344,8 @@ if [ "$platform" = "linux/amd64" ] || [ "$platform" = "linux/arm64" ]; then
cmd="docker buildx build \
--target gotenberg-aws-lambda-chromium \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-aws-lambda-chromium) \
--platform $platform \
--load \
${tags_aws_lambda_chromium_flags[*]} \
@@ -283,6 +356,8 @@ if [ "$platform" = "linux/amd64" ] || [ "$platform" = "linux/arm64" ]; then
cmd="docker buildx build \
--target gotenberg-aws-lambda-libreoffice \
--build-arg GOTENBERG_VERSION=$version \
--build-arg APT_SNAPSHOT=$apt_snapshot \
$(cache_flags gotenberg-aws-lambda-libreoffice) \
--platform $platform \
--load \
${tags_aws_lambda_libreoffice_flags[*]} \

View File

@@ -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

View File

@@ -59,7 +59,14 @@ RUN go build -o gotenberg -ldflags "-s -w -X 'github.com/gotenberg/gotenberg/v8/
# ----------------------------------------------
FROM debian:13-slim AS custom-jre-stage
RUN apt-get update -qq \
# APT_SNAPSHOT busts every layer below it when CI rotates the value, weekly.
# Without it a persistent build cache turns the unpinned "apt-get upgrade" into
# a cache hit and the published image keeps shipping the package versions that
# were current when the cache was first populated.
ARG APT_SNAPSHOT=""
RUN echo "apt snapshot: $APT_SNAPSHOT" \
&& apt-get update -qq \
&& apt-get upgrade -yqq \
&& DEBIAN_FRONTEND=noninteractive apt-get install -y -qq --no-install-recommends default-jdk-headless binutils
@@ -114,9 +121,15 @@ FROM base-image-stage AS common-stage
ARG GOTENBERG_USER_GID=1001
ARG GOTENBERG_USER_UID=1001
# See the note on APT_SNAPSHOT in custom-jre-stage. Declaring it here covers
# every "apt-get upgrade" in the gotenberg, gotenberg-chromium and
# gotenberg-libreoffice targets too, since all three branch from this stage.
ARG APT_SNAPSHOT=""
# Create a non-root user.
# All processes in the Docker container will run with this dedicated user.
RUN groupadd --gid "$GOTENBERG_USER_GID" gotenberg \
RUN echo "apt snapshot: $APT_SNAPSHOT" \
&& groupadd --gid "$GOTENBERG_USER_GID" gotenberg \
&& useradd --uid "$GOTENBERG_USER_UID" --gid gotenberg --shell /bin/bash --home /home/gotenberg --no-create-home gotenberg \
&& mkdir /home/gotenberg \
&& chown gotenberg: /home/gotenberg

View File

@@ -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}"

View File

@@ -1,70 +0,0 @@
package gotenberg
import (
"context"
"errors"
"fmt"
"time"
"github.com/dlclark/regexp2"
)
// ErrFiltered happens if a value is filtered by the [FilterDeadline] function.
var ErrFiltered = errors.New("value filtered")
// FilterDeadline checks if the given value is allowed and not denied according
// to regex patterns. The allowed list uses OR semantics (value must match at
// least one pattern). The denied list uses OR semantics (value is denied if it
// matches any pattern). It returns a [context.DeadlineExceeded] if it takes
// too long to process.
func FilterDeadline(allowed, denied []*regexp2.Regexp, s string, deadline time.Time) error {
if len(allowed) > 0 {
matched := false
for _, pattern := range allowed {
// FIXME: not ideal to compile everytime, but is there another way to create a clone?
clone := regexp2.MustCompile(pattern.String(), 0)
clone.MatchTimeout = time.Until(deadline)
ok, err := clone.MatchString(s)
if err != nil {
if time.Now().After(deadline) {
return context.DeadlineExceeded
}
return fmt.Errorf("'%s' cannot handle '%s': %w", clone.String(), s, err)
}
if ok {
matched = true
break
}
}
if !matched {
return fmt.Errorf("'%s' does not match any expression from the allowed list: %w", s, ErrFiltered)
}
}
if len(denied) > 0 {
for _, pattern := range denied {
clone := regexp2.MustCompile(pattern.String(), 0)
clone.MatchTimeout = time.Until(deadline)
ok, err := clone.MatchString(s)
if err != nil {
if time.Now().After(deadline) {
return context.DeadlineExceeded
}
return fmt.Errorf("'%s' cannot handle '%s': %w", clone.String(), s, err)
}
if ok {
return fmt.Errorf("'%s' matches the expression from the denied list: %w", s, ErrFiltered)
}
}
}
return nil
}

View File

@@ -1,117 +0,0 @@
package gotenberg
import (
"context"
"errors"
"testing"
"time"
"github.com/dlclark/regexp2"
)
func TestFilterDeadline(t *testing.T) {
for _, tc := range []struct {
scenario string
allowed []*regexp2.Regexp
denied []*regexp2.Regexp
s string
deadline time.Time
expectError bool
expectedError error
}{
{
scenario: "DeadlineExceeded (allowed)",
allowed: []*regexp2.Regexp{regexp2.MustCompile("foo", 0)},
denied: nil,
s: "foo",
deadline: time.Now().Add(time.Duration(-1) * time.Hour),
expectError: true,
expectedError: context.DeadlineExceeded,
},
{
scenario: "ErrFiltered (allowed, no match)",
allowed: []*regexp2.Regexp{regexp2.MustCompile("foo", 0)},
denied: nil,
s: "bar",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: true,
expectedError: ErrFiltered,
},
{
scenario: "DeadlineExceeded (denied)",
allowed: nil,
denied: []*regexp2.Regexp{regexp2.MustCompile("foo", 0)},
s: "foo",
deadline: time.Now().Add(time.Duration(-1) * time.Hour),
expectError: true,
expectedError: context.DeadlineExceeded,
},
{
scenario: "ErrFiltered (denied)",
allowed: nil,
denied: []*regexp2.Regexp{regexp2.MustCompile("foo", 0)},
s: "foo",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: true,
expectedError: ErrFiltered,
},
{
scenario: "success (empty lists)",
allowed: nil,
denied: nil,
s: "foo",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: false,
},
{
scenario: "multi-pattern allow list, second matches",
allowed: []*regexp2.Regexp{regexp2.MustCompile("^https://", 0), regexp2.MustCompile("^file:///tmp/", 0)},
denied: nil,
s: "file:///tmp/abc/index.html",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: false,
},
{
scenario: "multi-pattern allow list, none matches",
allowed: []*regexp2.Regexp{regexp2.MustCompile("^https://", 0), regexp2.MustCompile("^ftp://", 0)},
denied: nil,
s: "file:///tmp/abc/index.html",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: true,
expectedError: ErrFiltered,
},
{
scenario: "multi-pattern deny list, second matches",
allowed: nil,
denied: []*regexp2.Regexp{regexp2.MustCompile("^ftp://", 0), regexp2.MustCompile("^file:.*", 0)},
s: "file:///etc/passwd",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: true,
expectedError: ErrFiltered,
},
{
scenario: "https URL passes deny list targeting file://",
allowed: nil,
denied: []*regexp2.Regexp{regexp2.MustCompile("^file:.*", 0)},
s: "https://example.com",
deadline: time.Now().Add(time.Duration(5) * time.Second),
expectError: false,
},
} {
t.Run(tc.scenario, func(t *testing.T) {
err := FilterDeadline(tc.allowed, tc.denied, tc.s, tc.deadline)
if tc.expectError && err == nil {
t.Fatal("expected an error but got none")
}
if !tc.expectError && err != nil {
t.Fatalf("expected no error but got: %v", err)
}
if tc.expectedError != nil && !errors.Is(err, tc.expectedError) {
t.Fatalf("expected error %v but got: %v", tc.expectedError, err)
}
})
}
}

View File

@@ -207,15 +207,38 @@ func (f *ParsedFlags) MustDeprecatedHumanReadableBytes(deprecated string, newNam
return f.MustHumanReadableBytes(newName)
}
// PatternMatchTimeout bounds a single match against an operator-supplied
// allow-list or deny-list pattern.
//
// regexp2 backtracks, and the strings matched against these patterns are
// client-controlled: a request URL, a CONNECT host. A pattern that backtracks
// catastrophically would otherwise burn a core for as long as the caller's
// deadline allows, which is --api-timeout (env API_TIMEOUT), 30 seconds by
// default. The ceiling mirrors the one the Chromium module already applies to
// the per-request extraHttpHeaders scope pattern.
//
// [ParsedFlags.MustRegexp] and [ParsedFlags.MustRegexpSlice] stamp this onto
// every pattern they compile, which is how all four production lists are
// built. Patterns compiled any other way keep regexp2's default of
// math.MaxInt64, which it treats as no timeout at all, so a hand-built slice
// must set this itself before reaching [DecideOutbound].
const PatternMatchTimeout = 250 * time.Millisecond
// MustRegexp returns the regular expression of a flag given by name.
// It panics if an error occurs.
//
// The returned expression carries [PatternMatchTimeout] and is safe to match
// on concurrently: callers must not compile a private copy per match.
func (f *ParsedFlags) MustRegexp(name string) *regexp2.Regexp {
val, err := f.GetString(name)
if err != nil {
panic(err)
}
return regexp2.MustCompile(val, 0)
re := regexp2.MustCompile(val, 0)
re.MatchTimeout = PatternMatchTimeout
return re
}
// MustDeprecatedRegexp returns the regular expression of a deprecated flag if
@@ -235,6 +258,9 @@ func (f *ParsedFlags) MustDeprecatedRegexp(deprecated string, newName string) *r
//
// Every allow-list and deny-list in Gotenberg is read through this method, so
// it is also where allow-list patterns are audited. See [AuditAllowList].
//
// The returned expressions carry [PatternMatchTimeout] and are safe to match
// on concurrently: callers must not compile a private copy per match.
func (f *ParsedFlags) MustRegexpSlice(name string) []*regexp2.Regexp {
vals := f.MustStringSlice(name)
@@ -246,7 +272,10 @@ func (f *ParsedFlags) MustRegexpSlice(name string) []*regexp2.Regexp {
continue
}
regexps = append(regexps, regexp2.MustCompile(val, 0))
re := regexp2.MustCompile(val, 0)
re.MatchTimeout = PatternMatchTimeout
regexps = append(regexps, re)
}
return regexps

View File

@@ -1054,3 +1054,36 @@ func TestEnvVarName(t *testing.T) {
})
}
}
func TestParsedFlags_RegexpMatchTimeout(t *testing.T) {
// [DecideOutbound] matches on these patterns directly instead of compiling
// a private copy per call, so the bound has to come from here. regexp2's
// own default is math.MaxInt64, which it treats as no
// timeout at all, so a pattern built without this stamp runs unbounded
// against a client-controlled string.
fs := flag.NewFlagSet("tests", flag.ContinueOnError)
fs.StringSlice("some-deny-list", []string{`^file:`, `^https?://`}, "")
fs.String("some-pattern", `^file:`, "")
err := fs.Parse(nil)
if err != nil {
t.Fatalf("expected no error but got: %v", err)
}
parsedFlags := ParsedFlags{FlagSet: fs}
regexps := parsedFlags.MustRegexpSlice("some-deny-list")
if len(regexps) != 2 {
t.Fatalf("expected 2 patterns but got %d", len(regexps))
}
for _, re := range regexps {
if re.MatchTimeout != PatternMatchTimeout {
t.Fatalf("pattern '%s' has MatchTimeout %s, expected %s", re.String(), re.MatchTimeout, PatternMatchTimeout)
}
}
if got := parsedFlags.MustRegexp("some-pattern").MatchTimeout; got != PatternMatchTimeout {
t.Fatalf("expected MustRegexp MatchTimeout %s but got %s", PatternMatchTimeout, got)
}
}

View File

@@ -27,6 +27,12 @@ import (
// example [::ffff:127.0.0.1]).
var ErrNonPublicIP = errors.New("non-public IP")
// ErrFiltered happens when a value is rejected by an allow-list or a
// deny-list, or when it cannot be validated and [DecideOutbound] fails closed.
// Callers map it to a generic 403: the specific reason stays in the operator
// logs so a client cannot probe the lists.
var ErrFiltered = errors.New("value filtered")
// ErrPublicIP indicates that an outbound URL targets an IP address that is
// reachable on the public internet. It is returned when a caller opts
// into denying public destinations via [WithDenyPublicIPs]; typical use
@@ -291,6 +297,15 @@ func DecideOutbound(ctx context.Context, rawURL string, allowList, denyList []*r
opt(&cfg)
}
// Each match is bounded by [PatternMatchTimeout] rather than by the
// remaining budget, so an already-spent deadline no longer surfaces from
// the match itself. Schemes that resolve a host still learn about it from
// resolveHost, but a non-matching file:// or data: URL returns before that
// point, so check it here to keep failing closed on every path.
if !time.Now().Before(deadline) {
return OutboundDecision{}, context.DeadlineExceeded
}
parsed, err := url.Parse(rawURL)
if err != nil {
return OutboundDecision{}, fmt.Errorf("parse URL %q: %w", rawURL, ErrFiltered)
@@ -314,15 +329,12 @@ func DecideOutbound(ctx context.Context, rawURL string, allowList, denyList []*r
allowMatched := false
if len(allowList) > 0 {
for _, pattern := range allowList {
clone := regexp2.MustCompile(pattern.String(), 0)
clone.MatchTimeout = time.Until(deadline)
ok, err := clone.MatchString(normalized)
ok, err := pattern.MatchString(normalized)
if err != nil {
if time.Now().After(deadline) {
return OutboundDecision{}, context.DeadlineExceeded
}
return OutboundDecision{}, fmt.Errorf("'%s' cannot handle '%s': %w", clone.String(), normalized, err)
return OutboundDecision{}, fmt.Errorf("'%s' cannot handle '%s': %w", pattern.String(), normalized, err)
}
if ok {
@@ -337,15 +349,12 @@ func DecideOutbound(ctx context.Context, rawURL string, allowList, denyList []*r
}
for _, pattern := range denyList {
clone := regexp2.MustCompile(pattern.String(), 0)
clone.MatchTimeout = time.Until(deadline)
ok, err := clone.MatchString(normalized)
ok, err := pattern.MatchString(normalized)
if err != nil {
if time.Now().After(deadline) {
return OutboundDecision{}, context.DeadlineExceeded
}
return OutboundDecision{}, fmt.Errorf("'%s' cannot handle '%s': %w", clone.String(), normalized, err)
return OutboundDecision{}, fmt.Errorf("'%s' cannot handle '%s': %w", pattern.String(), normalized, err)
}
if ok {
@@ -392,9 +401,8 @@ func DecideOutbound(ctx context.Context, rawURL string, allowList, denyList []*r
}
// FilterOutboundURL validates that rawURL is acceptable for an outbound
// request from Gotenberg. It is the URL-aware replacement for
// [FilterDeadline] and should be preferred for any new code that filters
// a URL before issuing or instructing an outbound request.
// request from Gotenberg. Prefer it for any new code that filters a URL
// before issuing or instructing an outbound request.
//
// The default behavior is permissive: the URL passes as long as it clears
// the regex allow-list and deny-list. Callers that need IP-class checks

View File

@@ -715,3 +715,54 @@ func TestNewOutboundHttpClient_NonPositiveTimeout(t *testing.T) {
t.Fatalf("timeout for a negative budget = %s, want a positive value so the client fails closed", got)
}
}
func TestDecideOutboundExpiredDeadline(t *testing.T) {
// Patterns are matched under the fixed PatternMatchTimeout rather than
// under the caller's remaining budget, so an expired deadline no longer
// surfaces from the match itself. Every scheme must still fail closed,
// including the ones that return before a host is resolved.
expired := time.Now().Add(-time.Second)
for _, rawURL := range []string{
"https://example.com/",
"file:///tmp/foo.html",
"data:text/html,hello",
} {
_, err := DecideOutbound(context.Background(), rawURL, nil, nil, expired)
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("DecideOutbound(%q) with an expired deadline = %v, want context.DeadlineExceeded", rawURL, err)
}
}
}
func TestDecideOutboundBoundsCatastrophicPatterns(t *testing.T) {
// A deny-list pattern that backtracks catastrophically, matched against a
// client-controlled URL. Before PatternMatchTimeout the ceiling was the
// caller's whole budget, so a 30s API_TIMEOUT bought a 30s CPU burn.
// The trailing "!" makes the match fail only after the nested quantifier
// has explored every way to split the run of "a"s.
pattern := regexp2.MustCompile(`^https://example\.com/(a+)+$`, 0)
pattern.MatchTimeout = PatternMatchTimeout
rawURL := "https://example.com/" + strings.Repeat("a", 40) + "!"
start := time.Now()
_, err := DecideOutbound(
context.Background(),
rawURL,
nil,
[]*regexp2.Regexp{pattern},
time.Now().Add(30*time.Second),
)
elapsed := time.Since(start)
if err == nil {
t.Fatal("expected an error from a catastrophic deny-list pattern")
}
// Generous headroom over the 250ms ceiling, still far below the 30s
// deadline the match would otherwise have been allowed to consume.
if elapsed > 5*time.Second {
t.Fatalf("match took %s, want it aborted near PatternMatchTimeout (%s)", elapsed, PatternMatchTimeout)
}
}

View File

@@ -309,33 +309,47 @@ func telemetryMiddleware(logger *slog.Logger, serverName, correlationIdHeader st
WriteBytes: c.Response().Size,
})...)
accessLogger := logger.
With(slog.String("log_type", "access")).
With(slog.String("correlation_id", correlationId)).
With(slog.String("remote_ip", c.RealIP())).
With(slog.String("host", c.Request().Host)).
With(slog.String("uri", c.Request().RequestURI)).
With(slog.String("method", c.Request().Method)).
With(slog.String("path", routePath)).
With(slog.String("referer", c.Request().Referer())).
With(slog.String("user_agent", c.Request().UserAgent())).
With(slog.Int("status", c.Response().Status)).
With(slog.Int64("latency", int64(finishTime.Sub(startTime)))).
With(slog.String("latency_human", finishTime.Sub(startTime).String())).
With(slog.Int64("bytes_in", c.Request().ContentLength)).
With(slog.Int64("bytes_out", c.Response().Size))
// Pick the level and message before building the record: err.Error
// walks a joined error chain, and the nil-error branch has no use
// for it.
level := slog.LevelInfo
msg := "request handled"
switch {
case err == nil:
accessLogger.InfoContext(ctx, "request handled")
case canceled:
// A client abort is expected, not a server failure; keep it
// visible but out of the error stream.
accessLogger.InfoContext(ctx, err.Error())
msg = err.Error()
default:
accessLogger.ErrorContext(ctx, err.Error())
level = slog.LevelError
msg = err.Error()
}
// One record rather than a chain of With calls. Each With clones
// the whole handler chain, and this logger fans out to a JSON
// handler and an OpenTelemetry bridge that is wired in even when no
// exporter is configured, so a 14-deep chain clones both sub-chains
// 14 times to emit a single line.
latency := finishTime.Sub(startTime)
logger.LogAttrs(ctx, level, msg,
slog.String("log_type", "access"),
slog.String("correlation_id", correlationId),
slog.String("remote_ip", c.RealIP()),
slog.String("host", c.Request().Host),
slog.String("uri", c.Request().RequestURI),
slog.String("method", c.Request().Method),
slog.String("path", routePath),
slog.String("referer", c.Request().Referer()),
slog.String("user_agent", c.Request().UserAgent()),
slog.Int("status", c.Response().Status),
slog.Int64("latency", int64(latency)),
slog.String("latency_human", latency.String()),
slog.Int64("bytes_in", c.Request().ContentLength),
slog.Int64("bytes_out", c.Response().Size),
)
additionalAttributes := []attribute.KeyValue{
semconvSrv.Route(routePath),
}

View File

@@ -323,7 +323,15 @@ func listenForEventResponseReceived(
return
}
logger.DebugContext(ctx, fmt.Sprintf("event EventResponseReceived fired for a resource: %+v", ev.Response))
// Formatting the whole response is the most expensive thing this
// listener does, and it runs per sub-resource on chromedp's single
// per-target event goroutine while that goroutine holds the mutex
// it also takes to dispatch command responses. At the default log
// level the result is discarded, so gate it on the level rather
// than let slog drop it after the fact.
if logger.Enabled(ctx, slog.LevelDebug) {
logger.DebugContext(ctx, fmt.Sprintf("event EventResponseReceived fired for a resource: %+v", ev.Response))
}
if slices.Contains(options.failOnResourceOnHttpStatusCode, ev.Response.Status) {
if !shouldCheckResourceHttpStatusCode(ev.Response.URL, normalizedIgnoreDomains) {

View 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
}

View 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])
}
}
}

View File

@@ -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),
}
}

View File

@@ -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)

View File

@@ -11,7 +11,6 @@ import (
"github.com/moby/moby/client"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/exec"
"github.com/testcontainers/testcontainers-go/network"
"github.com/testcontainers/testcontainers-go/wait"
)
@@ -21,11 +20,13 @@ import (
const testcontainersLabel = "org.testcontainers"
// PruneOrphanedNetworks removes dangling networks created by the test suite.
// Each scenario spins a dedicated network, and a failed container start can
// leak one before teardown records it. Leaked networks consume Docker's
// predefined address pools until none remain and every later scenario fails
// with "all predefined address pools have been fully subnetted". Call this
// before a run and between retries to reclaim the subnets.
// Scenarios no longer create one: the Gotenberg container is reached over its
// mapped port and the host-side helper over host.docker.internal, so the
// default bridge suffices. This stays as cheap insurance against networks
// leaked by an older suite version or an interrupted run, which consume
// Docker's predefined address pools until none remain and every later
// scenario fails with "all predefined address pools have been fully
// subnetted".
//
// Only unused networks bearing the testcontainers label are removed, so
// running containers and operator networks are never affected.
@@ -99,17 +100,18 @@ func applyDefaultEnv(env map[string]string) map[string]string {
return env
}
func startGotenbergContainer(ctx context.Context, env map[string]string) (*testcontainers.DockerNetwork, testcontainers.Container, error) {
// startGotenbergContainer starts a Gotenberg container on Docker's default
// bridge. No dedicated network is created: the suite addresses the container
// through container.Host plus its mapped port, and the container reaches the
// host-side webhook and static file server through the host.docker.internal
// alias below, so a per-scenario network would carry no traffic while still
// consuming one of Docker's predefined subnets.
func startGotenbergContainer(ctx context.Context, env map[string]string) (testcontainers.Container, error) {
ctx, cancel := context.WithTimeout(ctx, 2*time.Minute)
defer cancel()
env = applyDefaultEnv(env)
n, err := network.New(ctx)
if err != nil {
return nil, nil, fmt.Errorf("create Gotenberg container network: %w", err)
}
healthPath := "/health"
if env["API_ROOT_PATH"] != "" {
healthPath = fmt.Sprintf("%shealth", env["API_ROOT_PATH"])
@@ -122,7 +124,6 @@ func startGotenbergContainer(ctx context.Context, env map[string]string) (*testc
HostConfigModifier: func(hostConfig *container.HostConfig) {
hostConfig.ExtraHosts = []string{"host.docker.internal:host-gateway"}
},
Networks: []string{n.Name},
WaitingFor: wait.ForHTTP(healthPath),
Env: env,
}
@@ -148,19 +149,10 @@ func startGotenbergContainer(ctx context.Context, env map[string]string) (*testc
}
}
// The network is already created. The scenario teardown only
// removes networks it knows about, and the caller discards n on
// error, so remove it here to avoid leaking a subnet on every
// failed start. Leaked networks accumulate until Docker's address
// pools are fully subnetted and all later scenarios fail.
if errRemove := n.Remove(ctx); errRemove != nil {
err = fmt.Errorf("%w (also failed to remove network: %v)", err, errRemove)
}
return nil, nil, err
return nil, err
}
return n, c, nil
return c, nil
}
func execCommandInIntegrationToolsContainer(ctx context.Context, cmd []string, path string) (string, error) {

View File

@@ -86,16 +86,15 @@ func findScenarioLine(filePath, name string) int {
}
type scenario struct {
resp *httptest.ResponseRecorder
concurrentResps []*httptest.ResponseRecorder
probeResps []*httptest.ResponseRecorder
sequentialResps []*httptest.ResponseRecorder
workdir string
teststoreDir string
gotenbergContainer testcontainers.Container
gotenbergContainerNetwork *testcontainers.DockerNetwork
server *server
hostPort int
resp *httptest.ResponseRecorder
concurrentResps []*httptest.ResponseRecorder
probeResps []*httptest.ResponseRecorder
sequentialResps []*httptest.ResponseRecorder
workdir string
teststoreDir string
gotenbergContainer testcontainers.Container
server *server
hostPort int
}
func (s *scenario) reset(ctx context.Context) error {
@@ -123,11 +122,10 @@ func (s *scenario) reset(ctx context.Context) error {
}
func (s *scenario) iHaveADefaultGotenbergContainer(ctx context.Context) error {
n, c, err := startGotenbergContainer(ctx, nil)
c, err := startGotenbergContainer(ctx, nil)
if err != nil {
return fmt.Errorf("create Gotenberg container: %s", err)
}
s.gotenbergContainerNetwork = n
s.gotenbergContainer = c
return nil
}
@@ -137,11 +135,10 @@ func (s *scenario) iHaveAGotenbergContainerWithTheFollowingEnvironmentVariables(
for _, row := range envTable.Rows {
env[row.Cells[0].Value] = row.Cells[1].Value
}
n, c, err := startGotenbergContainer(ctx, env)
c, err := startGotenbergContainer(ctx, env)
if err != nil {
return fmt.Errorf("create Gotenberg container: %s", err)
}
s.gotenbergContainerNetwork = n
s.gotenbergContainer = c
return nil
}
@@ -1812,12 +1809,6 @@ func InitializeScenario(ctx *godog.ScenarioContext) {
return ctx, fmt.Errorf("terminate Gotenberg container: %w", errTerminate)
}
}
if s.gotenbergContainerNetwork != nil {
errRemove := s.gotenbergContainerNetwork.Remove(ctx)
if errRemove != nil {
return ctx, fmt.Errorf("remove Gotenberg container network: %w", errRemove)
}
}
return ctx, nil
})
ctx.After(func(ctx context.Context, sc *godog.Scenario, err error) (context.Context, error) {