From f094903f75a0dcbec2a611c1a3cf2bbb23ea70cf Mon Sep 17 00:00:00 2001 From: Tiago Silva Date: Mon, 14 Aug 2023 17:22:44 +0100 Subject: [PATCH] Fix data race when calling Uploader's `Close` and `Serve` simultaneously (#30360) * Fix data race when calling Uploader's `Close` and `Serve` at the same time This PR fixes a data race that was introduced by #30211 while trying to fix the uploader write after being closed. The data race happens when `Serve` and `Close` methods are called at the same time - i.e. `Serve` is spawned by `RegisterFunc` and `Close` is called via `OnExit` function. There is no guarentee that when `OnExit` is called `Serve` is already spawn and running. This creates a data race when `Close` blocks on `wg.Wait()` but `Serve` still tries to call `wg.Add(1)` into the waitgroup. To prevent the data race, we added a `isClosing` field that controls when the server is marked to be closing and we call `wg.Add(1)` if `isClosing=false` and under locking to prevent the race between `Close` and `Serve`. Signed-off-by: Tiago Silva * add todo message --------- Signed-off-by: Tiago Silva --- lib/events/filesessions/fileasync.go | 28 ++++++++++++++++++++++++---- 1 file changed, 24 insertions(+), 4 deletions(-) diff --git a/lib/events/filesessions/fileasync.go b/lib/events/filesessions/fileasync.go index 2c0ab7ac7e9..187b1b0713b 100644 --- a/lib/events/filesessions/fileasync.go +++ b/lib/events/filesessions/fileasync.go @@ -110,7 +110,6 @@ func NewUploader(cfg UploaderConfig) (*Uploader, error) { eventsCh: make(chan events.UploadEvent, cfg.ConcurrentUploads), eventPreparer: &events.NoOpPreparer{}, } - return uploader, nil } @@ -129,14 +128,21 @@ type Uploader struct { cfg UploaderConfig log *log.Entry - eventsCh chan events.UploadEvent - closeC chan struct{} - wg sync.WaitGroup + eventsCh chan events.UploadEvent + closeC chan struct{} + wg sync.WaitGroup + mu sync.Mutex + isClosing bool eventPreparer *events.NoOpPreparer } func (u *Uploader) Close() { + // TODO(tigrato): prevent close to be called before Serve starts. + u.mu.Lock() + u.isClosing = true + u.mu.Unlock() + close(u.closeC) // wait for all uploads to finish u.wg.Wait() @@ -167,8 +173,22 @@ func (u *Uploader) checkSessionError(sessionID session.ID) (bool, error) { // Serve runs the uploader until stopped func (u *Uploader) Serve(ctx context.Context) error { + // Check if close operation is already in progress. + // We need to do this because Serve is spawned in a goroutine + // and Close can be called before Serve starts which ends up in a data + // race because Close is waiting for wg to be 0 and Serve is adding to wg. + // To avoid this, we check if Close is already in progress and return + // immediately. If Close is not in progress, we add to wg under the mutex + // lock to ensure that Close can't reach wg.Wait() before Serve adds to wg. + u.mu.Lock() + if u.isClosing { + u.mu.Unlock() + return nil + } u.wg.Add(1) + u.mu.Unlock() defer u.wg.Done() + u.log.Infof("uploader will scan %v every %v", u.cfg.ScanDir, u.cfg.ScanPeriod.String()) backoff, err := retryutils.NewLinear(retryutils.LinearConfig{ Step: u.cfg.ScanPeriod,