mirror of
https://github.com/certimate-go/certimate.git
synced 2026-09-24 23:10:13 +08:00
feat: enhance workflow logging
This commit is contained in:
@@ -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(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -214,7 +214,7 @@ const WorkflowRunLogs = ({ runId, runStatus }: { runId: string; runStatus: strin
|
||||
}
|
||||
|
||||
return (
|
||||
<div className="flex space-x-2 text-xs" style={{ wordBreak: "break-word" }}>
|
||||
<div className="flex space-x-2" style={{ wordBreak: "break-word" }}>
|
||||
{showTimestamp ? <div className="whitespace-nowrap text-stone-400">[{dayjs(record.timestamp).format("YYYY-MM-DD HH:mm:ss")}]</div> : <></>}
|
||||
<div
|
||||
className={mergeCls(
|
||||
@@ -323,7 +323,7 @@ const WorkflowRunLogs = ({ runId, runStatus }: { runId: string; runStatus: strin
|
||||
<span className="font-mono text-stone-400">{`#${group.id}\u00A0`}</span>
|
||||
<span>{group.name}</span>
|
||||
</div>
|
||||
<div className="flex flex-col space-y-1">{group.records.map((record) => renderRecord(record))}</div>
|
||||
<div className="flex flex-col space-y-1 text-xs">{group.records.map((record) => renderRecord(record))}</div>
|
||||
</div>
|
||||
);
|
||||
})}
|
||||
|
||||
@@ -74,6 +74,7 @@ export const get = async (id: string) => {
|
||||
.collection(COLLECTION_NAME_CERTIFICATE)
|
||||
.getOne<CertificateModel>(id, {
|
||||
expand: ["workflowRef"].join(","),
|
||||
fields: ["*", "expand.workflowRef.id", "expand.workflowRef.name", "expand.workflowRef.description"].join(","),
|
||||
requestKey: null,
|
||||
});
|
||||
};
|
||||
|
||||
@@ -63,6 +63,15 @@ export const get = async (id: string) => {
|
||||
.collection(COLLECTION_NAME_WORKFLOW)
|
||||
.getOne<WorkflowModel>(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,
|
||||
});
|
||||
};
|
||||
|
||||
@@ -48,6 +48,7 @@ export const get = async (id: string) => {
|
||||
.collection(COLLECTION_NAME_WORKFLOW_RUN)
|
||||
.getOne<WorkflowRunModel>(id, {
|
||||
expand: ["workflowRef"].join(","),
|
||||
fields: ["*", "expand.workflowRef.id", "expand.workflowRef.name", "expand.workflowRef.description"].join(","),
|
||||
requestKey: null,
|
||||
});
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user