diff --git a/config.yml b/config.yml index 6703168..d81a003 100644 --- a/config.yml +++ b/config.yml @@ -26,10 +26,10 @@ database: timeMaintainDisabled: false # (可选)是否完全关闭时间更新特性,为true时CreatedAt/UpdatedAt/DeletedAt都将失效 black_deacon: - type: "pgsql" - host: "116.204.74.41" - port: "15432" + host: "localhost" + port: "5432" user: "postgres" - pass: "Bjang09@686^*^" + pass: "123456" name: "black-deacon" prefix: "black_deacon_" # (可选)表名前缀 role: "master" @@ -48,7 +48,7 @@ database: redis: default: - address: 116.204.74.41:6379 + address: localhost:6379 db: 0 idleTimeout: "60s" #连接最大空闲时间,使用时间字符串例如30s/1m/1d maxConnLifetime: "90s" #连接最长存活时间,使用时间字符串例如30s/1m/1d @@ -59,13 +59,13 @@ redis: maxActive: 100 consul: - address: 116.204.74.41:8500 + address: localhost:8500 jaeger: - addr: 116.204.74.41:4318 + addr: localhost:4318 -# 文件上传服务地址,与oss模块minio中的endpoint一致 -filePrefix: "http://116.204.74.41:9000" +# 文件上传服务地址,cdn访问地址 +filePrefix: "http://cdn.redpowerfuture.com" -model-asynch: - addr: "127.0.0.1:8001" +# 文件上传服务地址,minio内网访问地址 +minioPrefix: "http://192.168.0.83:9000" diff --git a/go.mod b/go.mod index c832e10..79be24f 100644 --- a/go.mod +++ b/go.mod @@ -3,13 +3,15 @@ module ai-agent go 1.26.0 require ( - gitea.redpowerfuture.com/red-future/common v0.0.23 + gitea.redpowerfuture.com/red-future/common v0.0.24 github.com/cloudwego/eino v0.9.5 + github.com/cloudwego/eino-examples v0.0.0-20260611092511-bd64846fbc1d github.com/cloudwego/eino-ext/components/model/qwen v0.1.9 github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.2 github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.2 github.com/gogf/gf/v2 v2.10.2 github.com/google/uuid v1.6.0 + github.com/stretchr/testify v1.11.1 github.com/tidwall/gjson v1.19.0 github.com/tidwall/sjson v1.2.5 go.opentelemetry.io/otel/trace v1.44.0 @@ -21,7 +23,7 @@ require ( github.com/bahlo/generic-list-go v0.2.0 // indirect github.com/buger/jsonparser v1.1.1 // indirect github.com/bwmarrin/snowflake v0.3.0 // indirect - github.com/bytedance/gopkg v0.1.3 // indirect + github.com/bytedance/gopkg v0.1.4 // indirect github.com/bytedance/sonic v1.15.0 // indirect github.com/bytedance/sonic/loader v0.5.0 // indirect github.com/cenkalti/backoff/v5 v5.0.3 // indirect @@ -66,7 +68,7 @@ require ( github.com/hashicorp/serf v0.10.1 // indirect github.com/json-iterator/go v1.1.12 // indirect github.com/klauspost/compress v1.18.2 // indirect - github.com/klauspost/cpuid/v2 v2.2.11 // indirect + github.com/klauspost/cpuid/v2 v2.3.0 // indirect github.com/lib/pq v1.10.9 // indirect github.com/magiconair/properties v1.8.10 // indirect github.com/mailru/easyjson v0.9.0 // indirect @@ -82,16 +84,16 @@ require ( github.com/olekukonko/errors v1.1.0 // indirect github.com/olekukonko/ll v0.0.9 // indirect github.com/olekukonko/tablewriter v1.1.0 // indirect - github.com/pelletier/go-toml/v2 v2.0.9 // indirect - github.com/pkg/errors v0.9.1 // indirect + github.com/pelletier/go-toml/v2 v2.2.4 // indirect + github.com/pkg/errors v0.9.2-0.20201214064552-5dd12d0cfe7f // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/r3labs/diff/v2 v2.15.1 // indirect - github.com/redis/go-redis/v9 v9.12.1 // indirect + github.com/redis/go-redis/v9 v9.17.2 // indirect github.com/rivo/uniseg v0.4.7 // indirect github.com/sirupsen/logrus v1.9.3 // indirect github.com/slongfield/pyfmt v0.0.0-20220222012616-ea85ff4c361f // indirect github.com/tidwall/match v1.1.1 // indirect - github.com/tidwall/pretty v1.2.0 // indirect + github.com/tidwall/pretty v1.2.1 // indirect github.com/tiger1103/gfast-token v1.0.10 // indirect github.com/twitchyliquid64/golang-asm v0.15.1 // indirect github.com/vcaesar/cedar v0.30.0 // indirect @@ -107,8 +109,8 @@ require ( go.opentelemetry.io/otel/metric v1.44.0 // indirect go.opentelemetry.io/otel/sdk v1.38.0 // indirect go.opentelemetry.io/proto/otlp v1.7.1 // indirect - golang.org/x/arch v0.11.0 // indirect - golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 // indirect + golang.org/x/arch v0.19.0 // indirect + golang.org/x/exp v0.0.0-20250718183923-645b1fa84792 // indirect golang.org/x/net v0.48.0 // indirect golang.org/x/sys v0.39.0 // indirect golang.org/x/text v0.32.0 // indirect diff --git a/go.sum b/go.sum index 0508072..0c5ec37 100644 --- a/go.sum +++ b/go.sum @@ -1,8 +1,8 @@ cloud.google.com/go v0.26.0/go.mod h1:aQUYkXzVsufM+DwF1aE+0xfcU+56JwCaLick0ClmMTw= -gitea.com/red-future/common v0.0.21 h1:8w30HmCVmFG/hphH3ODJs1KxDEGmRpq+/PXI0pQjJKc= -gitea.com/red-future/common v0.0.21/go.mod h1:6/nqIucVzmjOyqDTIq71feYBXXFNBy0rFwzaQ0/Ueoo= gitea.redpowerfuture.com/red-future/common v0.0.23 h1:xieoA00iKOCDm5SO9iXn+cSyMKBAlZwI0fuEVPWrHLg= gitea.redpowerfuture.com/red-future/common v0.0.23/go.mod h1:50U1Xi+Ie56z09S5LQbZvaken0Mxv3OeS9LgR7U/ZRY= +gitea.redpowerfuture.com/red-future/common v0.0.24 h1:sXxhnmDmCgn+KwH/3gDnhAtAQ7FCmf/5AsMfvxRmri0= +gitea.redpowerfuture.com/red-future/common v0.0.24/go.mod h1:50U1Xi+Ie56z09S5LQbZvaken0Mxv3OeS9LgR7U/ZRY= github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= github.com/BurntSushi/toml v1.5.0 h1:W5quZX/G/csjUnuI8SUYlsHs9M38FC7znL0lIO+DvMg= github.com/BurntSushi/toml v1.5.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= @@ -38,6 +38,8 @@ github.com/bwmarrin/snowflake v0.3.0 h1:xm67bEhkKh6ij1790JB83OujPR5CzNe8QuQqAgIS github.com/bwmarrin/snowflake v0.3.0/go.mod h1:NdZxfVWX+oR6y2K0o6qAYv6gIOP9rjG0/E9WsDpxqwE= github.com/bytedance/gopkg v0.1.3 h1:TPBSwH8RsouGCBcMBktLt1AymVo2TVsBVCY4b6TnZ/M= github.com/bytedance/gopkg v0.1.3/go.mod h1:576VvJ+eJgyCzdjS+c4+77QF3p7ubbtiKARP3TxducM= +github.com/bytedance/gopkg v0.1.4 h1:oZnQwnX82KAIWb7033bEwtxvTqXcYMxDBaQxo5JJHWM= +github.com/bytedance/gopkg v0.1.4/go.mod h1:v1zWfPm21Fb+OsyXN2VAHdL6TBb2L88anLQgdyje6R4= github.com/bytedance/mockey v1.3.0 h1:ONLRdvhqmCfr9rTasUB8ZKCfvbdD2tohOg4u+4Q/ed0= github.com/bytedance/mockey v1.3.0/go.mod h1:1BPHF9sol5R1ud/+0VEHGQq/+i2lN+GTsr3O2Q9IENY= github.com/bytedance/sonic v1.15.0 h1:/PXeWFaR5ElNcVE84U0dOHjiMHQOwNIx3K4ymzh/uSE= @@ -58,14 +60,16 @@ github.com/clbanning/mxj/v2 v2.7.0/go.mod h1:hNiWqW14h+kc+MdF9C6/YoRfjEJoR3ou6tn github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= github.com/cloudwego/base64x v0.1.6 h1:t11wG9AECkCDk5fMSoxmufanudBtJ+/HemLstXDLI2M= github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gEHfghB2IPU= -github.com/cloudwego/eino v0.8.13 h1:z5dhaZNN8TWZbP/lgKxGmF26Ii8fPeUlQCGV/NTtms0= -github.com/cloudwego/eino v0.8.13/go.mod h1:+2N4nsMPxA6kGBHpH+75JuTfEcGprAMTdsZESrShKpU= github.com/cloudwego/eino v0.9.5 h1:0Nftjx9gPek/2S/hzm38LVxSjk5/6mqRr3I9VKrKvm4= github.com/cloudwego/eino v0.9.5/go.mod h1:OBD1mrkfkt/pJa4rkg1P0VnaMeOVl7l8IAdEqY//3IQ= +github.com/cloudwego/eino-examples v0.0.0-20260611092511-bd64846fbc1d h1:NrAxhU58S5SgK5YbrYnaOQrQFwAz3x4/0qg46JM8Eo4= +github.com/cloudwego/eino-examples v0.0.0-20260611092511-bd64846fbc1d/go.mod h1:VVmcWGhnLIxLkrQAaCoOtQoEbcvOCrMQvbXgbo9O34Q= github.com/cloudwego/eino-ext/components/model/qwen v0.1.9 h1:xCz/mp43JeWqupjPR3zLRArmwC6P29/6lTwbwh1yzYM= github.com/cloudwego/eino-ext/components/model/qwen v0.1.9/go.mod h1:slTGTuhzkzhNavf+1UtUg1FvUSA31iNAF+rq1mT4SnI= github.com/cloudwego/eino-ext/libs/acl/openai v0.1.17 h1:EeVcR1TslRA2IdNW1h/2LaGbPlffwGhQm99jM3zWZiI= github.com/cloudwego/eino-ext/libs/acl/openai v0.1.17/go.mod h1:Zkcx6DPTR2NfWmtSXbhItswGw6hqUezNPhNcke0pOG8= +github.com/cloudwego/hertz v0.10.5 h1:N4oBqAJShSjYQm2Jfr0ryTzzJ9fnY9qSvIvUBMkoFWg= +github.com/cloudwego/hertz v0.10.5/go.mod h1:Im9u6rUa1v2mL2HiDKKJoof/CPQ3mPBBpT92v67Cetg= github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -116,20 +120,14 @@ github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY= github.com/gofrs/uuid v3.2.0+incompatible/go.mod h1:b2aQJv3Z4Fp6yNu3cdSllBxTCLRxnplIgP/c0N/04lM= -github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.0 h1:39+jbTenm7KBj4hO2C8ANAxVHpX/7OuRDs1VcGC9ylA= -github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.0/go.mod h1:B0s0fVzn0W220E8UTpSGzrrGKsop5KcB90twBeLCiz0= github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.2 h1:u8EpP24GkprogROnJ7htMov9Fc66pTP1eVYrWxiCYOs= github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.2/go.mod h1:GmvM3r8GVByVMi4RD2+MCs5+CfxVXPMeT8mVDkAaAXE= -github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.0 h1:N/F9CuDdUZLoM1nVRqrDE/33pDZuhVxpNY4wYdeIaBs= -github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.0/go.mod h1:x6uoJGfZOtirIRQls8xUlYzC6f7T/eULPUa9er368X0= github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.2 h1:iTQegT+lEg/wDKvj2mi3W1wrdrwFarjokf88EXVVgu4= github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.2/go.mod h1:ZRw3GNz5cq4uYrW4TPSVyrYWaoqzujKdWro/AOcGBaE= github.com/gogf/gf/contrib/registry/consul/v2 v2.9.5 h1:eUqwJ/qNH8lJ6yssiqskazgp1ACQuNU6zXlLOZVuXTQ= github.com/gogf/gf/contrib/registry/consul/v2 v2.9.5/go.mod h1:sjQyMry9+0POYZCA6lHXBxO77WoNKkruJpRB4xKqk5k= github.com/gogf/gf/contrib/trace/otlphttp/v2 v2.9.5 h1:tHUEZYB5GTqEYYVDYnlGobf1xISARKDE4KHVlgjwTec= github.com/gogf/gf/contrib/trace/otlphttp/v2 v2.9.5/go.mod h1:cfzTn2HS9RDX8f5pUVkbGxUWcSosouqfNQ1G6cY0V88= -github.com/gogf/gf/v2 v2.10.0 h1:rzDROlyqGMe/eM6dCalSR8dZOuMIdLhmxKSH1DGhbFs= -github.com/gogf/gf/v2 v2.10.0/go.mod h1:Svl1N+E8G/QshU2DUbh/3J/AJauqCgUnxHurXWR4Qx0= github.com/gogf/gf/v2 v2.10.2 h1:46IO0Uc8e85/FqdftJFskfDejJLBL0JBnGS5qOftUu8= github.com/gogf/gf/v2 v2.10.2/go.mod h1:Svl1N+E8G/QshU2DUbh/3J/AJauqCgUnxHurXWR4Qx0= github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ= @@ -243,8 +241,8 @@ github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.18.2 h1:iiPHWW0YrcFgpBYhsA6D1+fqHssJscY/Tm/y2Uqnapk= github.com/klauspost/compress v1.18.2/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4= -github.com/klauspost/cpuid/v2 v2.2.11 h1:0OwqZRYI2rFrjS4kvkDnqJkKHdHaRnCm68/DY4OxRzU= -github.com/klauspost/cpuid/v2 v2.2.11/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= +github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= @@ -314,12 +312,13 @@ github.com/onsi/gomega v1.5.0/go.mod h1:ex+gbHU/CVuBBDIJjb2X0qEXbFg53c61hWP/1Cpa github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= github.com/pascaldekloe/goe v0.1.0 h1:cBOtyMzM9HTpWjXfbbunk26uA6nG3a8n06Wieeh0MwY= github.com/pascaldekloe/goe v0.1.0/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= -github.com/pelletier/go-toml/v2 v2.0.9 h1:uH2qQXheeefCCkuBBSLi7jCiSmj3VRh2+Goq2N7Xxu0= -github.com/pelletier/go-toml/v2 v2.0.9/go.mod h1:tJU2Z3ZkXwnxa4DPO899bsyIoywizdUvyaeZurnPPDc= +github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= +github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= +github.com/pkg/errors v0.9.2-0.20201214064552-5dd12d0cfe7f h1:lJqhwddJVYAkyp72a4pwzMClI20xTwL7miDdm2W/KBM= +github.com/pkg/errors v0.9.2-0.20201214064552-5dd12d0cfe7f/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= @@ -339,13 +338,13 @@ github.com/prometheus/procfs v0.0.2/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsT github.com/prometheus/procfs v0.0.8/go.mod h1:7Qr8sr6344vo1JqZ6HhLceV9o3AJ1Ff+GxbHq6oeK9A= github.com/r3labs/diff/v2 v2.15.1 h1:EOrVqPUzi+njlumoqJwiS/TgGgmZo83619FNDB9xQUg= github.com/r3labs/diff/v2 v2.15.1/go.mod h1:I8noH9Fc2fjSaMxqF3G2lhDdC0b+JXCfyx85tWFM9kc= -github.com/redis/go-redis/v9 v9.12.1 h1:k5iquqv27aBtnTm2tIkROUDp8JBXhXZIVu1InSgvovg= -github.com/redis/go-redis/v9 v9.12.1/go.mod h1:huWgSWd8mW6+m0VPhJjSSQ+d6Nh1VICQ6Q5lHuCH/Iw= +github.com/redis/go-redis/v9 v9.17.2 h1:P2EGsA4qVIM3Pp+aPocCJ7DguDHhqrXNhVcEp4ViluI= +github.com/redis/go-redis/v9 v9.17.2/go.mod h1:u410H11HMLoB+TP67dz8rL9s6QW2j76l0//kSOd3370= github.com/rivo/uniseg v0.2.0/go.mod h1:J6wj4VEh+S6ZtnVlnTBMWIodfgj8LQOQFoIToxlJtxc= github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ= github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88= -github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII= -github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/rollbar/rollbar-go v1.0.2/go.mod h1:AcFs5f0I+c71bpHlXNNDbOWJiKwjFDtISeXco0L5PKQ= github.com/ryanuber/columnize v0.0.0-20160712163229-9b3edd62028f/go.mod h1:sm1tb6uqfes/u+d4ooFouqFdy9/2g9QGwK3SQygK0Ts= github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529 h1:nn5Wsu0esKSJiIVhscUtVbo7ada43DJhG55ua/hjS5I= @@ -384,8 +383,9 @@ github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU= github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc= github.com/tidwall/match v1.1.1 h1:+Ho715JplO36QYgwN9PGYNhgZvoUSc9X2c80KVTi+GA= github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= -github.com/tidwall/pretty v1.2.0 h1:RWIZEg2iJ8/g6fDDYzMpobmaoGh5OLl4AXtGUGPcqCs= github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= +github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= +github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= github.com/tidwall/sjson v1.2.5 h1:kLy8mja+1c9jlljvWTlSazM7cKDRfJuR/bOJhcY5NcY= github.com/tidwall/sjson v1.2.5/go.mod h1:Fvgq9kS/6ociJEDnK0Fk1cpYF4FIW6ZF7LAe+6jwd28= github.com/tiger1103/gfast-token v1.0.10 h1:fNiBE/Dq5iTHvTGlCx3DmXa2o4hr0NtumFpffZ39k6s= @@ -411,28 +411,20 @@ go.mongodb.org/mongo-driver/v2 v2.5.0 h1:yXUhImUjjAInNcpTcAlPHiT7bIXhshCTL3jVBkF go.mongodb.org/mongo-driver/v2 v2.5.0/go.mod h1:yOI9kBsufol30iFsl1slpdq1I0eHPzybRWdyYUs8K/0= go.opencensus.io v0.23.0 h1:gqCw0LfLxScz8irSi8exQc7fyQ0fKQU/qnC/X8+V/1M= go.opencensus.io v0.23.0/go.mod h1:XItmlyltB5F7CS4xOC1DcqMoFqwtC6OG2xF7mCv7P7E= -go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= -go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/otel v1.38.0 h1:RkfdswUDRimDg0m2Az18RKOsnI8UDzppJAtj01/Ymk8= -go.opentelemetry.io/otel v1.38.0/go.mod h1:zcmtmQ1+YmQM9wrNsTGV/q/uyusom3P8RxwExxkZhjM= go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0 h1:GqRJVj7UmLjCVyVJ3ZFLdPRmhDUp2zFmQe3RHIOsw24= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.38.0/go.mod h1:ri3aaHSmCTVYu2AWv44YMauwAQc0aqI9gHKIcSbI1pU= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0 h1:aTL7F04bJHUlztTsNGJ2l+6he8c+y/b//eR0jjjemT4= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0/go.mod h1:kldtb7jDTeol0l3ewcmd8SDvx3EmIE7lyvqbasU3QC4= -go.opentelemetry.io/otel/metric v1.38.0 h1:Kl6lzIYGAh5M159u9NgiRkmoMKjvbsKtYRwgfrA6WpA= -go.opentelemetry.io/otel/metric v1.38.0/go.mod h1:kB5n/QoRM8YwmUahxvI3bO34eVtQf2i4utNVLr9gEmI= go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= go.opentelemetry.io/otel/sdk v1.38.0 h1:l48sr5YbNf2hpCUj/FoGhW9yDkl+Ma+LrVl8qaM5b+E= go.opentelemetry.io/otel/sdk v1.38.0/go.mod h1:ghmNdGlVemJI3+ZB5iDEuk4bWA3GkTpW+DOoZMYBVVg= go.opentelemetry.io/otel/sdk/metric v1.38.0 h1:aSH66iL0aZqo//xXzQLYozmWrXxyFkBJ6qT5wthqPoM= go.opentelemetry.io/otel/sdk/metric v1.38.0/go.mod h1:dg9PBnW9XdQ1Hd6ZnRz689CbtrUp0wMMs9iPcgT9EZA= -go.opentelemetry.io/otel/trace v1.38.0 h1:Fxk5bKrDZJUH+AMyyIXGcFAPah0oRcT+LuNtJrmcNLE= -go.opentelemetry.io/otel/trace v1.38.0/go.mod h1:j1P9ivuFsTceSWe1oY+EeW3sc+Pp42sO++GHkg4wwhs= go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= go.opentelemetry.io/proto/otlp v1.7.1 h1:gTOMpGDb0WTBOP8JaO72iL3auEZhVmAQg4ipjOVAtj4= @@ -441,8 +433,8 @@ go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/mock v0.5.0 h1:KAMbZvZPyBPWgD14IrIQ38QCyjwpvVVV6K/bHl1IwQU= go.uber.org/mock v0.5.0/go.mod h1:ge71pBPLYDk7QIi1LupWxdAykm7KIEFchiOqd6z7qMM= -golang.org/x/arch v0.11.0 h1:KXV8WWKCXm6tRpLirl2szsO5j/oOODwZf4hATmGVNs4= -golang.org/x/arch v0.11.0/go.mod h1:FEVrYAQjsQXMVJ1nsMoVVXPZg6p2JE2mx8psSWTDQys= +golang.org/x/arch v0.19.0 h1:LmbDQUodHThXE+htjrnmVD73M//D9GTH6wFZjyDkjyU= +golang.org/x/arch v0.19.0/go.mod h1:bdwinDaKcfZUGpH09BB7ZmOfhalA8lQdzl62l8gGWsk= golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20190923035154-9ee001bba392/go.mod h1:/lpIB1dKB+9EgE3H3cr1v9wB50oz8l4C4h62xy7jSTY= @@ -451,8 +443,8 @@ golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPh golang.org/x/crypto v0.46.0 h1:cKRW/pmt1pKAfetfu+RCEvjvZkA9RimPbh7bhFjGVBU= golang.org/x/crypto v0.46.0/go.mod h1:Evb/oLKmMraqjZ2iQTwDwvCtJkczlDuTmdJXoZVzqU0= golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= -golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 h1:nDVHiLt8aIbd/VzvPWN6kSOPE7+F/fNFDSXLVYkE/Iw= -golang.org/x/exp v0.0.0-20250305212735-054e65f0b394/go.mod h1:sIifuuw/Yco/y6yb6+bDNfyeQ/MdPUy/hKEMYQV17cM= +golang.org/x/exp v0.0.0-20250718183923-645b1fa84792 h1:R9PFI6EUdfVKgwKjZef7QIwGcBKu86OEFpJ9nUEP2l4= +golang.org/x/exp v0.0.0-20250718183923-645b1fa84792/go.mod h1:A+z0yzpGtvnG90cToK5n2tu8UJVP2XUATh+r+sfOOOc= golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE= golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU= golang.org/x/lint v0.0.0-20190313153728-d0100b6bd8b3/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= diff --git a/update.sql b/update.sql index dcd7ba4..2f30a81 100644 --- a/update.sql +++ b/update.sql @@ -415,7 +415,7 @@ CREATE TABLE IF NOT EXISTS black_deacon_flow_execution ( node_input_params JSONB DEFAULT '[]'::JSONB, output_params JSONB DEFAULT '[]'::JSONB, error_message TEXT DEFAULT '', -- 错误信息 - trace_id VARCHAR(64) DEFAULT '' -- 跟踪ID + trace_id VARCHAR(64) DEFAULT '', -- 跟踪ID session_id VARCHAR(64) DEFAULT '' -- 会话ID ); diff --git a/workflow/consts/node/node_template.go b/workflow/consts/node/node_template.go index 7db8827..30fe692 100644 --- a/workflow/consts/node/node_template.go +++ b/workflow/consts/node/node_template.go @@ -16,15 +16,13 @@ const ( NodeNameAudioModel = "生成音频" NodeNameBatchModel = "批量处理一起返回" NodeNameDataConversionModel = "参数转换" - NodeNameSenseOptimizeModel = "语义优化" - NodeNameStoryOptimizeModel = "分镜优化" - NodeNameScriptOptimizeModel = "剧本优化" NodeNameModel = "模型" NodeNameMerge = "结果合并" NodeNameDataMerge = "结果汇集" NodeNameJudge = "条件判断" NodeNameLoop = "循环" NodeNameForm = "表单" + NodeSubFlow = "子流程" NodeNameHttp = "HTTP(S)接口" NodeNameCustomNode = "自定义节点" NodeNameSystemSum = "系统-结果汇总" @@ -50,16 +48,12 @@ type NodeType string const ( // 组件 - NodeTypeTextModel NodeType = "text_model" - NodeTypeImageModel NodeType = "image_model" - NodeTypeVideoModel NodeType = "video_model" - NodeTypeAudioModel NodeType = "audio_model" - NodeTypeBatchModel NodeType = "batch_model" - + NodeTypeTextModel NodeType = "text_model" + NodeTypeImageModel NodeType = "image_model" + NodeTypeVideoModel NodeType = "video_model" + NodeTypeAudioModel NodeType = "audio_model" + NodeTypeBatchModel NodeType = "batch_model" NodeTypeDataConversionModel NodeType = "data_conversion_model" - NodeTypeSenseOptimizeModel NodeType = "sense_optimize_model" - NodeTypeStoryOptimizeModel NodeType = "story_optimize_model" - NodeTypeScriptOptimizeModel NodeType = "script_optimize_model" // 基础 NodeTypeModel NodeType = "model" NodeTypeMerge NodeType = "merge" @@ -67,6 +61,7 @@ const ( NodeTypeJudge NodeType = "judge" NodeTypeForm NodeType = "form" NodeTypeIntent NodeType = "intent" + NodeTypeSubFlow NodeType = "sub_flow" NodeTypeHttp NodeType = "http" // 自定义 NodeTypeCustomNode NodeType = "custom_node" diff --git a/workflow/model/dto/flow/flow_execution_dto.go b/workflow/model/dto/flow/flow_execution_dto.go index 6b8f628..2ffe70c 100644 --- a/workflow/model/dto/flow/flow_execution_dto.go +++ b/workflow/model/dto/flow/flow_execution_dto.go @@ -128,12 +128,11 @@ type ComposeCallbackReq struct { type ModelCallbackReq struct { g.Meta `path:"/modelCallback" method:"post" tags:"提示词处理" summary:"model-gateway 回调" dc:"model-gateway 成功后 GET 回调:callbackUrl/{bizName}"` - TaskId string `p:"task_id" json:"task_id" v:"required#task_id不能为空" dc:"网关任务ID"` - State int `p:"state" json:"state" dc:"网关任务状态"` - OssFile string `p:"oss_file" json:"oss_file" dc:"结果文件地址"` - FileType string `p:"file_type" json:"file_type" dc:"结果文件类型"` - Messages map[string]any `json:"messages"` - ErrorMsg string `json:"error_msg"` + TaskId string `p:"task_id" json:"task_id" v:"required#task_id不能为空" dc:"网关任务ID"` + State int `p:"state" json:"state" dc:"网关任务状态"` + OssFile string `p:"oss_file" json:"oss_file" dc:"结果文件地址"` + FileType string `p:"file_type" json:"file_type" dc:"结果文件类型"` + ErrorMsg string `json:"error_msg"` } type VideoCallbackReq struct { diff --git a/workflow/model/entity/flow_user.go b/workflow/model/entity/flow_user.go index 4cb344f..bb92c2f 100644 --- a/workflow/model/entity/flow_user.go +++ b/workflow/model/entity/flow_user.go @@ -23,6 +23,7 @@ type FlowNode struct { PromptContent string `json:"promptContent"` IsSaveFile bool `json:"isSaveFile"` InputSource []FlowNodeInputSource `json:"inputSource"` // 前端指定:来源节点ID + SubConfig *SubFlowConfig `json:"subConfig"` FormConfig []node.NodeFormField `json:"formConfig"` ModelConfig node.ModelItem `json:"modelConfig"` OutputConfig []node.NodeFormField `json:"outputConfig"` @@ -30,9 +31,23 @@ type FlowNode struct { } type FlowNodeInputSource struct { - NodeId string `json:"nodeId"` - QuoteOutput bool `json:"quoteOutput"` - Field []string `json:"field"` + NodeId string `json:"nodeId"` + QuoteOutput bool `json:"quoteOutput"` + Field []string `json:"field"` + FieldMap []FlowField `json:"fieldMap"` +} + +type FlowField struct { + Key string `json:"key"` + Value string `json:"value"` + Desc string `json:"desc"` +} + +// SubFlowConfig 子流程节点配置 +type SubFlowConfig struct { + FlowId int64 `json:"flowId"` + MaxConcurrency int `json:"maxConcurrency"` // 子流程并发数 + InputSource []FlowNodeInputSource `json:"inputSource"` // 前端指定:来源节点ID } type FlowEdge struct { diff --git a/workflow/service/flow/flow_execution_service.go b/workflow/service/flow/flow_execution_service.go index f85d309..ea5a74a 100644 --- a/workflow/service/flow/flow_execution_service.go +++ b/workflow/service/flow/flow_execution_service.go @@ -285,6 +285,13 @@ func (s *flowExecutionService) Execute(ctx context.Context, req *flowDto.Execute cancel() }() + //getRes, err := FlowUserService.Get(ctx, &flowDto.GetFlowUserReq{ + // Id: req.FlowId, + //}) + //if err != nil { + // return nil, err + //} + nodeInputParams := ExtractFlowNodeFrom(req.FlowContent) flowInfo, err := flowDao.FlowExecutionDao.Get(ctx, &flowDto.GetFlowExecutionReq{ SessionId: req.SessionId, }) @@ -302,7 +309,7 @@ func (s *flowExecutionService) Execute(ctx context.Context, req *flowDto.Execute r.NodeGroupId = nodeGroupId r.TriggerType = flow.FlowExecutionTriggerTypeManual.Code() r.FlowContent = req.FlowContent - r.NodeInputParams = req.NodeInputParams + //r.NodeInputParams = nodeInputParams r.SessionId = req.SessionId r.Status = flow.FlowExecutionStatusRunning.Code() span := trace.SpanFromContext(ctx) @@ -385,7 +392,7 @@ func (s *flowExecutionService) Execute(ctx context.Context, req *flowDto.Execute // ✅【第3步】构建 ConfigMap // ========================================================================= configMap := make(map[string]*entity.FlowNode) - for _, cfg := range req.NodeInputParams { + for _, cfg := range nodeInputParams { configMap[cfg.Id] = cfg } for _, i := range nodeList { @@ -430,8 +437,7 @@ func (s *flowExecutionService) Execute(ctx context.Context, req *flowDto.Execute return } -// BuildGraphFromFlowContent 根据前端保存的工作流JSON,自动构建执行图 -func BuildGraphFromFlowContent(ctx context.Context, flowContent *entity.FlowInfo) ([]entity.FlowNode, compose.Runnable[any, any], error) { +func BuildGraph(ctx context.Context, flowContent *entity.FlowInfo) ([]entity.FlowNode, *compose.Graph[any, any]) { // 注册自定义合并函数:处理 *flowDto.FlowExecutionInput 类型合并 // 由于 ConfigMap 是 map 引用类型,所有并行分支修改已经写入共享内存 // 直接返回第一个实例即可,所有修改都已经可见 @@ -550,15 +556,18 @@ func BuildGraphFromFlowContent(ctx context.Context, flowContent *entity.FlowInfo _ = graph.AddEdge(e.From, e.To) } } + return nodeList, graph +} + +// BuildGraphFromFlowContent 根据前端保存的工作流JSON,自动构建执行图 +func BuildGraphFromFlowContent(ctx context.Context, flowContent *entity.FlowInfo) ([]entity.FlowNode, compose.Runnable[any, any], error) { + nodeList, graph := BuildGraph(ctx, flowContent) compile, err := graph.Compile(ctx, compose.WithGraphName("auto_build_workflow")) return nodeList, compile, err } // -------------------------- 节点自动注册器(核心分发) -------------------------- func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNode) { - nodeID := flowNode.Id - code := flowNode.NodeCode - // 通用包装:全程入参都是 *FlowExecutionInput wrapLambda := func(lambda func(ctx context.Context, input any) (any, error)) func(ctx context.Context, input any) (any, error) { return func(ctx context.Context, input any) (any, error) { @@ -569,9 +578,9 @@ func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNod } configMap := execInput.ConfigMap - currentConfig := configMap[nodeID] + currentConfig := configMap[flowNode.Id] if currentConfig == nil { - return nil, fmt.Errorf("节点%s无配置", nodeID) + return nil, fmt.Errorf("节点%s无配置", flowNode.Id) } // 获取入参 - 适配切片类型:遍历所有来源节点 @@ -604,7 +613,7 @@ func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNod nodeExecutionId, err := nodeDao.NodeExecutionDao.Insert(ctx, &nodeDto.CreateNodeExecutionReq{ FlowExecutionId: execInput.ExecutionId, - NodeId: nodeID, + NodeId: flowNode.Id, NodeName: flowNode.Name, NodeGroupId: execInput.NodeGroupId, InputParamsPath: ossResult.FileURL, @@ -613,7 +622,7 @@ func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNod if err != nil { // 记录失败到已执行列表 execInput.ExecutedNodes = append(execInput.ExecutedNodes, flowDto.ExecutedNode{ - NodeId: nodeID, + NodeId: flowNode.Id, Status: node.NodeExecutionStatusFailed.Code(), }) return nil, err @@ -622,18 +631,9 @@ func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNod // 执行节点 _, err = lambda(ctx, realInput) durationMs := time.Since(startTime).Milliseconds() - // 上传OSS(每条独立上传) - ossResult1, err := Upload(ctx, &dto.UploadFileBytesReq{ - FileBytes: gconv.Bytes(gconv.String(realInput)), - FileName: fmt.Sprintf("nodeInput:%v.txt", time.Now().UnixMilli()), - }) - if err != nil { - return nil, err - } updateReq := &nodeDto.UpdateNodeExecutionReq{ - Id: nodeExecutionId, - OutputParamsPath: ossResult1.FileURL, - DurationMs: durationMs, + Id: nodeExecutionId, + DurationMs: durationMs, } if err != nil { // 执行失败,更新状态 @@ -642,18 +642,26 @@ func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNod _, _ = nodeDao.NodeExecutionDao.Update(ctx, updateReq) // 记录失败到已执行列表 execInput.ExecutedNodes = append(execInput.ExecutedNodes, flowDto.ExecutedNode{ - NodeId: nodeID, + NodeId: flowNode.Id, Status: node.NodeExecutionStatusFailed.Code(), }) return nil, err } - + // 上传OSS(每条独立上传) + ossResult1, err := Upload(ctx, &dto.UploadFileBytesReq{ + FileBytes: gconv.Bytes(gconv.String(realInput)), + FileName: fmt.Sprintf("nodeInput:%v.txt", time.Now().UnixMilli()), + }) + if err != nil { + return nil, err + } + updateReq.OutputParamsPath = ossResult1.FileURL // 执行成功,更新状态 updateReq.Status = node.NodeExecutionStatusSuccess.Code() _, _ = nodeDao.NodeExecutionDao.Update(ctx, updateReq) // 记录成功到已执行列表 execInput.ExecutedNodes = append(execInput.ExecutedNodes, flowDto.ExecutedNode{ - NodeId: nodeID, + NodeId: flowNode.Id, Status: node.NodeExecutionStatusSuccess.Code(), }) @@ -661,35 +669,37 @@ func registerNodeToGraph(graph *compose.Graph[any, any], flowNode entity.FlowNod return execInput, nil } } - switch code { + switch flowNode.NodeCode { case "__start__": - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(StartLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(StartLambda))) case node.NodeTypeSystemSum: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(SummaryLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(SummaryLambda))) case node.NodeTypeTextModel: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(TextModelLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(TextModelLambda))) case node.NodeTypeImageModel: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(ImageModelLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(ImageModelLambda))) case node.NodeTypeVideoModel: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(VideoModelLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(VideoModelLambda))) case node.NodeTypeAudioModel: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(AudioModelLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(AudioModelLambda))) case node.NodeTypeBatchModel: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(BatchModelLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(BatchModelLambda))) case node.NodeTypeDataConversionModel: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(DataConversionLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(DataConversionLambda))) case node.NodeTypeCustomNode: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(CustomLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(CustomLambda))) case node.NodeTypeForm: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(FormLambda))) - case node.NodeTypeIntent: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(IntentLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(FormLambda))) + //case node.NodeTypeIntent: + // _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(IntentLambda))) case node.NodeTypeMerge: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(MergeLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(MergeLambda))) case node.NodeTypeDataMerge: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(DataMergeLambda)), compose.WithGraphCompileOptions(compose.WithNodeTriggerMode(compose.AllPredecessor))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(DataMergeLambda)), compose.WithGraphCompileOptions(compose.WithNodeTriggerMode(compose.AllPredecessor))) + case node.NodeTypeSubFlow: + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(SubFlowLambda))) case node.NodeTypeHttp: - _ = graph.AddLambdaNode(nodeID, compose.InvokableLambda(wrapLambda(HttpLambda))) + _ = graph.AddLambdaNode(flowNode.Id, compose.InvokableLambda(wrapLambda(HttpLambda))) } } diff --git a/workflow/service/flow/lambda_node.go b/workflow/service/flow/lambda_node.go index 55d4ddc..ce9f21d 100644 --- a/workflow/service/flow/lambda_node.go +++ b/workflow/service/flow/lambda_node.go @@ -9,6 +9,7 @@ import ( "ai-agent/workflow/model/dto" fileDto "ai-agent/workflow/model/dto/file" flowDto "ai-agent/workflow/model/dto/flow" + "ai-agent/workflow/model/entity" "context" "fmt" "strconv" @@ -18,6 +19,8 @@ import ( "gitea.redpowerfuture.com/red-future/common/db/gfdb" "gitea.redpowerfuture.com/red-future/common/utils" + "github.com/cloudwego/eino-examples/compose/batch/batch" + "github.com/cloudwego/eino/compose" "github.com/gogf/gf/v2/database/gdb" "github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/util/gconv" @@ -31,8 +34,78 @@ func FormLambda(ctx context.Context, input any) (any, error) { return input, nil } -func IntentLambda(ctx context.Context, input any) (any, error) { - return input, nil +func SubFlowLambda(ctx context.Context, input any) (any, error) { + // 1. 类型断言(和其他节点保持一致的入参结构) + nodeExecInput, ok := input.(*flowDto.NodeExecutionInput) + if !ok { + return nil, fmt.Errorf("子流程节点入参类型错误,期望*flowDto.NodeExecutionInput,实际%T", input) + } + // 2. 解析子流程配置 + subFlowConfig := nodeExecInput.Config.SubConfig + if subFlowConfig == nil { + return nil, fmt.Errorf("子流程节点缺少配置") + } + getRes, err := FlowUserService.Get(ctx, &flowDto.GetFlowUserReq{ + Id: subFlowConfig.FlowId, + }) + if err != nil { + return nil, err + } + // 3. 编译子流程Graph(复用现有 BuildGraphFromFlowContent 逻辑) + nodeList, subGraph := BuildGraph(ctx, getRes.FlowContent) + // 4. 构建子流程Workflow(绑定START/END,和示例对齐) + innerWorkflow := compose.NewWorkflow[*flowDto.FlowExecutionInput, *flowDto.FlowExecutionInput]() + // 挂载子图节点并绑定全局START + innerWorkflow.AddGraphNode("sub_flow_graph", subGraph).AddInput(compose.START) + // 绑定子图输出到全局END + innerWorkflow.End().AddInput("sub_flow_graph") + // 5. 构建BatchNode(批量执行子流程,复用示例逻辑) + batchNode := batch.NewBatchNode(&batch.NodeConfig[*flowDto.FlowExecutionInput, *flowDto.FlowExecutionInput]{ + Name: fmt.Sprintf("sub_flow_batch_%s", nodeExecInput.Config.Id), + InnerTask: innerWorkflow, + MaxConcurrency: subFlowConfig.MaxConcurrency, + }) + + skillName, from, userFrom := BuildParam(nodeExecInput) + fmt.Printf("skillName: %s, from: %s, userFrom: %s\n", skillName, from, userFrom) + + // 6. 提取批量输入(从全局入参中获取) + batchInputs := make([]*flowDto.FlowExecutionInput, 0) + + nodeInputParams := ExtractFlowNodeFrom(getRes.FlowContent) + configMap := make(map[string]*entity.FlowNode) + for _, cfg := range nodeInputParams { + configMap[cfg.Id] = cfg + } + for _, i := range nodeList { + configMap[i.Id] = &i + } + // ========================================================================= + // ✅【第4步】构建全局执行入参(现在 schemaMap 是有值的!) + // ========================================================================= + execInput := &flowDto.FlowExecutionInput{ + NodeGroupId: nodeExecInput.Global.NodeGroupId, + IsDialogue: nodeExecInput.Global.IsDialogue, + ExecutionId: nodeExecInput.Global.ExecutionId, + ConfigMap: configMap, + SessionId: nodeExecInput.Global.SessionId, + Desc: nodeExecInput.Global.Desc, + SkillName: nodeExecInput.Global.SkillName, + FileUrl: nodeExecInput.Global.FileUrl, + } + batchInputs = append(batchInputs, execInput) + + // 7. 执行批量子流程 + batchOutput, err := batchNode.Invoke(ctx, batchInputs) + if err != nil { + return nil, fmt.Errorf("执行子流程BatchNode失败: %v", err) + } + for idx, singleSubResult := range batchOutput { + fmt.Printf("【批量任务%d 最终消息条数】: %v\n", idx+1, singleSubResult) + } + // 8. 保存子流程执行结果到当前节点输出 + //nodeExecInput.Config.OutputResult = append(nodeExecInput.Config.OutputResult, batchOutput) + return nodeExecInput, nil } // JudgeLambda 分支判断核心:读取IntentLambda的输出 → 返回目标节点ID做路由 @@ -139,62 +212,75 @@ func BatchModelLambda(ctx context.Context, input any) (any, error) { } } } - // 结果按索引存放,保证顺序 + // 结果按索引存放,切片不同下标并发写无竞争,不用锁 res := make([][]node.NodeFormField, len(reqMap)) var wg sync.WaitGroup - // 用一个通道标记是否完成 - done := make(chan struct{}) - // 错误只存一个 - var execErr error - // 并发执行 + subCtx, cancel := context.WithCancel(ctx) + defer cancel() + + // 缓冲1错误通道,仅接收第一个错误 + errCh := make(chan error, 1) + + // 并发执行任务 for idx, item := range reqMap { wg.Add(1) go func(idx int, userItem map[string]any) { defer wg.Done() + // 上下文已取消则直接退出 + select { + case <-subCtx.Done(): + return + default: + } + singleUserFrom := []map[string]any{userItem} - output, err := TextNode(ctx, nodeInput, skillName, from, singleUserFrom) + output, err := TextNode(subCtx, nodeInput, skillName, from, singleUserFrom) if err != nil { - // 并发安全赋值错误 - if execErr == nil { - execErr = err + // 仅第一个错误写入通道 + select { + case errCh <- err: + cancel() // 触发全局取消,其他协程快速退出 + default: } return } - - // 直接按原索引写,顺序绝对正确 res[idx] = output }(idx, item) } - // 后台等待所有协程完成,然后关闭 done 通道 + // 任务全部结束后关闭错误通道 go func() { wg.Wait() - close(done) + close(errCh) }() - // 等待全部完成 - <-done - - // 如果有错误,直接返回 - if execErr != nil { - return nil, execErr + // ========== 修正后的等待逻辑 ========== + var execErr error + select { + // 优先捕获业务错误 + case execErr = <-errCh: + if execErr != nil { + // 收到真实业务错误,等待剩余协程收尾后返回 + wg.Wait() + return nil, execErr + } + // execErr == nil 代表通道关闭、无任何错误,走到下方返回完整结果 + case <-subCtx.Done(): + // 上下文被取消,阻塞读完errCh,确认是否存在业务错误 + execErr = <-errCh } - // 全局自增 i + // 拼接输出结果 var globalIndex int var outputRes []node.NodeFormField for _, items := range res { for _, item := range items { - // 1. 拿到原来的 Field:例如 "text_content:2:0" oldField := item.Field - // 2. 找到最后一个 : 的位置 if idx := strings.LastIndex(oldField, ":"); idx != -1 { - // 3. 截断前面部分,拼接上新的 globalIndex item.Field = oldField[:idx+1] + fmt.Sprint(globalIndex) } - // Label 同理 oldLabel := item.Label if idx := strings.LastIndex(oldLabel, ":"); idx != -1 { item.Label = oldLabel[:idx+1] + fmt.Sprint(globalIndex) @@ -290,6 +376,7 @@ func VideoModelLambda(ctx context.Context, input any) (any, error) { if err != nil { return nil, err } + newS := strings.ReplaceAll(urlPrefix, g.Cfg().MustGet(ctx, "filePrefix").String(), g.Cfg().MustGet(ctx, "minioPrefix").String()) outputRes := make([]node.NodeFormField, 0) if nodeInput.Config.IsSaveFile { @@ -302,7 +389,7 @@ func VideoModelLambda(ctx context.Context, input any) (any, error) { } outputRes = append(outputRes, node.NodeFormField{ Field: fmt.Sprintf("concat_video_url:content:%d", 0), - Value: urlPrefix + msg.FileURL, + Value: newS + msg.FileURL, Label: fmt.Sprintf("视频内容:content:%d", 0), Type: "string", }) diff --git a/workflow/service/flow/lambda_node_util.go b/workflow/service/flow/lambda_node_util.go index 63020a6..923dc08 100644 --- a/workflow/service/flow/lambda_node_util.go +++ b/workflow/service/flow/lambda_node_util.go @@ -258,16 +258,15 @@ func waitGatewayResult(ctx context.Context, taskId string) (map[string]any, erro if task.State == 3 || !g.IsEmpty(task.ErrorMsg) { return nil, fmt.Errorf("模型执行失败:%s", task.ErrorMsg) } - if g.IsEmpty(task.Messages) { + if g.IsEmpty(task.OssFile) { return nil, fmt.Errorf("模型返回结果为空") } // 获取远程文件内容 - //file, err := GetFileBytesFromURL(ctx, task.OssFile) - //if err != nil { - // return nil, err - //} - //task.Messages = gconv.Map(file) - return task.Messages, nil + file, err := GetFileBytesFromURL(ctx, task.OssFile) + if err != nil { + return nil, err + } + return gconv.Map(file), nil } // updateTokenCount updates the token count in node execution @@ -296,6 +295,9 @@ func GetModelResult(ctx context.Context, sessionId string, nodeInput *flowDto.No if err != nil { return nil, err } + if composeResult.Status != "success" { + return nil, fmt.Errorf("模型提示词构建错误") + } modelInfo, err := GetModelInfo(ctx, &flowDto.GetModelInfoReq{ModelName: nodeInput.Config.ModelConfig.ModelName}) if err != nil { @@ -345,31 +347,41 @@ func GetModelResult(ctx context.Context, sessionId string, nodeInput *flowDto.No taskIdList := make([]string, len(composeResult.Messages.Rounds)) for idx, item := range composeResult.Messages.Rounds { - var taskId string - taskId, err = createGatewayTaskOnly(ctx, composeResult.EpicycleId, nodeInput.Config.ModelConfig.ModelName, item) + taskId, err := createGatewayTaskOnly(ctx, composeResult.EpicycleId, nodeInput.Config.ModelConfig.ModelName, item) if err != nil { return nil, err } taskIdList[idx] = taskId } + // 全局共享子上下文,实现一处报错全部终止 + subCtx, globalCancel := context.WithCancel(ctx) + defer globalCancel() // 函数退出兜底释放 + var wg sync.WaitGroup errChan := make(chan error, len(taskIdList)) + // 加互斥锁保护结果map + var mu sync.Mutex + for idx, taskId := range taskIdList { wg.Add(1) go func(idx int, taskId string) { defer wg.Done() - var taskResult map[string]any - taskResult, err = waitGatewayResult(ctx, taskId) + taskResult, err := waitGatewayResult(subCtx, taskId) if err != nil { errChan <- err + globalCancel() // 全局取消,所有协程收到ctx取消信号快速退出 return } + // 加锁写入map,解决并发竞态 + mu.Lock() mapTaskResult[idx] = taskResult + mu.Unlock() + updateTokenCount(ctx, nodeInput.NodeExecutionId, modelInfo.Model.ResponseTokenField, taskResult) }(idx, taskId) } @@ -377,8 +389,15 @@ func GetModelResult(ctx context.Context, sessionId string, nodeInput *flowDto.No wg.Wait() close(errChan) - if len(errChan) > 0 { - return nil, <-errChan + // 收集全部错误,而非只读一条 + var errs []error + for len(errChan) > 0 { + errs = append(errs, <-errChan) + } + + if len(errs) > 0 { + // 返回第一个错误;如需汇总所有错误可拼接 + return nil, errs[0] } } } else { @@ -490,16 +509,17 @@ func VideoConcat(ctx context.Context, videoUrls []string) (r any, err error) { } func GetFileBytesFromURL(ctx context.Context, fileUrl string) ([]byte, error) { + newS := strings.ReplaceAll(fileUrl, g.Cfg().MustGet(ctx, "filePrefix").String(), g.Cfg().MustGet(ctx, "minioPrefix").String()) // 使用 GoFrame 客户端(自带超时、追踪、日志等能力) - resp, err := g.Client().Get(ctx, fileUrl) + resp, err := g.Client().Get(ctx, newS) if err != nil { - return nil, gerror.Wrapf(err, "failed to request url: %s", fileUrl) + return nil, gerror.Wrapf(err, "failed to request url: %s", newS) } defer resp.Close() // 校验状态码 if resp.StatusCode != http.StatusOK { - return nil, gerror.Newf("request failed with status code: %d, url: %s", resp.StatusCode, fileUrl) + return nil, gerror.Newf("request failed with status code: %d, url: %s", resp.StatusCode, newS) } // 读取全部内容 diff --git a/workflow/service/node/node_library_service.go b/workflow/service/node/node_library_service.go index 6781b29..c7c553d 100644 --- a/workflow/service/node/node_library_service.go +++ b/workflow/service/node/node_library_service.go @@ -93,36 +93,23 @@ func (s *nodeLibraryService) GetNodeLibrary(ctx context.Context, req *nodeDto.Wo FormConfig: []node.NodeFormField{}, ModelConfig: []node.ModelItem{}, }, - //{ - // NodeCode: node.NodeTypeSenseOptimizeModel, - // NodeName: node.NodeNameSenseOptimizeModel, - // ModelType: node.ModelTypeText, - // SkillOption: false, - // FormConfig: []node.NodeFormField{}, - // ModelConfig: []node.ModelItem{}, - //}, - //{ - // NodeCode: node.NodeTypeStoryOptimizeModel, - // NodeName: node.NodeNameStoryOptimizeModel, - // ModelType: node.ModelTypeText, - // SkillOption: false, - // FormConfig: []node.NodeFormField{}, - // ModelConfig: []node.ModelItem{}, - //}, - //{ - // NodeCode: node.NodeTypeScriptOptimizeModel, - // NodeName: node.NodeNameScriptOptimizeModel, - // ModelType: node.ModelTypeText, - // SkillOption: false, - // FormConfig: []node.NodeFormField{}, - // ModelConfig: []node.ModelItem{}, - //}, }, }, { Group: node.NodeGroupBase, Label: node.NodeGroupNameBase, Items: []node.NodeItem{ + { + NodeCode: node.NodeTypeSubFlow, + NodeName: node.NodeSubFlow, + SkillOption: false, + PromptOption: false, + IsSaveFile: false, + FormConfig: []node.NodeFormField{ + {Field: "maxConcurrency", Label: "最大并发数", Type: "input", Required: true}, + }, + ModelConfig: []node.ModelItem{}, + }, { NodeCode: node.NodeTypeDataConversionModel, NodeName: node.NodeNameDataConversionModel,