diff --git a/lifecycle.go b/lifecycle.go index 09588ca..2065869 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 @@ -118,6 +127,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, @@ -125,18 +136,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 } @@ -148,6 +169,8 @@ func (lc *Lifecycle) Progress(line []byte) { return } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ Type: types.EventTypeProgress, @@ -155,8 +178,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) } } } @@ -170,6 +201,8 @@ func (lc *Lifecycle) Failed(result *types.ChildResult) { return } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ @@ -178,11 +211,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) } - } } @@ -219,6 +258,8 @@ func (lc *Lifecycle) Completed(result *types.ChildResult) { return } + doneCh := make(chan *handlerResponse, len(handlers)) + for _, v := range handlers { event := &types.Event{ @@ -232,10 +273,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) + } } }