Stop waiting for count writes after a request's deadline (closes #224)
check / check (push) Waiting to run
check / check (push) Waiting to run
TestService_Get_ReturnsByItsDeadline failed on a busy host because its request, past its deadline, still waited for its miss count to be written, a database write with no deadline; in pixad that write waits for the one database connection every request shares. The hit, miss, upstream fetch and transform counts are now each written in a goroutine of their own, which the request waits for only until its deadline; a count not written by then is written after it returns. Shutdown waits for those writes within ShutdownTimeout and reports any left unfinished. The test phase also runs go test with -parallel 4: on a busy host, as many tests at once as there are CPUs wait so long to be scheduled that a timed request can fail. Model: opus-5-5
This commit is contained in:
@@ -9,6 +9,7 @@ import (
|
||||
"log/slog"
|
||||
"net/url"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/dustin/go-humanize"
|
||||
@@ -36,6 +37,9 @@ type Service struct {
|
||||
// variantsInProgress lets the requests that miss the same variant at the
|
||||
// same time share one fetch and one transcode.
|
||||
variantsInProgress singleflight.Group
|
||||
// countWrites holds the count writes writeCount has started, for
|
||||
// WaitForCountWrites.
|
||||
countWrites sync.WaitGroup
|
||||
}
|
||||
|
||||
// ServiceConfig holds configuration for the image service.
|
||||
@@ -161,8 +165,9 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
||||
s.log.Error("failed to get cached variant", "key", result.CacheKey, "error", err)
|
||||
// Fall through to re-process
|
||||
} else {
|
||||
// Counted also when the request context has ended meanwhile
|
||||
s.cache.IncrementStats(context.WithoutCancel(ctx), true, 0)
|
||||
s.writeCount(ctx, func(ctx context.Context) {
|
||||
s.cache.IncrementStats(ctx, true, 0)
|
||||
})
|
||||
|
||||
return &ImageResponse{
|
||||
Content: reader,
|
||||
@@ -179,7 +184,9 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
||||
// failed or the request context has ended meanwhile
|
||||
response, err := s.processOrWait(ctx, req)
|
||||
|
||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0)
|
||||
s.writeCount(ctx, func(ctx context.Context) {
|
||||
s.cache.IncrementStats(ctx, false, 0)
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -208,6 +215,25 @@ func (s *Service) WaitForProcessing(ctx context.Context) int {
|
||||
return s.processor.WaitForProcessing(ctx)
|
||||
}
|
||||
|
||||
// WaitForCountWrites waits until every count a request has started writing
|
||||
// to the database is written, or until ctx ends, and reports whether they
|
||||
// all were.
|
||||
func (s *Service) WaitForCountWrites(ctx context.Context) bool {
|
||||
written := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
s.countWrites.Wait()
|
||||
close(written)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-written:
|
||||
return true
|
||||
case <-ctx.Done():
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// ValidateRequest validates the request signature if required.
|
||||
func (s *Service) ValidateRequest(req *ImageRequest) error {
|
||||
// Check if host is allowed (no signature required)
|
||||
@@ -254,6 +280,38 @@ func (s *Service) GenerateSignedURL(
|
||||
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
|
||||
}
|
||||
|
||||
// writeCount runs write, which adds to a count in the database, in a
|
||||
// goroutine of its own. The write gets ctx without its cancellation or
|
||||
// deadline, so an ended request is still counted. writeCount waits for it
|
||||
// until ctx's deadline, even if ctx is cancelled first, so a request
|
||||
// returns with its count written unless its deadline has passed; past
|
||||
// that, it does not wait for the database connection, which every request
|
||||
// shares.
|
||||
func (s *Service) writeCount(ctx context.Context, write func(context.Context)) {
|
||||
written := make(chan struct{})
|
||||
|
||||
s.countWrites.Go(func() {
|
||||
defer close(written)
|
||||
|
||||
write(context.WithoutCancel(ctx))
|
||||
})
|
||||
|
||||
deadline, hasDeadline := ctx.Deadline()
|
||||
if !hasDeadline {
|
||||
<-written
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
deadlineTimer := time.NewTimer(time.Until(deadline))
|
||||
defer deadlineTimer.Stop()
|
||||
|
||||
select {
|
||||
case <-written:
|
||||
case <-deadlineTimer.C:
|
||||
}
|
||||
}
|
||||
|
||||
// errPanicked is returned when processing a variant panicked.
|
||||
var errPanicked = errors.New("panic while processing image")
|
||||
|
||||
@@ -452,8 +510,9 @@ func (s *Service) fetchAndProcess(
|
||||
sourceData, err := io.ReadAll(fetchResult.Content)
|
||||
fetchBytes := int64(len(sourceData))
|
||||
|
||||
// Counted also when the request context has ended meanwhile
|
||||
s.cache.IncrementUpstreamFetch(context.WithoutCancel(ctx), fetchBytes)
|
||||
s.writeCount(ctx, func(ctx context.Context) {
|
||||
s.cache.IncrementUpstreamFetch(ctx, fetchBytes)
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to read upstream response: %w", err)
|
||||
@@ -534,8 +593,7 @@ func (s *Service) processAndStore(
|
||||
|
||||
processDuration := time.Since(processStart)
|
||||
|
||||
// Counted also when the request context has ended meanwhile
|
||||
s.cache.IncrementTransformCount(context.WithoutCancel(ctx))
|
||||
s.writeCount(ctx, s.cache.IncrementTransformCount)
|
||||
|
||||
// Read processed content
|
||||
processedData, err := io.ReadAll(processResult.Content)
|
||||
|
||||
Reference in New Issue
Block a user