Files
ai-agent/workflow/service/flow/exec_lifecycle.go
19904408334 d699f7ce14 feat(workflow): 增加工作流计费与执行生命周期管理
- 新增计费模块:执行开始建单、终态结算/取消/失败处理,支持按条/按秒/按token计费
- 新增执行生命周期跟踪:优雅关停时取消运行中执行并等待落库
- 新增异步任务等待/通知机制(Wait/Notify)
- 重构执行记录落库与进度上报,统一失败分类与重试语义
- 重命名文件:async_task.go→async.go、flow_checkpoint_store.go→exec_checkpoint.go、flow_graph_util.go→exec_record.go
- 更新 .gitignore 与数据库密码配置
2026-09-03 13:22:22 +08:00

160 lines
6.1 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package flow
import (
"context"
"errors"
"sync"
"sync/atomic"
"time"
"ai-agent/workflow/consts/flow"
sessionDao "ai-agent/workflow/dao/session"
"ai-agent/workflow/model/entity"
"github.com/gogf/gf/v2/frame/g"
"github.com/google/uuid"
)
// 运行中执行跟踪:优雅关停时取消全部执行(含脱离连接的恢复执行)并等待落完终态再退出,
// 避免"程序停止 → exec 仍卡在 status=1"。
//
// 背景:WS 执行(handleExecute)的 execCtx 派生自连接 closeCtxClose() 取消连接即随之中止落终态;
// 但恢复执行(recoverExecution)的 ctx 经 context.WithoutCancel 脱离连接,程序关停时若不显式取消,
// 恢复中的 exec 不会落终态(仍 status=1),重启要等心跳陈旧 60s 才能再捞。这里登记所有运行中执行,
// SetShuttingDown 统一取消、main 等待全部落库后再退出。
var (
execRunMu sync.Mutex
execRuns = make(map[string]context.CancelFunc)
)
// trackExecRun 登记一次运行中执行,返回 finish 在落完终态(recordWorkflow 后)调用。
// 每次运行唯一 key(重试/断点续跑复用同一 execId 也不冲突)。
func trackExecRun(cancel context.CancelFunc) (finish func()) {
key := uuid.NewString()
execRunMu.Lock()
execRuns[key] = cancel
execRunMu.Unlock()
return func() {
execRunMu.Lock()
delete(execRuns, key)
execRunMu.Unlock()
}
}
// cancelAllExecRuns 优雅关停:取消所有运行中执行。恢复执行 ctx 经 WithoutCancel 脱离连接,
// 必须显式取消才能随 WS 执行一起落终态(status=3/retryable=1)。
func cancelAllExecRuns() {
execRunMu.Lock()
cancels := make([]context.CancelFunc, 0, len(execRuns))
for _, c := range execRuns {
cancels = append(cancels, c)
}
execRunMu.Unlock()
for _, c := range cancels {
c()
}
}
// WaitExecRunsDrain 等待所有运行中执行落完终态(限时),供 main 优雅关停收尾后退出进程。
// 关停后不再启动新执行(execute/reExecute/recoverExecution 顶部有 IsShuttingDown 守卫),
// 因此运行中集合只减不增,轮询安全。
func WaitExecRunsDrain(timeout time.Duration) {
deadline := time.Now().Add(timeout)
for {
execRunMu.Lock()
n := len(execRuns)
execRunMu.Unlock()
if n == 0 {
return
}
if time.Now().After(deadline) {
g.Log().Warningf(context.Background(), "优雅关停等待执行落库超时,剩余 %d 个执行", n)
return
}
time.Sleep(100 * time.Millisecond)
}
}
// errExecAlreadyRunning 用户触发时该执行已在运行(本节点/其它节点后台恢复),不新建执行
var errExecAlreadyRunning = errors.New("工作流正在执行中,不重复执行")
// errInterruptedByShutdown 程序优雅关停导致连接 ctx 取消时的错误标记(区别于用户主动取消)。
// 落库为 status=3 + retryable=1,下次启动恢复扫描捞起续跑。
var errInterruptedByShutdown = errors.New("程序关停中断")
// shuttingDown 优雅关停标记:程序收到退出信号后置位。
// 用于区分"程序关停导致的 WS 连接取消"与"用户主动取消",避免前者被误分类为不可重试。
var shuttingDown atomic.Bool
// SetShuttingDown 置位优雅关停标记并取消所有运行中执行(main 信号处理在 Close() 前调用)。
// 恢复执行 ctx 经 WithoutCancel 脱离连接,必须在此显式取消,否则程序关停时恢复中的 exec
// 不会落终态(仍 status=1);取消后 BuildExecution 随之中止、走恢复错误分支写 status=3/retryable=1。
func SetShuttingDown() {
shuttingDown.Store(true)
cancelAllExecRuns()
}
// IsShuttingDown 是否处于优雅关停
func IsShuttingDown() bool {
return shuttingDown.Load()
}
// shouldRetry 错误分类(retryable 终局语义与 retry_count 预算见《工作流执行并发仲裁设计.md》§5):
// 用户取消不重试;计费门禁拦截不重试;其余程序报错重试
func shouldRetry(err error) bool {
return err != nil && !errors.Is(err, context.Canceled) && !errors.Is(err, errExecAlreadyRunning) &&
!errors.Is(err, errBillingGateBlocked)
}
// isRecoverable 判定可恢复(恢复侧谓词,见《工作流执行并发仲裁设计.md》§2/§5):僵尸运行中(status=1 且心跳陈旧)或可重试失败(status=3 且 retryable=1 且未耗尽)。
// 与 ListRecoverable SQL 判定一致;last_heartbeat=0(老数据/默认)视为陈旧。
func isRecoverable(exec *entity.ExecWorkflow, nowMs int64) bool {
if exec == nil || exec.Status == nil {
return false
}
status := *exec.Status
switch status {
case *flow.FlowExecutionStatusRunning.Code():
return exec.LastHeartbeat < nowMs-int64(heartbeatStaleAfter/time.Millisecond)
case *flow.FlowExecutionStatusFailed.Code():
return exec.Retryable == 1 && exec.RetryCount < execMaxRetryCount
default:
return false
}
}
// startHeartbeat 后台心跳 goroutine:每 30s touch last_heartbeat,返回 stop 函数。
// 覆盖正常执行与恢复执行,崩溃前最后一次心跳即崩溃近似时间戳(心跳=在跑活体标记/租约语义见《工作流执行并发仲裁设计.md》§3)。
// onLeaseLost:心跳连续失败达到陈旧阈值(租约丢失)时回调——调用方应取消执行 ctx,
// 使心跳与执行同生共死,避免"心跳已过期但执行还活着"的窗口被其它节点恢复导致双跑。
func startHeartbeat(ctx context.Context, execId int64, onLeaseLost func()) func() {
stopCh := make(chan struct{})
done := make(chan struct{})
go func() {
defer close(done)
ticker := time.NewTicker(heartbeatInterval)
defer ticker.Stop()
var consecutiveFail int
for {
select {
case <-ticker.C:
if err := sessionDao.ExecWorkflowDao.TouchHeartbeat(ctx, execId); err != nil {
g.Log().Warningf(ctx, "心跳落库失败 execId=%d: %v", execId, err)
consecutiveFail++
if onLeaseLost != nil && consecutiveFail >= heartbeatMaxFail {
g.Log().Errorf(ctx, "心跳连续失败 %d 次,租约丢失,中止执行 execId=%d", consecutiveFail, execId)
onLeaseLost()
}
} else {
consecutiveFail = 0
}
case <-stopCh:
return
case <-ctx.Done():
return
}
}
}()
return func() { close(stopCh); <-done }
}