diff --git a/internal/workflow/dispatcher/dispatcher.go b/internal/workflow/dispatcher/dispatcher.go index 8cf598dcc..e5c8de825 100644 --- a/internal/workflow/dispatcher/dispatcher.go +++ b/internal/workflow/dispatcher/dispatcher.go @@ -52,6 +52,8 @@ type workflowDispatcher struct { workflowRepo workflowRepository workflowRunRepo workflowRunRepository workflowLogRepo workflowLogRepository + + logger *slog.Logger } var _ WorkflowDispatcher = (*workflowDispatcher)(nil) @@ -99,12 +101,12 @@ func (wd *workflowDispatcher) Start(ctx context.Context, runId string) error { defer wd.taskMtx.Unlock() if _, exists := wd.processingTasks[runId]; exists { - return errors.New("workflow run is already processing") + return fmt.Errorf("workflow run %s is already processing", runId) } for _, pendingRunId := range wd.pendingRunQueue { if pendingRunId == runId { - return errors.New("workflow run is already in the queue") + return fmt.Errorf("workflow run %s is already in the queue", runId) } } @@ -122,7 +124,7 @@ func (wd *workflowDispatcher) Cancel(ctx context.Context, runId string) error { if err != nil { return err } else if workflowRun.Status != domain.WorkflowRunStatusTypePending && workflowRun.Status != domain.WorkflowRunStatusTypeProcessing { - return errors.New("workflow run is already completed") + return fmt.Errorf("workflow run #%s is already completed", workflowRun.Id) } workflow, err := wd.workflowRepo.GetById(ctx, workflowRun.WorkflowId) @@ -146,6 +148,8 @@ func (wd *workflowDispatcher) Cancel(ctx context.Context, runId string) error { if task, exists := wd.processingTasks[runId]; exists { task.cancel() delete(wd.processingTasks, runId) + + wd.logger.Info(fmt.Sprintf("workflow run #%s was canceled", task.RunId)) } for i, pendingRunId := range wd.pendingRunQueue { @@ -167,7 +171,7 @@ func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { defer func() { if r := recover(); r != nil { slog.Default().Warn(fmt.Sprintf("workflow dispatcher panic: %v, stack trace: %s", r, string(debug.Stack())), slog.Any("workflowId", task.WorkflowId), slog.Any("runId", task.RunId)) - app.GetLogger().Error(fmt.Sprintf("workflow dispatcher panic: %v", r), slog.Any("workflowId", task.WorkflowId), slog.Any("runId", task.RunId)) + wd.logger.Error(fmt.Sprintf("workflow dispatcher panic: %v", r), slog.Any("workflowId", task.WorkflowId), slog.Any("runId", task.RunId)) if workflowRun != nil { workflowRun.Status = domain.WorkflowRunStatusTypeFailed @@ -189,7 +193,7 @@ func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { // 查询运行实体,并级联更新状态 if run, err := wd.workflowRunRepo.GetById(task.ctx, task.RunId); err != nil { if !(errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded)) { - app.GetLogger().Error(fmt.Sprintf("failed to get workflow run #%s record", task.RunId), slog.Any("error", err)) + wd.logger.Error(fmt.Sprintf("failed to get workflow run #%s record", task.RunId), slog.Any("error", err)) } return } else { @@ -248,7 +252,7 @@ func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { logsBuf = append(logsBuf, log) if _, err := wd.workflowLogRepo.Save(ctx, &log); err != nil { - app.GetLogger().Error(err.Error()) + wd.logger.Error(err.Error()) } return nil @@ -267,15 +271,16 @@ func (wd *workflowDispatcher) tryExecuteAsync(task *taskInfo) { logsBuf = append(logsBuf, log) if _, err := wd.workflowLogRepo.Save(ctx, &log); err != nil { - app.GetLogger().Error(err.Error()) + wd.logger.Error(err.Error()) } return nil }) // 执行工作流 - app.GetLogger().Info(fmt.Sprintf("start to invoke workflow run #%s", task.RunId)) + wd.logger.Info(fmt.Sprintf("workflow run #%s was started", task.RunId)) we.Invoke(task.ctx, workflowRun.WorkflowId, workflowRun.Id, workflowRun.Graph) + wd.logger.Info(fmt.Sprintf("workflow run #%s was stopped", task.RunId)) } func (wd *workflowDispatcher) tryNextAsync() { @@ -284,7 +289,7 @@ func (wd *workflowDispatcher) tryNextAsync() { for i, pendingRunId := range wd.pendingRunQueue { workflowRun, err := wd.workflowRunRepo.GetById(context.Background(), pendingRunId) if err != nil { - app.GetLogger().Error(fmt.Sprintf("failed to get workflow run #%s record", pendingRunId), slog.Any("error", err)) + wd.logger.Error(fmt.Sprintf("failed to get workflow run #%s record", pendingRunId), slog.Any("error", err)) continue } @@ -296,7 +301,11 @@ func (wd *workflowDispatcher) tryNextAsync() { } } - if !hasSameWorkflowTask && len(wd.processingTasks) < wd.concurrency { + if hasSameWorkflowTask { + wd.logger.Warn(fmt.Sprintf("workflow run #%s is pending, because tasks that belonging to the same workflow already exists", pendingRunId)) + } else if len(wd.processingTasks) >= wd.concurrency { + wd.logger.Warn(fmt.Sprintf("workflow run #%s is pending, because the maximum concurrency limit has been reached", pendingRunId)) + } else { wd.taskMtx.RUnlock() wd.taskMtx.Lock() defer wd.taskMtx.Unlock() @@ -323,5 +332,7 @@ func newWorkflowDispatcher() WorkflowDispatcher { workflowRepo: repository.NewWorkflowRepository(), workflowRunRepo: repository.NewWorkflowRunRepository(), workflowLogRepo: repository.NewWorkflowLogRepository(), + + logger: app.GetLogger(), } } diff --git a/internal/workflow/engine/engine.go b/internal/workflow/engine/engine.go index c7a875ab3..3b98c6be4 100644 --- a/internal/workflow/engine/engine.go +++ b/internal/workflow/engine/engine.go @@ -28,8 +28,6 @@ type WorkflowEngine interface { } type workflowEngine struct { - logger *slog.Logger - executors map[NodeType]NodeExecutor hooksMtx sync.RWMutex @@ -42,6 +40,8 @@ type workflowEngine struct { onNodeLoggingHooks [](func(ctx context.Context, node *Node, log logging.Record) error) wfoutputRepo workflowOutputRepository + + logger *slog.Logger } var _ WorkflowEngine = (*workflowEngine)(nil) @@ -285,9 +285,9 @@ func (we *workflowEngine) fireOnNodeLoggingHooks(ctx context.Context, node *Node func NewWorkflowEngine() WorkflowEngine { engine := &workflowEngine{ - logger: app.GetLogger(), executors: make(map[NodeType]NodeExecutor), wfoutputRepo: repository.NewWorkflowOutputRepository(), + logger: app.GetLogger(), } engine.executors[NodeTypeStart] = newStartNodeExecutor() engine.executors[NodeTypeEnd] = newEndNodeExecutor() diff --git a/internal/workflow/event.go b/internal/workflow/event.go index 1af72a9ef..58b1660f3 100644 --- a/internal/workflow/event.go +++ b/internal/workflow/event.go @@ -3,6 +3,7 @@ package workflow import ( "context" "fmt" + "log/slog" "github.com/pocketbase/pocketbase/core" @@ -75,7 +76,8 @@ func onWorkflowRecordCreateOrUpdate(ctx context.Context, record *core.Record) er }) }) if err != nil { - return fmt.Errorf("add cron job failed: %w", err) + app.GetLogger().Error(fmt.Sprintf("failed to register cron job for workflow #%s", record.Id), slog.Any("error", err)) + return fmt.Errorf("failed to add cron job: %w", err) } return nil diff --git a/internal/workflow/service.go b/internal/workflow/service.go index 96cad7a05..2ce5506db 100644 --- a/internal/workflow/service.go +++ b/internal/workflow/service.go @@ -62,7 +62,7 @@ func (s *WorkflowService) InitSchedule(ctx context.Context) error { }) }) if err != nil { - app.GetLogger().Error(fmt.Sprintf("failed to add workflow #%s to scheduler: %w", workflow.Id), slog.Any("error", err)) + app.GetLogger().Error(fmt.Sprintf("failed to register cron job for workflow #%s", workflow.Id), slog.Any("error", err)) errs = append(errs, err) } diff --git a/ui/src/components/workflow/WorkflowRunDetail.tsx b/ui/src/components/workflow/WorkflowRunDetail.tsx index 36ee20e2b..e5a6780f5 100644 --- a/ui/src/components/workflow/WorkflowRunDetail.tsx +++ b/ui/src/components/workflow/WorkflowRunDetail.tsx @@ -214,7 +214,7 @@ const WorkflowRunLogs = ({ runId, runStatus }: { runId: string; runStatus: strin } return ( -
+
{showTimestamp ?
[{dayjs(record.timestamp).format("YYYY-MM-DD HH:mm:ss")}]
: <>}
{`#${group.id}\u00A0`} {group.name}
-
{group.records.map((record) => renderRecord(record))}
+
{group.records.map((record) => renderRecord(record))}
); })} diff --git a/ui/src/repository/certificate.ts b/ui/src/repository/certificate.ts index 4cca965d8..c92987394 100644 --- a/ui/src/repository/certificate.ts +++ b/ui/src/repository/certificate.ts @@ -74,6 +74,7 @@ export const get = async (id: string) => { .collection(COLLECTION_NAME_CERTIFICATE) .getOne(id, { expand: ["workflowRef"].join(","), + fields: ["*", "expand.workflowRef.id", "expand.workflowRef.name", "expand.workflowRef.description"].join(","), requestKey: null, }); }; diff --git a/ui/src/repository/workflow.ts b/ui/src/repository/workflow.ts index e9ee43d92..b492c64a8 100644 --- a/ui/src/repository/workflow.ts +++ b/ui/src/repository/workflow.ts @@ -63,6 +63,15 @@ export const get = async (id: string) => { .collection(COLLECTION_NAME_WORKFLOW) .getOne(id, { expand: ["lastRunRef"].join(","), + fields: [ + "*", + "expand.lastRunRef.id", + "expand.lastRunRef.status", + "expand.lastRunRef.trigger", + "expand.lastRunRef.startedAt", + "expand.lastRunRef.endedAt", + "expand.lastRunRef.error", + ].join(","), requestKey: null, }); }; diff --git a/ui/src/repository/workflowRun.ts b/ui/src/repository/workflowRun.ts index 72ef9bf16..76984e04a 100644 --- a/ui/src/repository/workflowRun.ts +++ b/ui/src/repository/workflowRun.ts @@ -48,6 +48,7 @@ export const get = async (id: string) => { .collection(COLLECTION_NAME_WORKFLOW_RUN) .getOne(id, { expand: ["workflowRef"].join(","), + fields: ["*", "expand.workflowRef.id", "expand.workflowRef.name", "expand.workflowRef.description"].join(","), requestKey: null, }); };