From 081438092c70b2d7feaf9c561c8a9efc2f4b7f94 Mon Sep 17 00:00:00 2001 From: Kristina Kovalevskaya Date: Fri, 6 Oct 2017 17:30:15 +0300 Subject: [PATCH] Execute all handlers for each lifecycle phase concurrently --- lifecycle.go | 78 ++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 63 insertions(+), 15 deletions(-) diff --git a/lifecycle.go b/lifecycle.go index 7b63fe3..7832ecf 100644 --- a/lifecycle.go +++ b/lifecycle.go @@ -15,6 +15,15 @@ var ( dResolverHost = "consul.service" ) +// HandlerResponse stores a result of Handler execution +type handlerResponse struct { + meta map[string]interface{} + err error + eType types.EventType + conf *types.HandlerConfig + ignoreErr bool +} + // Lifecycle implements a Lifecycle that calls multiple lifecycles for an event. type Lifecycle struct { ctx *types.Context @@ -100,6 +109,8 @@ func (lc *Lifecycle) Begin(ctx *types.Context) error { return nil } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ Type: types.EventTypeBegin, @@ -107,18 +118,28 @@ func (lc *Lifecycle) Begin(ctx *types.Context) error { Timestamp: time.Now().UnixNano(), } - meta, err := v.Handle(event) - if err != nil { - if v.conf.IgnoreErrors { - log.Printf("[ERROR] phase=%s handler=%s %v", event.Type, v.conf.Type, err) + go func(v *phaseHandler) { + meta, err := v.Handle(event) + doneCh <- &handlerResponse{meta: meta, err: err, eType: event.Type, conf: v.conf, ignoreErr: v.conf.IgnoreErrors} + }(v) + } + + var failedHandlers int + for i := 0; i < len(handlers); i++ { + res := <-doneCh + if res.err != nil { + log.Printf("[ERROR] phase=%s handler=%s %v", res.eType, res.conf.Type, res.err) + if res.ignoreErr { continue } - return err + failedHandlers++ } - - lc.applyContext(meta, v.conf) + lc.applyContext(res.meta, res.conf) } + if failedHandlers > 0 { + return fmt.Errorf("%d handlers were finished with error", failedHandlers) + } return nil } @@ -130,6 +151,8 @@ func (lc *Lifecycle) Progress(line []byte) { return } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ Type: types.EventTypeProgress, @@ -137,8 +160,16 @@ func (lc *Lifecycle) Progress(line []byte) { Data: line, Timestamp: time.Now().UnixNano(), } - if _, err := v.Handle(event); err != nil { - log.Printf("[ERROR] phase=%s handler=%s %v", event.Type, v.conf.Type, err) + go func(v *phaseHandler) { + _, err := v.Handle(event) + doneCh <- &handlerResponse{err: err, eType: event.Type, conf: v.conf} + }(v) + } + + for i := 0; i < len(handlers); i++ { + res := <-doneCh + if res.err != nil { + log.Printf("[ERROR] phase=%s handler=%s %v", res.eType, res.conf.Type, res.err) } } } @@ -152,6 +183,8 @@ func (lc *Lifecycle) Failed(result *types.ChildResult) { return } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ @@ -160,11 +193,17 @@ func (lc *Lifecycle) Failed(result *types.ChildResult) { Data: result, Timestamp: time.Now().UnixNano(), } + go func(v *phaseHandler) { + _, err := v.Handle(event) + doneCh <- &handlerResponse{err: err, eType: event.Type, conf: v.conf} + }(v) + } - if _, err := v.Handle(event); err != nil { - log.Printf("[ERROR] phase=%s handler=%s %v", event.Type, v.conf.Type, err) + for i := 0; i < len(handlers); i++ { + res := <-doneCh + if res.err != nil { + log.Printf("[ERROR] phase=%s handler=%s %v", res.eType, res.conf.Type, res.err) } - } } @@ -176,6 +215,8 @@ func (lc *Lifecycle) Completed(result *types.ChildResult) { return } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ @@ -189,10 +230,17 @@ func (lc *Lifecycle) Completed(result *types.ChildResult) { // event.Data = result.Stdout //} - if _, err := v.Handle(event); err != nil { - log.Printf("[ERROR] phase=%s handler=%s %v", event.Type, v.conf.Type, err) - } + go func(v *phaseHandler) { + _, err := v.Handle(event) + doneCh <- &handlerResponse{err: err, eType: event.Type, conf: v.conf} + }(v) + } + for i := 0; i < len(handlers); i++ { + res := <-doneCh + if res.err != nil { + log.Printf("[ERROR] phase=%s handler=%s %v", res.eType, res.conf.Type, res.err) + } } }