@@ -41,6 +41,7 @@ import (
4141 "go.opentelemetry.io/otel"
4242 "go.opentelemetry.io/otel/trace"
4343 "golang.org/x/sync/errgroup"
44+ "golang.org/x/sync/semaphore"
4445
4546 "github.com/owncloud/reva/v2/pkg/appctx"
4647 "github.com/owncloud/reva/v2/pkg/autoprop"
@@ -133,7 +134,7 @@ type Decomposedfs struct {
133134 spaceTypeIndex * spaceidindex.Index
134135
135136 // commitLimiter caps concurrent async blob commits at NumConsumers.
136- commitLimiter chan struct {}
137+ commitLimiter * semaphore. Weighted
137138
138139 log * zerolog.Logger
139140}
@@ -281,7 +282,7 @@ func New(o *options.Options, aspects aspects.Aspects, log *zerolog.Logger) (stor
281282 o .Events .NumConsumers = 1
282283 }
283284
284- fs .commitLimiter = make ( chan struct {}, o .Events .NumConsumers )
285+ fs .commitLimiter = semaphore . NewWeighted ( int64 ( o .Events .NumConsumers ) )
285286
286287 for i := 0 ; i < o .Events .NumConsumers ; i ++ {
287288 go fs .Postprocessing (ch )
@@ -303,23 +304,32 @@ func (fs *Decomposedfs) Postprocessing(ch <-chan events.Event) {
303304// finalizeWithRetry commits the staged bytes, retrying blobstore failures with
304305// capped exponential backoff.
305306func (fs * Decomposedfs ) finalizeWithRetry (ctx context.Context , session * upload.OcisSession , log * zerolog.Logger ) error {
307+ maxAttempts := fs .o .Events .CommitMaxRetries + 1
306308 backoff := fs .o .Events .CommitRetryBackoff
307309 var err error
308- for attempt := 0 ; attempt <= fs . o . Events . CommitMaxRetries ; attempt ++ {
310+ for attempt := 1 ; attempt <= maxAttempts ; attempt ++ {
309311 if err = session .Finalize (ctx ); err == nil {
310312 return nil
311313 }
312- if attempt == fs .o .Events .CommitMaxRetries {
314+ ev := log .Warn ().Err (err ).
315+ Int ("attempt" , attempt ).Int ("maxAttempts" , maxAttempts ).
316+ Str ("spaceid" , session .SpaceID ()).Str ("nodeid" , session .NodeID ())
317+ if attempt == maxAttempts {
318+ ev .Msg ("blob commit failed, giving up" )
313319 break
314320 }
321+ // clamp before use: this bounds backoff to maxCommitRetryBackoff every
322+ // iteration, so the backoff *= 2 below can never grow past 2*max and
323+ // cannot overflow the int64 duration.
315324 backoff = min (backoff , maxCommitRetryBackoff )
316- log . Warn (). Err ( err ). Int ( "attempt" , attempt + 1 ) .Dur ("backoff" , backoff ).Msg ("blob commit failed, retrying after backoff" )
325+ ev .Dur ("backoff" , backoff ).Msg ("blob commit failed, retrying after backoff" )
317326 timer := time .NewTimer (backoff )
318327 select {
319328 case <- ctx .Done ():
320329 timer .Stop ()
321330 return ctx .Err ()
322331 case <- timer .C :
332+ // backoff elapsed, fall through to the next attempt
323333 }
324334 backoff *= 2
325335 }
@@ -443,8 +453,11 @@ func (fs *Decomposedfs) processEvent(evCtx context.Context, event events.Event,
443453 // commit re-uploads the whole file and can block for the retry window;
444454 // run it detached (bounded by commitLimiter) to not stall the consumer
445455 go func () {
446- fs .commitLimiter <- struct {}{}
447- defer func () { <- fs .commitLimiter }()
456+ if err := fs .commitLimiter .Acquire (ctx , 1 ); err != nil {
457+ sublog .Error ().Err (err ).Msg ("could not acquire commit slot" )
458+ return
459+ }
460+ defer fs .commitLimiter .Release (1 )
448461
449462 if err := fs .finalizeWithRetry (ctx , session , & sublog ); err != nil {
450463 sublog .Error ().Err (err ).Msg ("could not finalize upload after retries, reverting to a recoverable failed state" )
0 commit comments