99 lines
3.8 KiB
Go
99 lines
3.8 KiB
Go
package main
|
||
|
||
import (
|
||
digitalhumanController "ai-agent/digital-human/controller"
|
||
"ai-agent/workflow/service/flow"
|
||
|
||
// 空导入:加载内置模型工具与工作流处理器(组合根统一装载,业务侧不感知实现位置)
|
||
_ "ai-agent/tools/builtin/current_time"
|
||
"ai-agent/workflow/consts/public"
|
||
workController "ai-agent/workflow/controller"
|
||
workflowController "ai-agent/workflow/controller/flow"
|
||
workflowNodeController "ai-agent/workflow/controller/node"
|
||
sessionController "ai-agent/workflow/controller/session"
|
||
workflowSkillController "ai-agent/workflow/controller/skill"
|
||
toolController "ai-agent/workflow/controller/tool"
|
||
_ "ai-agent/workflow/service/flow/processor/builtin/media"
|
||
_ "ai-agent/workflow/service/flow/processor/builtin/split_batch"
|
||
_ "ai-agent/workflow/service/flow/processor/builtin/split_segment"
|
||
_ "ai-agent/workflow/service/flow/processor/builtin/split_shots_pipeline"
|
||
"context"
|
||
"os"
|
||
"os/signal"
|
||
"syscall"
|
||
"time"
|
||
|
||
"gitea.redpowerfuture.com/red-future/common/http"
|
||
"gitea.redpowerfuture.com/red-future/common/jaeger"
|
||
gmq "github.com/bjang03/gmq/core/gmq"
|
||
"github.com/bjang03/gmq/mq"
|
||
_ "github.com/gogf/gf/contrib/drivers/pgsql/v2"
|
||
_ "github.com/gogf/gf/contrib/nosql/redis/v2"
|
||
"github.com/gogf/gf/v2/frame/g"
|
||
)
|
||
|
||
func main() {
|
||
ctx := context.Background()
|
||
defer jaeger.ShutDown(ctx)
|
||
|
||
// 注册HTTP路由
|
||
http.Httpserver.BindHandler("/httpNodeCallback", workflowController.FlowCallBack.HttpNodeCallback)
|
||
|
||
http.RouteRegister([]interface{}{
|
||
//digitalhuman相关接口
|
||
digitalhumanController.Audio, // 语音相关接口
|
||
digitalhumanController.CustomVoice, // 自定义语音相关接口
|
||
digitalhumanController.DigitalHuman, // 数字人相关接口
|
||
digitalhumanController.Video, // 视频相关接口
|
||
digitalhumanController.AsyncTask, // 异步任务相关接口
|
||
workController.CreationInfo,
|
||
workflowController.FlowExecution,
|
||
workflowController.FlowUser,
|
||
workflowController.FlowTemplate,
|
||
workflowNodeController.NodeLibrary,
|
||
workflowNodeController.NodePrompt,
|
||
workflowSkillController.SkillTemplate,
|
||
workflowSkillController.SkillUser,
|
||
toolController.Tool,
|
||
sessionController.Session,
|
||
})
|
||
//workflow.ExternalInterruptDemo()
|
||
//err := activePullService.ActivePullService.AllList(ctx)
|
||
//if err != nil {
|
||
// g.Log().Error(ctx, "ActivePullService err: %v", err)
|
||
//}
|
||
|
||
gmq.GmqRegister(public.GmqMsgPluginsName, &mq.NatsConn{
|
||
NatsConfig: mq.NatsConfig{
|
||
Addr: g.Config().MustGet(ctx, "nats.addr").String(),
|
||
Port: g.Config().MustGet(ctx, "nats.port").String(),
|
||
Username: g.Config().MustGet(ctx, "nats.username").String(),
|
||
Password: g.Config().MustGet(ctx, "nats.password").String(),
|
||
},
|
||
})
|
||
|
||
// 启动工作流恢复扫描(首次立即扫 + 周期扫,多节点靠 Redis 锁去重)
|
||
flow.StartRecoveryLoop(ctx)
|
||
|
||
// 监听退出信号,执行优雅关闭
|
||
sigCh := make(chan os.Signal, 1)
|
||
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM, syscall.SIGQUIT)
|
||
shutdownDone := make(chan struct{})
|
||
go func() {
|
||
<-sigCh
|
||
// 先置关停标记再 Close:Close 会取消各连接 ctx,正在执行的 exec 会收到 context.Canceled;
|
||
// 标记让错误分类识别为"程序关停中断"(retryable=1),下次启动恢复扫描捞起续跑,
|
||
// 而非误判为"用户已终止执行"(retryable=0,永不恢复)
|
||
// SetShuttingDown 同时取消所有运行中执行(含脱离连接的恢复执行),保证全部落终态
|
||
flow.SetShuttingDown()
|
||
flow.SessionWsService.Close()
|
||
// 等待运行中执行落完终态再退出,避免进程退出时 exec 仍卡 status=1;
|
||
// 结束后返回 main 触发 deferred 资源清理,进程自然退出(不再 select{} 挂死)
|
||
flow.WaitExecRunsDrain(15 * time.Second)
|
||
close(shutdownDone)
|
||
}()
|
||
|
||
// 保持应用运行;收到关停信号并落完终态后自然退出
|
||
<-shutdownDone
|
||
}
|