Skip to content
This repository was archived by the owner on May 22, 2026. It is now read-only.
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
78 changes: 63 additions & 15 deletions lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -118,25 +127,37 @@ 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,
Meta: ctx.Meta,
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
}

Expand All @@ -148,15 +169,25 @@ func (lc *Lifecycle) Progress(line []byte) {
return
}

doneCh := make(chan *handlerResponse, len(handlers))

for _, v := range handlers {
event := &types.Event{
Type: types.EventTypeProgress,
Meta: lc.ctx.Meta,
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)
}
}
}
Expand All @@ -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{
Expand All @@ -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)
}

}
}

Expand Down Expand Up @@ -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{
Expand All @@ -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)
}
}
}

Expand Down