Bound concurrent image processing and upstream fetches (closes #64)
check / check (push) Successful in 17s
check / check (push) Successful in 17s
Nothing bounded total in-flight work, so a burst of cache misses across hosts could exhaust memory. Two settings now do: max_concurrent_processing (default the CPUs Go uses) and upstream_connections (default 64, beside the per-host limit). A request that finds either full waits up to 10 seconds, then gets 503; a slot is released on every path. No request holds source bytes while it waits: a cached source is read only after the processing slot is taken, and a fetched one only while it holds its upstream connection. libvips runs one worker thread per image with its operation cache off. Both waits count toward downstream_timeout, as the README says. Model: opus-5-5
This commit was merged in pull request #148.
This commit is contained in:
@@ -7,7 +7,9 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"runtime"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/davidbyttow/govips/v2/vips"
|
||||
)
|
||||
@@ -17,11 +19,21 @@ import (
|
||||
//nolint:gochecknoglobals // package-level sync.Once for one-time vips init
|
||||
var vipsOnce sync.Once
|
||||
|
||||
// initVips initializes libvips with quiet logging.
|
||||
// initVips initializes libvips with quiet logging, one worker thread per
|
||||
// image and no operation cache. Process already works on one image per CPU
|
||||
// by default, so more threads per image would only compete for the CPUs.
|
||||
// Each request decodes different source bytes, so the operation cache
|
||||
// would rarely be hit and would hold memory outside MaxConcurrentProcessing;
|
||||
// repeated requests are served from pixa's disk cache instead.
|
||||
func initVips() {
|
||||
vipsOnce.Do(func() {
|
||||
vips.LoggingSettings(nil, vips.LogLevelError)
|
||||
vips.Startup(nil)
|
||||
vips.Startup(&vips.Config{
|
||||
ConcurrencyLevel: 1,
|
||||
MaxCacheSize: 0,
|
||||
MaxCacheMem: 0,
|
||||
MaxCacheFiles: 0,
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
@@ -106,9 +118,23 @@ var ErrInputDataTooLarge = errors.New("input data exceeds maximum allowed size")
|
||||
// not supported.
|
||||
var ErrUnsupportedOutputFormat = errors.New("unsupported output format")
|
||||
|
||||
// ErrTooManyImages is returned when MaxConcurrentProcessing images are being
|
||||
// processed and none finishes within ProcessingWaitTimeout.
|
||||
var ErrTooManyImages = errors.New("too many images being processed at once")
|
||||
|
||||
// ProcessingWaitTimeout is how long Process waits for a free slot when
|
||||
// MaxConcurrentProcessing images are already being processed.
|
||||
const ProcessingWaitTimeout = 10 * time.Second
|
||||
|
||||
// ImageProcessor implements image transformation using libvips via govips.
|
||||
type ImageProcessor struct {
|
||||
maxInputBytes int64
|
||||
// processingSemaphore has one slot per image that may be processed at
|
||||
// once. Process holds a slot from before it reads its input until it
|
||||
// returns, so the input, the decoded image and the output all count.
|
||||
processingSemaphore chan struct{}
|
||||
// processingWaitTimeout is ProcessingWaitTimeout; tests shorten it.
|
||||
processingWaitTimeout time.Duration
|
||||
}
|
||||
|
||||
// Params holds configuration for creating an ImageProcessor.
|
||||
@@ -117,6 +143,9 @@ type Params struct {
|
||||
// MaxInputBytes is the maximum allowed input size in bytes.
|
||||
// If <= 0, DefaultMaxInputBytes is used.
|
||||
MaxInputBytes int64
|
||||
// MaxConcurrentProcessing is the most images processed at once.
|
||||
// If <= 0, the number of CPUs Go uses (runtime.GOMAXPROCS(0)) is used.
|
||||
MaxConcurrentProcessing int
|
||||
}
|
||||
|
||||
// New creates a new image processor with the given parameters.
|
||||
@@ -129,17 +158,34 @@ func New(params Params) *ImageProcessor {
|
||||
maxInputBytes = DefaultMaxInputBytes
|
||||
}
|
||||
|
||||
maxConcurrentProcessing := params.MaxConcurrentProcessing
|
||||
if maxConcurrentProcessing <= 0 {
|
||||
maxConcurrentProcessing = runtime.GOMAXPROCS(0)
|
||||
}
|
||||
|
||||
return &ImageProcessor{
|
||||
maxInputBytes: maxInputBytes,
|
||||
maxInputBytes: maxInputBytes,
|
||||
processingSemaphore: make(chan struct{}, maxConcurrentProcessing),
|
||||
processingWaitTimeout: ProcessingWaitTimeout,
|
||||
}
|
||||
}
|
||||
|
||||
// Process transforms an image according to the request.
|
||||
// Process transforms an image according to the request. When
|
||||
// MaxConcurrentProcessing images are already being processed, it waits up
|
||||
// to ProcessingWaitTimeout for one to finish, then fails with
|
||||
// ErrTooManyImages.
|
||||
func (p *ImageProcessor) Process(
|
||||
_ context.Context,
|
||||
ctx context.Context,
|
||||
input io.Reader,
|
||||
req *Request,
|
||||
) (*Result, error) {
|
||||
release, err := p.acquireSlot(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
defer release()
|
||||
|
||||
// Read input with a size limit to prevent unbounded memory consumption.
|
||||
// We read at most maxInputBytes+1 so we can detect if the input exceeds
|
||||
// the limit without consuming additional memory.
|
||||
@@ -285,6 +331,29 @@ func FormatToMIME(format Format) string {
|
||||
}
|
||||
}
|
||||
|
||||
// acquireSlot takes a slot in processingSemaphore, waiting at most
|
||||
// processingWaitTimeout for one to free up, and returns the func that gives
|
||||
// it back. A free slot is taken even when ctx has ended; only the wait for
|
||||
// one stops when ctx ends, as the rest of Process does not check ctx.
|
||||
func (p *ImageProcessor) acquireSlot(ctx context.Context) (func(), error) {
|
||||
release := func() { <-p.processingSemaphore }
|
||||
|
||||
select {
|
||||
case p.processingSemaphore <- struct{}{}:
|
||||
return release, nil
|
||||
default:
|
||||
}
|
||||
|
||||
select {
|
||||
case p.processingSemaphore <- struct{}{}:
|
||||
return release, nil
|
||||
case <-time.After(p.processingWaitTimeout):
|
||||
return nil, ErrTooManyImages
|
||||
case <-ctx.Done():
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
// detectFormat returns the format string from a vips image.
|
||||
func (p *ImageProcessor) detectFormat(img *vips.ImageRef) string {
|
||||
format := img.Format()
|
||||
|
||||
@@ -0,0 +1,302 @@
|
||||
package imageprocessor
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"testing/iotest"
|
||||
"time"
|
||||
)
|
||||
|
||||
// errTestReadFailed is the error the unreadable test input returns.
|
||||
var errTestReadFailed = errors.New("test input cannot be read")
|
||||
|
||||
// readingCounter counts the Process calls reading their input at the same
|
||||
// time and remembers the most there ever were.
|
||||
type readingCounter struct {
|
||||
mu sync.Mutex
|
||||
reading int
|
||||
most int
|
||||
}
|
||||
|
||||
func (c *readingCounter) start() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
c.reading++
|
||||
c.most = max(c.most, c.reading)
|
||||
}
|
||||
|
||||
func (c *readingCounter) stop() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
c.reading--
|
||||
}
|
||||
|
||||
func (c *readingCounter) mostReading() int {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
return c.most
|
||||
}
|
||||
|
||||
// gatedReader is a Process input. Its first Read counts the call in,
|
||||
// reports it on entered and blocks until gate is closed; it counts the call
|
||||
// out when it returns io.EOF. Process reads its input only while it holds a
|
||||
// processing slot, so the count never goes above MaxConcurrentProcessing.
|
||||
type gatedReader struct {
|
||||
data *bytes.Reader
|
||||
gate <-chan struct{}
|
||||
entered chan<- struct{}
|
||||
counter *readingCounter
|
||||
started bool
|
||||
}
|
||||
|
||||
func (r *gatedReader) Read(p []byte) (int, error) {
|
||||
if !r.started {
|
||||
r.started = true
|
||||
r.counter.start()
|
||||
|
||||
r.entered <- struct{}{}
|
||||
|
||||
<-r.gate
|
||||
}
|
||||
|
||||
n, err := r.data.Read(p)
|
||||
if errors.Is(err, io.EOF) {
|
||||
r.counter.stop()
|
||||
}
|
||||
|
||||
return n, err
|
||||
}
|
||||
|
||||
// smallJPEGRequest asks for a 5x5 JPEG.
|
||||
func smallJPEGRequest() *Request {
|
||||
return &Request{
|
||||
Size: Size{Width: 5, Height: 5},
|
||||
Format: FormatJPEG,
|
||||
Quality: 85,
|
||||
FitMode: FitCover,
|
||||
}
|
||||
}
|
||||
|
||||
// processInBackground runs Process on reader in a new goroutine and sends
|
||||
// its error on results.
|
||||
func processInBackground(
|
||||
proc *ImageProcessor, reader *gatedReader, results chan<- error,
|
||||
) {
|
||||
go func() {
|
||||
result, err := proc.Process(context.Background(), reader, smallJPEGRequest())
|
||||
if err == nil {
|
||||
_ = result.Content.Close()
|
||||
}
|
||||
|
||||
results <- err
|
||||
}()
|
||||
}
|
||||
|
||||
// waitForEntries fails the test unless count Process calls report on
|
||||
// entered within a few seconds.
|
||||
func waitForEntries(t *testing.T, entered <-chan struct{}, count int) {
|
||||
t.Helper()
|
||||
|
||||
for range count {
|
||||
select {
|
||||
case <-entered:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("Process calls did not start reading their input")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewDefaultsMaxConcurrentProcessingToCPUs(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
for _, limit := range []int{0, -1} {
|
||||
proc := New(Params{MaxConcurrentProcessing: limit})
|
||||
if got := cap(proc.processingSemaphore); got != runtime.GOMAXPROCS(0) {
|
||||
t.Errorf("MaxConcurrentProcessing %d: %d slots, want %d, one per CPU",
|
||||
limit, got, runtime.GOMAXPROCS(0))
|
||||
}
|
||||
}
|
||||
|
||||
proc := New(Params{MaxConcurrentProcessing: 3})
|
||||
if got := cap(proc.processingSemaphore); got != 3 {
|
||||
t.Errorf("MaxConcurrentProcessing 3: %d slots, want 3", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestProcessNeverExceedsMaxConcurrentProcessing starts more Process calls
|
||||
// than MaxConcurrentProcessing allows and holds the first ones inside
|
||||
// Process until the test lets them go. No more than the limit may be
|
||||
// working at once, and the calls held back must wait for a slot and then
|
||||
// succeed.
|
||||
func TestProcessNeverExceedsMaxConcurrentProcessing(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const (
|
||||
limit = 2
|
||||
calls = 6
|
||||
)
|
||||
|
||||
proc := New(Params{MaxConcurrentProcessing: limit})
|
||||
input := createTestJPEG(t, 50, 50)
|
||||
|
||||
counter := &readingCounter{}
|
||||
gate := make(chan struct{})
|
||||
entered := make(chan struct{}, calls)
|
||||
results := make(chan error, calls)
|
||||
|
||||
openGate := sync.OnceFunc(func() { close(gate) })
|
||||
t.Cleanup(openGate)
|
||||
|
||||
for range calls {
|
||||
processInBackground(proc, &gatedReader{
|
||||
data: bytes.NewReader(input), gate: gate, entered: entered,
|
||||
counter: counter,
|
||||
}, results)
|
||||
}
|
||||
|
||||
waitForEntries(t, entered, limit)
|
||||
|
||||
// A call beyond the limit would start reading its input now.
|
||||
select {
|
||||
case <-entered:
|
||||
t.Fatalf("a Process call started while %d were already working", limit)
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
|
||||
openGate()
|
||||
|
||||
for range calls {
|
||||
err := <-results
|
||||
if err != nil {
|
||||
t.Errorf("Process() error = %v, want nil once a slot is free", err)
|
||||
}
|
||||
}
|
||||
|
||||
if most := counter.mostReading(); most > limit {
|
||||
t.Errorf("%d Process calls worked at once, want at most %d", most, limit)
|
||||
}
|
||||
}
|
||||
|
||||
// TestProcessWaitsThenFailsWhenNoSlotFrees holds the only slot and checks
|
||||
// that another call waits the whole wait timeout, then fails with
|
||||
// ErrTooManyImages instead of processing anyway.
|
||||
func TestProcessWaitsThenFailsWhenNoSlotFrees(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
proc := New(Params{MaxConcurrentProcessing: 1})
|
||||
proc.processingWaitTimeout = 100 * time.Millisecond
|
||||
|
||||
input := createTestJPEG(t, 10, 10)
|
||||
|
||||
gate := make(chan struct{})
|
||||
entered := make(chan struct{}, 1)
|
||||
held := make(chan error, 1)
|
||||
|
||||
openGate := sync.OnceFunc(func() { close(gate) })
|
||||
t.Cleanup(openGate)
|
||||
|
||||
processInBackground(proc, &gatedReader{
|
||||
data: bytes.NewReader(input), gate: gate, entered: entered,
|
||||
counter: &readingCounter{},
|
||||
}, held)
|
||||
waitForEntries(t, entered, 1)
|
||||
|
||||
start := time.Now()
|
||||
|
||||
_, err := proc.Process(context.Background(), bytes.NewReader(input),
|
||||
smallJPEGRequest())
|
||||
if !errors.Is(err, ErrTooManyImages) {
|
||||
t.Fatalf("Process() error = %v, want ErrTooManyImages", err)
|
||||
}
|
||||
|
||||
if waited := time.Since(start); waited < proc.processingWaitTimeout {
|
||||
t.Errorf("Process() failed after %v, before waiting %v",
|
||||
waited, proc.processingWaitTimeout)
|
||||
}
|
||||
|
||||
openGate()
|
||||
|
||||
err = <-held
|
||||
if err != nil {
|
||||
t.Errorf("Process() holding the slot: error = %v, want nil", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestProcessReleasesSlotOnError checks that Process gives its slot back
|
||||
// when it fails, whether it fails early or late: with one slot, the slot
|
||||
// must be free after the failure and the next call must succeed.
|
||||
func TestProcessReleasesSlotOnError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
valid := createTestJPEG(t, 10, 10)
|
||||
|
||||
unsupported := smallJPEGRequest()
|
||||
unsupported.Format = "bmp"
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
input io.Reader
|
||||
req *Request
|
||||
// want is the error Process must return; nil means any error.
|
||||
want error
|
||||
}{
|
||||
{
|
||||
name: "input cannot be read",
|
||||
input: iotest.ErrReader(errTestReadFailed),
|
||||
req: smallJPEGRequest(),
|
||||
want: errTestReadFailed,
|
||||
},
|
||||
{
|
||||
name: "input over the byte limit",
|
||||
input: bytes.NewReader(createTestJPEG(t, 800, 600)),
|
||||
req: smallJPEGRequest(),
|
||||
want: ErrInputDataTooLarge,
|
||||
},
|
||||
{
|
||||
name: "input not an image",
|
||||
input: strings.NewReader("not an image"),
|
||||
req: smallJPEGRequest(),
|
||||
},
|
||||
{
|
||||
name: "output format not supported",
|
||||
input: bytes.NewReader(valid),
|
||||
req: unsupported,
|
||||
want: ErrUnsupportedOutputFormat,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
proc := New(Params{MaxInputBytes: 4096, MaxConcurrentProcessing: 1})
|
||||
proc.processingWaitTimeout = 100 * time.Millisecond
|
||||
|
||||
_, err := proc.Process(context.Background(), tc.input, tc.req)
|
||||
if err == nil || (tc.want != nil && !errors.Is(err, tc.want)) {
|
||||
t.Fatalf("Process() error = %v, want %v", err, tc.want)
|
||||
}
|
||||
|
||||
if held := len(proc.processingSemaphore); held != 0 {
|
||||
t.Fatalf("slot still held after the error: %d held", held)
|
||||
}
|
||||
|
||||
result, err := proc.Process(context.Background(), bytes.NewReader(valid),
|
||||
smallJPEGRequest())
|
||||
if err != nil {
|
||||
t.Fatalf("Process() after the error = %v, want nil", err)
|
||||
}
|
||||
|
||||
_ = result.Content.Close()
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user