diff --git a/automation/service/session.go b/automation/service/session.go index fab79d7c3..e74a58093 100644 --- a/automation/service/session.go +++ b/automation/service/session.go @@ -29,9 +29,10 @@ type ( } spawn struct { - session chan *wfexec.Session - graph *wfexec.Graph - trace bool + workflowID uint64 + session chan *wfexec.Session + graph *wfexec.Graph + trace bool } sessionAccessController interface { @@ -161,7 +162,7 @@ func (svc *session) Start(g *wfexec.Graph, i auth.Identifiable, ssp types.Sessio var ( ctx = auth.SetIdentityToContext(context.Background(), i) - ses = svc.spawn(g, ssp.Trace) + ses = svc.spawn(g, ssp.WorkflowID, ssp.Trace) ) ses.CreatedAt = *now() @@ -203,8 +204,13 @@ func (svc *session) Resume(sessionID, stateID uint64, i auth.Identifiable, input // // We need initial context for the session because we want to catch all cancellations or timeouts from there // and not from any potential HTTP requests or similar temporary context that can prematurely destroy a workflow session -func (svc *session) spawn(g *wfexec.Graph, trace bool) (ses *types.Session) { - s := &spawn{make(chan *wfexec.Session, 1), g, trace} +func (svc *session) spawn(g *wfexec.Graph, workflowID uint64, trace bool) (ses *types.Session) { + s := &spawn{ + workflowID: workflowID, + session: make(chan *wfexec.Session, 1), + graph: g, + trace: trace, + } // Send new-session request svc.spawnQueue <- s @@ -230,13 +236,16 @@ func (svc *session) Watch(ctx context.Context) { case s := <-svc.spawnQueue: wfexecSessLog := zap.NewNop() if svc.opt.ExecDebug { - wfexecSessLog = svc.log.Named("exec") + wfexecSessLog = svc.log. + Named("exec"). + With(zap.Uint64("workflowID", s.workflowID)) } // s.session <- wfexec.NewSession(ctx, s.graph, wfexec.SetHandler(svc.stateChangeHandler(ctx)), + wfexec.SetWorkflowID(s.workflowID), wfexec.SetLogger(wfexecSessLog), wfexec.SetDumpStacktraceOnPanic(svc.opt.ExecDebug), ) diff --git a/pkg/wfexec/session.go b/pkg/wfexec/session.go index 652ffb19b..1c0b179e7 100644 --- a/pkg/wfexec/session.go +++ b/pkg/wfexec/session.go @@ -18,6 +18,8 @@ type ( // Session identifier id uint64 + workflowID uint64 + // steps graph g *Graph @@ -417,7 +419,12 @@ func (s *Session) worker(ctx context.Context) { <-s.execLock if st.err != nil { - st.err = fmt.Errorf("step %d execution failed: %w", st.step.ID(), st.err) + st.err = fmt.Errorf( + "workflow %d step %d execution failed: %w", + s.workflowID, + st.step.ID(), + st.err, + ) } status := s.Status() @@ -745,6 +752,12 @@ func SetHandler(fn StateChangeHandler) sessionOpt { } } +func SetWorkflowID(workflowID uint64) sessionOpt { + return func(s *Session) { + s.workflowID = workflowID + } +} + func SetLogger(log *zap.Logger) sessionOpt { return func(s *Session) { s.log = log