Workflow execution logging is now configurable
Use WORKFLOW_EXEC_DEBUG=true to enable it
This commit is contained in:
@@ -74,7 +74,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config)
|
||||
// store interface exposed or generated inside app package
|
||||
DefaultStore = s
|
||||
|
||||
DefaultLogger = log.Named("service")
|
||||
DefaultLogger = log.Named("workflow")
|
||||
|
||||
{
|
||||
tee := zap.NewNop()
|
||||
@@ -92,7 +92,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config)
|
||||
|
||||
DefaultAccessControl = AccessControl(rbac.Global())
|
||||
|
||||
DefaultSession = Session(DefaultLogger.Named("session"))
|
||||
DefaultSession = Session(DefaultLogger.Named("session"), c.Workflow)
|
||||
DefaultWorkflow = Workflow(DefaultLogger.Named("workflow"))
|
||||
DefaultTrigger = Trigger(DefaultLogger.Named("trigger"), c.Workflow)
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"github.com/cortezaproject/corteza-server/pkg/auth"
|
||||
"github.com/cortezaproject/corteza-server/pkg/errors"
|
||||
"github.com/cortezaproject/corteza-server/pkg/expr"
|
||||
"github.com/cortezaproject/corteza-server/pkg/options"
|
||||
"github.com/cortezaproject/corteza-server/pkg/sentry"
|
||||
"github.com/cortezaproject/corteza-server/pkg/wfexec"
|
||||
"github.com/cortezaproject/corteza-server/store"
|
||||
@@ -20,6 +21,7 @@ type (
|
||||
store store.Storer
|
||||
actionlog actionlog.Recorder
|
||||
ac sessionAccessController
|
||||
opt options.WorkflowOpt
|
||||
log *zap.Logger
|
||||
mux *sync.RWMutex
|
||||
pool map[uint64]*types.Session
|
||||
@@ -40,9 +42,10 @@ type (
|
||||
WaitFn func(ctx context.Context) (*expr.Vars, wfexec.SessionStatus, types.Stacktrace, error)
|
||||
)
|
||||
|
||||
func Session(log *zap.Logger) *session {
|
||||
func Session(log *zap.Logger, opt options.WorkflowOpt) *session {
|
||||
return &session{
|
||||
log: log,
|
||||
opt: opt,
|
||||
actionlog: DefaultActionlog,
|
||||
store: DefaultStore,
|
||||
ac: DefaultAccessControl,
|
||||
@@ -225,11 +228,17 @@ func (svc *session) Watch(ctx context.Context) {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case s := <-svc.spawnQueue:
|
||||
wfexecSessLog := zap.NewNop()
|
||||
if svc.opt.ExecDebug {
|
||||
wfexecSessLog = svc.log.Named("exec")
|
||||
}
|
||||
|
||||
//
|
||||
s.session <- wfexec.NewSession(ctx,
|
||||
s.graph,
|
||||
wfexec.SetHandler(svc.stateChangeHandler(ctx)),
|
||||
wfexec.SetLogger(svc.log),
|
||||
wfexec.SetLogger(wfexecSessLog),
|
||||
wfexec.SetDumpStacktraceOnPanic(svc.opt.ExecDebug),
|
||||
)
|
||||
// case time for a pool cleanup
|
||||
// @todo cleanup pool when sessions are complete
|
||||
|
||||
@@ -10,14 +10,16 @@ package options
|
||||
|
||||
type (
|
||||
WorkflowOpt struct {
|
||||
Register bool `env:"WORKFLOW_REGISTER"`
|
||||
Register bool `env:"WORKFLOW_REGISTER"`
|
||||
ExecDebug bool `env:"WORKFLOW_EXEC_DEBUG"`
|
||||
}
|
||||
)
|
||||
|
||||
// Workflow initializes and returns a WorkflowOpt with default values
|
||||
func Workflow() (o *WorkflowOpt) {
|
||||
o = &WorkflowOpt{
|
||||
Register: true,
|
||||
Register: true,
|
||||
ExecDebug: false,
|
||||
}
|
||||
|
||||
fill(o)
|
||||
|
||||
@@ -5,4 +5,9 @@ props:
|
||||
- name: register
|
||||
type: bool
|
||||
default: true
|
||||
description: Registers enabled and valid workflows and executes them whe ntriggere
|
||||
description: Registers enabled and valid workflows and executes them when triggered
|
||||
|
||||
- name: execDebug
|
||||
type: bool
|
||||
default: false
|
||||
description: Enables verbose logging for workflow execution
|
||||
|
||||
+21
-2
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/cortezaproject/corteza-server/pkg/id"
|
||||
"github.com/cortezaproject/corteza-server/pkg/logger"
|
||||
"go.uber.org/zap"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
@@ -56,6 +57,8 @@ type (
|
||||
// debug logger
|
||||
log *zap.Logger
|
||||
|
||||
dumpStacktraceOnPanic bool
|
||||
|
||||
eventHandler StateChangeHandler
|
||||
}
|
||||
|
||||
@@ -165,6 +168,7 @@ func NewSession(ctx context.Context, g *Graph, oo ...sessionOpt) *Session {
|
||||
mux: &sync.RWMutex{},
|
||||
|
||||
log: zap.NewNop(),
|
||||
|
||||
eventHandler: func(SessionStatus, *State, *Session) {
|
||||
// noop
|
||||
},
|
||||
@@ -490,13 +494,22 @@ func (s *Session) exec(ctx context.Context, st *State) (err error) {
|
||||
return
|
||||
}
|
||||
|
||||
var perr error
|
||||
|
||||
// normalize error and set it to state
|
||||
switch reason := reason.(type) {
|
||||
case error:
|
||||
s.qErr <- fmt.Errorf("step %d crashed: %w", st.step.ID(), reason)
|
||||
perr = fmt.Errorf("step %d crashed: %w", st.step.ID(), reason)
|
||||
default:
|
||||
s.qErr <- fmt.Errorf("step %d crashed: %v", st.step.ID(), reason)
|
||||
perr = fmt.Errorf("step %d crashed: %v", st.step.ID(), reason)
|
||||
}
|
||||
|
||||
if s.dumpStacktraceOnPanic {
|
||||
fmt.Printf("Error: %v\n", perr)
|
||||
println(string(debug.Stack()))
|
||||
}
|
||||
|
||||
s.qErr <- perr
|
||||
}()
|
||||
|
||||
var (
|
||||
@@ -738,6 +751,12 @@ func SetLogger(log *zap.Logger) sessionOpt {
|
||||
}
|
||||
}
|
||||
|
||||
func SetDumpStacktraceOnPanic(dump bool) sessionOpt {
|
||||
return func(s *Session) {
|
||||
s.dumpStacktraceOnPanic = dump
|
||||
}
|
||||
}
|
||||
|
||||
func (ss Steps) hash() map[Step]bool {
|
||||
out := make(map[Step]bool)
|
||||
for _, s := range ss {
|
||||
|
||||
@@ -10,7 +10,6 @@ package automation
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
atypes "github.com/cortezaproject/corteza-server/automation/types"
|
||||
"github.com/cortezaproject/corteza-server/pkg/expr"
|
||||
"github.com/cortezaproject/corteza-server/pkg/wfexec"
|
||||
|
||||
Reference in New Issue
Block a user