make image.Close wait for all images Submit completion
This commit is contained in:
committed by
Umputun
parent
7151698b18
commit
4605d567ae
@@ -33,11 +33,12 @@ import (
|
||||
type Service struct {
|
||||
ServiceParams
|
||||
|
||||
store Store
|
||||
wg sync.WaitGroup
|
||||
submitCh chan submitReq
|
||||
once sync.Once
|
||||
term int32 // term value used atomically to detect emergency termination
|
||||
store Store
|
||||
wg sync.WaitGroup
|
||||
submitCh chan submitReq
|
||||
once sync.Once
|
||||
term int32 // term value used atomically to detect emergency termination
|
||||
submitCount int32 // atomic increment for counting submitted images
|
||||
}
|
||||
|
||||
// ServiceParams contains externally adjustable parameters of Service
|
||||
@@ -88,10 +89,9 @@ func (s *Service) Submit(idsFn func() []string) {
|
||||
s.wg.Add(1)
|
||||
go func() {
|
||||
defer s.wg.Done()
|
||||
atomic.StoreInt32(&s.term, 1)
|
||||
for req := range s.submitCh {
|
||||
// wait for TTL expiration with emergency pass on term
|
||||
for atomic.LoadInt32(&s.term) == 1 && time.Since(req.TS) <= s.TTL/2 { // commit on a half of TTL
|
||||
for atomic.LoadInt32(&s.term) == 0 && time.Since(req.TS) <= s.TTL/2 { // commit on a half of TTL
|
||||
time.Sleep(time.Millisecond * 10) // small sleep to relive busy wait but keep reactive for term (close)
|
||||
}
|
||||
for _, id := range req.idsFn() {
|
||||
@@ -99,12 +99,13 @@ func (s *Service) Submit(idsFn func() []string) {
|
||||
log.Printf("[WARN] failed to commit image %s", id)
|
||||
}
|
||||
}
|
||||
atomic.StoreInt32(&s.term, 1) // indicates completion of ids commits
|
||||
atomic.AddInt32(&s.submitCount, -1)
|
||||
}
|
||||
log.Printf("[INFO] image submitter terminated")
|
||||
}()
|
||||
})
|
||||
|
||||
atomic.AddInt32(&s.submitCount, 1)
|
||||
s.submitCh <- submitReq{idsFn: idsFn, TS: time.Now()}
|
||||
}
|
||||
|
||||
@@ -159,18 +160,15 @@ func (s *Service) Close(ctx context.Context) {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
if atomic.LoadInt32(&s.term) == 1 { // set to 1 by commit goroutine after everything waited on TTL sent
|
||||
if atomic.LoadInt32(&s.submitCount) == 0 {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// terminate commit goroutine using s.term only if it is started
|
||||
if atomic.LoadInt32(&s.term) == 1 {
|
||||
atomic.StoreInt32(&s.term, 0) // enforce non-delayed commits for all ids left in submitCh
|
||||
waitForTerm(ctx)
|
||||
}
|
||||
atomic.StoreInt32(&s.term, 1) // enforce non-delayed commits for all ids left in submitCh
|
||||
waitForTerm(ctx)
|
||||
|
||||
if s.submitCh != nil {
|
||||
close(s.submitCh)
|
||||
|
||||
@@ -119,12 +119,13 @@ func TestService_Submit(t *testing.T) {
|
||||
store.AssertNumberOfCalls(t, "Commit", 0)
|
||||
time.Sleep(time.Millisecond * 150)
|
||||
store.AssertNumberOfCalls(t, "Commit", 5)
|
||||
svc.Close(context.TODO())
|
||||
}
|
||||
|
||||
func TestService_Close(t *testing.T) {
|
||||
store := MockStore{}
|
||||
store.On("Commit", mock.Anything, mock.Anything).Times(5).Return(nil)
|
||||
svc := Service{store: &store, ServiceParams: ServiceParams{ImageAPI: "/blah/", TTL: time.Millisecond * 500}}
|
||||
svc := Service{store: &store, ServiceParams: ServiceParams{ImageAPI: "/blah/", TTL: time.Hour * 24}}
|
||||
svc.Submit(func() []string { return []string{"id1", "id2", "id3"} })
|
||||
svc.Submit(func() []string { return []string{"id4", "id5"} })
|
||||
svc.Submit(nil)
|
||||
|
||||
Reference in New Issue
Block a user