diff --git a/pkg/esign/completion_certificate_worker.go b/pkg/esign/completion_certificate_worker.go index ce30fbe0d..fca4dd584 100644 --- a/pkg/esign/completion_certificate_worker.go +++ b/pkg/esign/completion_certificate_worker.go @@ -115,12 +115,11 @@ LOOP: case <-time.After(w.interval): // From there we should not accept cancelations anymore. nonCancelableCtx := context.WithoutCancel(ctx) - w.recoverStaleCertificateRows(nonCancelableCtx) for { - if err := w.processNext(nonCancelableCtx, sem, &wg); err != nil { + if err := w.processNext(ctx, sem, &wg); err != nil { if !errors.Is(err, coredata.ErrResourceNotFound) { - w.logger.ErrorCtx(nonCancelableCtx, "cannot process certificate", log.Error(err)) + w.logger.ErrorCtx(ctx, "cannot process certificate", log.Error(err)) } break } @@ -132,26 +131,29 @@ LOOP: func (w *CompletionCertificateWorker) processNext(ctx context.Context, sem chan struct{}, wg *sync.WaitGroup) error { select { case sem <- struct{}{}: - case <-ctx.Done(): + case <-ctx.Done(): // FIXME: this will never be fired return ctx.Err() } var ( signature coredata.ElectronicSignature now = time.Now() + + // From there we should not accept cancelations anymore. + nonCancelableCtx = context.WithoutCancel(ctx) ) if err := w.pg.WithTx( - ctx, + nonCancelableCtx, func(tx pg.Conn) error { - if err := signature.LoadNextCompletedWithoutCertificateForUpdate(ctx, tx); err != nil { + if err := signature.LoadNextCompletedWithoutCertificateForUpdate(nonCancelableCtx, tx); err != nil { return err } scope := coredata.NewScopeFromObjectID(signature.ID) signature.CertificateProcessingStartedAt = &now signature.UpdatedAt = now - if err := signature.Update(ctx, tx, scope); err != nil { + if err := signature.Update(nonCancelableCtx, tx, scope); err != nil { return fmt.Errorf("cannot update signature: %w", err) } @@ -169,9 +171,9 @@ func (w *CompletionCertificateWorker) processNext(ctx context.Context, sem chan scope := coredata.NewScopeFromObjectID(signature.ID) - if err := w.generateAndCommit(ctx, &signature); err != nil { - if err := w.handleCertFailure(ctx, &signature, scope, err); err != nil { - w.logger.ErrorCtx(ctx, "cannot handle certificate failure", log.Error(err)) + if err := w.generateAndCommit(nonCancelableCtx, &signature); err != nil { + if err := w.handleCertFailure(nonCancelableCtx, &signature, scope, err); err != nil { + w.logger.ErrorCtx(nonCancelableCtx, "cannot handle certificate failure", log.Error(err)) } } }(signature) diff --git a/pkg/esign/sealing_worker.go b/pkg/esign/sealing_worker.go index fb8f29e29..d23a01e62 100644 --- a/pkg/esign/sealing_worker.go +++ b/pkg/esign/sealing_worker.go @@ -114,10 +114,10 @@ LOOP: case <-time.After(w.interval): // From there we should not accept cancelations anymore. nonCancelableCtx := context.WithoutCancel(ctx) - w.recoverStaleRows(nonCancelableCtx) + for { - if err := w.processNext(nonCancelableCtx, sem, &wg); err != nil { + if err := w.processNext(ctx, sem, &wg); err != nil { if !errors.Is(err, coredata.ErrResourceNotFound) { w.logger.ErrorCtx(nonCancelableCtx, "cannot claim signature", log.Error(err)) } @@ -139,12 +139,15 @@ func (w *SealingWorker) processNext(ctx context.Context, sem chan struct{}, wg * var ( signature = coredata.ElectronicSignature{} now = time.Now() + + // From there we should not accept cancelations anymore. + nonCancelableCtx = context.WithoutCancel(ctx) ) if err := w.pg.WithTx( - ctx, + nonCancelableCtx, func(tx pg.Conn) error { - if err := signature.LoadNextAcceptedForUpdateSkipLocked(ctx, tx); err != nil { + if err := signature.LoadNextAcceptedForUpdateSkipLocked(nonCancelableCtx, tx); err != nil { return err } @@ -153,7 +156,7 @@ func (w *SealingWorker) processNext(ctx context.Context, sem chan struct{}, wg * signature.AttemptCount++ signature.LastAttemptedAt = &now signature.UpdatedAt = now - if err := signature.Update(ctx, tx, coredata.NewNoScope()); err != nil { + if err := signature.Update(nonCancelableCtx, tx, coredata.NewNoScope()); err != nil { return fmt.Errorf("cannot update signature: %w", err) } @@ -169,9 +172,9 @@ func (w *SealingWorker) processNext(ctx context.Context, sem chan struct{}, wg * defer wg.Done() defer func() { <-sem }() - if err := w.sealAndCommit(ctx, &signature); err != nil { - if err := w.failSignature(ctx, &signature, err); err != nil { - w.logger.ErrorCtx(ctx, "cannot fail signature", log.Error(err)) + if err := w.sealAndCommit(nonCancelableCtx, &signature); err != nil { + if err := w.failSignature(nonCancelableCtx, &signature, err); err != nil { + w.logger.ErrorCtx(nonCancelableCtx, "cannot fail signature", log.Error(err)) } } }(signature)