mirror of
https://github.com/gotenberg/gotenberg.git
synced 2026-10-07 21:13:18 +01:00
Compare commits
6 Commits
cddaa0fa57
...
b16ce08da7
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b16ce08da7 | ||
|
|
fcfd590169 | ||
|
|
b39b8c76aa | ||
|
|
06ed58b6e7 | ||
|
|
34b7b4845e | ||
|
|
ab3832d9e5 |
6
.github/actions/build-test-push/action.yml
vendored
6
.github/actions/build-test-push/action.yml
vendored
@@ -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'
|
||||
|
||||
75
.github/actions/build-test-push/build.sh
vendored
75
.github/actions/build-test-push/build.sh
vendored
@@ -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[*]} \
|
||||
|
||||
1
Makefile
1
Makefile
@@ -87,6 +87,7 @@ LOG_STD_FORMAT=auto
|
||||
LOG_STD_ENABLE_GCP_FIELDS=false
|
||||
LOG_STD_LEVEL_CASE=lower
|
||||
PDFENGINES_DISABLE_ROUTES=false
|
||||
PDFENGINES_MAX_CONCURRENCY=1
|
||||
PDFENGINES_MERGE_ENGINES=qpdf,pdfcpu,pdftk
|
||||
PDFENGINES_SPLIT_ENGINES=pdfcpu,qpdf,pdftk
|
||||
PDFENGINES_FLATTEN_ENGINES=qpdf
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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}"
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user