Count interrupted misses and the upstream bytes they read (closes #56)
check / check (push) Successful in 3m2s
check / check (push) Successful in 3m2s
The miss and transform counters are written with context.WithoutCancel, so a client disconnect or the request timeout during or after the work no longer loses them. A failed read of the upstream body now returns the bytes read before the error, so an over-size or cut-off body still moves the upstream fetch counters. processFromSourceOrFetch passes the cached source's length directly instead of through a local named fetchBytes. Model: opus-5-5
This commit is contained in:
@@ -156,12 +156,13 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
||||
}
|
||||
|
||||
// Cache miss - process the cached source or fetch it, then count the
|
||||
// miss with the bytes it fetched from upstream, also when it failed
|
||||
// miss with the bytes it fetched from upstream, also when it failed or
|
||||
// the request context has ended meanwhile
|
||||
cacheKey := CacheKey(req)
|
||||
|
||||
response, fetchedBytes, err := s.processFromSourceOrFetch(ctx, req, cacheKey)
|
||||
|
||||
s.cache.IncrementStats(ctx, false, fetchedBytes)
|
||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, fetchedBytes)
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -283,10 +284,7 @@ func (s *Service) processFromSourceOrFetch(
|
||||
s.log.Warn("source lookup failed", "error", err)
|
||||
}
|
||||
|
||||
var (
|
||||
sourceData []byte
|
||||
fetchBytes int64
|
||||
)
|
||||
var sourceData []byte
|
||||
|
||||
if contentHash != "" {
|
||||
s.log.Debug("using cached source", "hash", contentHash)
|
||||
@@ -299,16 +297,16 @@ func (s *Service) processFromSourceOrFetch(
|
||||
}
|
||||
|
||||
// Process using cached source; nothing was fetched from upstream
|
||||
fetchBytes = int64(len(sourceData))
|
||||
|
||||
resp, err := s.processAndStore(ctx, req, cacheKey, sourceData, fetchBytes)
|
||||
resp, err := s.processAndStore(
|
||||
ctx, req, cacheKey, sourceData, int64(len(sourceData)),
|
||||
)
|
||||
|
||||
return resp, 0, err
|
||||
}
|
||||
|
||||
// fetchAndProcess fetches from upstream, processes, and caches the result.
|
||||
// It also returns the number of bytes fetched from upstream once the
|
||||
// response has been read, including when a later step fails.
|
||||
// It also returns the number of bytes read from upstream, including when
|
||||
// reading the response or a later step fails.
|
||||
func (s *Service) fetchAndProcess(
|
||||
ctx context.Context,
|
||||
req *ImageRequest,
|
||||
@@ -339,13 +337,13 @@ func (s *Service) fetchAndProcess(
|
||||
|
||||
// Read and validate the source content
|
||||
sourceData, err := io.ReadAll(fetchResult.Content)
|
||||
fetchBytes := int64(len(sourceData))
|
||||
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("failed to read upstream response: %w", err)
|
||||
return nil, fetchBytes, fmt.Errorf("failed to read upstream response: %w", err)
|
||||
}
|
||||
|
||||
// Calculate download bitrate
|
||||
fetchBytes := int64(len(sourceData))
|
||||
|
||||
var downloadRate string
|
||||
|
||||
if fetchResult.FetchDurationMs > 0 {
|
||||
@@ -410,7 +408,8 @@ func (s *Service) processAndStore(
|
||||
|
||||
processDuration := time.Since(processStart)
|
||||
|
||||
s.cache.IncrementTransformCount(ctx)
|
||||
// Counted also when the request context has ended meanwhile
|
||||
s.cache.IncrementTransformCount(context.WithoutCancel(ctx))
|
||||
|
||||
// Read processed content
|
||||
processedData, err := io.ReadAll(processResult.Content)
|
||||
|
||||
Reference in New Issue
Block a user