Skip to content

Commit 93c0607

Browse files
committed
feat(worker): 实现基于进程的 Worker 调度与管理
- 新增 WorkerScheduler 接口替代原 WorkerRuntime - 新增 WorkerEnvType 类型定义 (process/docker/kubevirt) - 实现 ProcessScheduler 支持 Worker 进程启动/停止/健康检查 - Worker 启动时通过 WebSocket 从 Server 获取 DigitalAssistant 配置 - Server 新增 getConfig 消息处理,从数据库查询配置返回给 Worker - DigitalAssistant Service 创建时自动启动对应的 Worker 进程 - 支持 --assistant-code 参数指定数字员工标识 架构流程: 用户创建 DigitalAssistant → 保存数据库 → 启动 Worker 进程 → Worker 连接 Server 获取配置 → 初始化 Agent Runtime → 就绪
1 parent 5c22462 commit 93c0607

9 files changed

Lines changed: 659 additions & 289 deletions

File tree

.gitignore

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,3 +40,4 @@ example/
4040

4141
# Prevent accidental Go source files in root directory (except go.mod/go.sum)
4242
# Note: Actual Go test files should be in backend/tests or appropriate package directories
43+
docs/superpowers

backend/cmd/singer/worker.go

Lines changed: 28 additions & 233 deletions
Original file line numberDiff line numberDiff line change
@@ -3,274 +3,69 @@ package main
33
import (
44
"context"
55
"fmt"
6-
"net"
7-
"net/http"
8-
"strings"
9-
"time"
106

11-
"github.com/gin-gonic/gin"
127
"github.com/insmtx/SingerOS/backend/config"
13-
"github.com/insmtx/SingerOS/backend/internal/agent"
14-
"github.com/insmtx/SingerOS/backend/internal/agent/externalcli"
15-
"github.com/insmtx/SingerOS/backend/internal/eventengine"
16-
"github.com/insmtx/SingerOS/backend/internal/infra/mq"
17-
singerMCP "github.com/insmtx/SingerOS/backend/mcp"
18-
"github.com/insmtx/SingerOS/backend/runtime/engines"
19-
"github.com/insmtx/SingerOS/backend/runtime/engines/builtin"
20-
"github.com/insmtx/SingerOS/backend/tools"
21-
skilltools "github.com/insmtx/SingerOS/backend/tools/skill"
8+
"github.com/insmtx/SingerOS/backend/internal/worker/client"
229
"github.com/spf13/cobra"
23-
"github.com/ygpkg/yg-go/lifecycle"
10+
ygconfig "github.com/ygpkg/yg-go/config"
2411
"github.com/ygpkg/yg-go/logs"
2512
)
2613

2714
var (
28-
workerConfigPath string
29-
workerServerAddr string
15+
workerConfigPath string
16+
workerServerAddr string
17+
workerAssistantCode string
3018
)
3119

3220
var workerCmd = &cobra.Command{
3321
Use: "worker",
3422
Short: "Start the SingerOS background worker",
3523
Long: `Start the background worker service for processing asynchronous tasks and events.`,
3624
Run: func(cmd *cobra.Command, args []string) {
37-
mcpServer, err := startWorkerMCPServer(workerServerAddr)
38-
if err != nil {
39-
logs.Fatalf("Failed to start worker MCP server: %v", err)
40-
return
41-
}
25+
ctx := cmd.Context()
4226

43-
cfg, err := loadWorkerConfig(workerConfigPath, workerServerAddr)
27+
worker, err := createWorker(ctx)
4428
if err != nil {
45-
logs.Fatalf("Failed to load config: %v", err)
29+
logs.Fatalf("Failed to create worker: %v", err)
4630
return
4731
}
4832

49-
natsUrl := "nats://nats:4222"
50-
if cfg.NATS != nil && cfg.NATS.URL != "" {
51-
natsUrl = cfg.NATS.URL
52-
}
53-
54-
subscriber, err := mq.NewPublisher(natsUrl)
55-
if err != nil {
56-
logs.Fatalf("Failed to create event subscriber: %v", err)
57-
return
33+
if err := worker.Start(ctx); err != nil {
34+
logs.Fatalf("Failed to start worker: %v", err)
5835
}
59-
60-
runtimeConfig, err := buildRuntimeConfig()
61-
if err != nil {
62-
logs.Fatalf("Failed to build runtime config: %v", err)
63-
return
64-
}
65-
66-
ctx, cancel := context.WithCancel(context.Background())
67-
runner, err := buildRuntimeRunner(ctx, cfg, runtimeConfig)
68-
if err != nil {
69-
cancel()
70-
logs.Fatalf("Failed to create agent runtime: %v", err)
71-
return
72-
}
73-
74-
orchestratorInstance := eventengine.NewOrchestrator(subscriber, runner)
75-
if err := orchestratorInstance.Start(ctx); err != nil {
76-
cancel()
77-
logs.Fatalf("Failed to start orchestrator: %v", err)
78-
return
79-
}
80-
logs.Info("Orchestrator started successfully")
81-
82-
lifecycle.Std().AddCloseFunc(func() error {
83-
cancel()
84-
return nil
85-
})
86-
lifecycle.Std().AddCloseFunc(func() error {
87-
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 5*time.Second)
88-
defer shutdownCancel()
89-
return mcpServer.Shutdown(shutdownCtx)
90-
})
91-
lifecycle.Std().AddCloseFunc(subscriber.Close)
92-
93-
logs.Info("Worker runtime initialized successfully")
94-
logs.Info("Worker service started")
95-
96-
lifecycle.Std().WaitExit()
97-
98-
logs.Info("Worker exited")
9936
},
10037
}
10138

10239
func init() {
10340
workerCmd.Flags().StringVar(&workerConfigPath, "config", "", "Configuration file path")
104-
workerCmd.Flags().StringVar(&workerServerAddr, "server-addr", ":8081", "Worker MCP server listen address for runtime CLI bootstrap")
41+
workerCmd.Flags().StringVar(&workerServerAddr, "server-addr", ":8080", "Server address for WebSocket connection")
42+
workerCmd.Flags().StringVar(&workerAssistantCode, "assistant-code", "", "Assistant code for configuration retrieval")
10543
rootCmd.AddCommand(workerCmd)
10644
}
10745

108-
func startWorkerMCPServer(addr string) (*http.Server, error) {
109-
if strings.TrimSpace(addr) == "" {
110-
addr = ":8081"
111-
}
112-
113-
listener, err := net.Listen("tcp", addr)
46+
func createWorker(ctx context.Context) (*client.Worker, error) {
47+
_, err := loadWorkerConfig()
11448
if err != nil {
115-
return nil, fmt.Errorf("listen on %s: %w", addr, err)
116-
}
117-
118-
r := gin.New()
119-
v1 := r.Group("/v1")
120-
singerMCP.RegisterRoutes(v1, singerMCP.NewServer())
121-
122-
server := &http.Server{
123-
Addr: addr,
124-
Handler: r,
125-
}
126-
127-
go func() {
128-
logs.Infof("Worker MCP server listening on %s", listener.Addr().String())
129-
if err := server.Serve(listener); err != nil && err != http.ErrServerClosed {
130-
logs.Errorf("Worker MCP server stopped unexpectedly: %v", err)
131-
}
132-
}()
133-
134-
return server, nil
135-
}
136-
137-
func loadWorkerConfig(configPath string, bootstrapAddr string) (*config.Config, error) {
138-
cfg, err := loadConfig(configPath)
139-
if err != nil {
140-
return nil, err
141-
}
142-
143-
bootstrapped, err := builtin.BootstrapCLIEngines(context.Background(), cfg.CLI, defaultCLIBootstrapOptions(bootstrapAddr))
144-
if err != nil {
145-
logs.Warnf("CLI bootstrap failed: %v", err)
146-
}
147-
if bootstrapped != nil {
148-
cfg.CLI = bootstrapped
49+
return nil, fmt.Errorf("load config: %w", err)
14950
}
15051

151-
return cfg, nil
52+
return client.NewWorker(ctx, &client.WorkerConfig{
53+
ServerAddr: workerServerAddr,
54+
AssistantCode: workerAssistantCode,
55+
SkillsDir: "",
56+
ToolsEnabled: true,
57+
})
15258
}
15359

154-
func defaultCLIBootstrapOptions(addr string) builtin.BootstrapOptions {
155-
return builtin.BootstrapOptions{
156-
MCP: engines.MCPServerConfig{
157-
Name: "singeros",
158-
URL: mcpURLFromAddr(addr),
159-
BearerToken: singerMCP.DefaultAuthToken(),
160-
},
161-
}
162-
}
163-
164-
func mcpURLFromAddr(addr string) string {
165-
host := "localhost"
166-
port := "8081"
167-
168-
if strings.TrimSpace(addr) != "" {
169-
if splitHost, splitPort, err := net.SplitHostPort(addr); err == nil {
170-
if splitHost != "" && splitHost != "0.0.0.0" && splitHost != "::" && splitHost != "[::]" {
171-
host = splitHost
172-
}
173-
if splitPort != "" {
174-
port = splitPort
175-
}
176-
} else if strings.HasPrefix(addr, ":") {
177-
port = strings.TrimPrefix(addr, ":")
178-
} else {
179-
host = addr
180-
}
181-
}
182-
183-
return fmt.Sprintf("http://%s:%s/v1/mcp", host, port)
184-
}
185-
186-
func buildRuntimeConfig() (agent.Config, error) {
187-
catalog, skillDir, err := skilltools.LoadDefaultCatalog()
188-
if err != nil {
189-
return agent.Config{}, fmt.Errorf("load skills: %w", err)
190-
}
191-
192-
logs.Infof("Loaded %d skills from %s for runtime", len(catalog.List()), skillDir)
193-
194-
toolRegistry, err := buildTooling(catalog)
195-
if err != nil {
196-
return agent.Config{}, err
197-
}
198-
199-
return agent.Config{
200-
SkillsCatalog: catalog,
201-
ToolRegistry: toolRegistry,
202-
}, nil
203-
}
204-
205-
func buildTooling(catalog *skilltools.Catalog) (*tools.Registry, error) {
206-
registry := tools.NewRegistry()
207-
208-
if err := skilltools.Register(registry, catalog); err != nil {
209-
return nil, fmt.Errorf("register skill use tool: %w", err)
210-
}
211-
212-
logs.Infof("Loaded %d tools for runtime", len(registry.List()))
213-
214-
return registry, nil
215-
}
216-
217-
func buildRuntimeRunner(ctx context.Context, cfg *config.Config, runtimeConfig agent.Config) (agent.Runner, error) {
218-
if cfg == nil {
219-
return nil, fmt.Errorf("config is required")
220-
}
221-
222-
router := agent.NewRuntimeRouter(agent.RuntimeKindSingerOS)
223-
registered := 0
224-
225-
if cfg.LLM != nil && cfg.LLM.APIKey != "" {
226-
switch cfg.LLM.Provider {
227-
case "", "openai":
228-
logs.Info("Registering SingerOS agent runtime")
229-
singerRunner, err := agent.NewAgent(ctx, cfg.LLM, runtimeConfig)
230-
if err != nil {
231-
return nil, err
232-
}
233-
if err := router.Register(agent.RuntimeKindSingerOS, singerRunner); err != nil {
234-
return nil, err
235-
}
236-
registered++
237-
default:
238-
logs.Warnf("Skipping SingerOS agent runtime for unsupported Eino chat model provider: %s", cfg.LLM.Provider)
239-
}
240-
}
241-
242-
cliRegistry, err := builtin.NewRegistryFromConfig(cfg.CLI)
243-
if err != nil {
244-
return nil, fmt.Errorf("create CLI engine registry: %w", err)
245-
}
246-
cliNames := cliRegistry.Names()
247-
for _, name := range cliNames {
248-
engine, ok := cliRegistry.Get(name)
249-
if !ok {
250-
continue
251-
}
252-
runner, err := externalcli.NewRunner(name, engine, cfg.LLM)
60+
func loadWorkerConfig() (*config.Config, error) {
61+
if workerConfigPath != "" {
62+
cfg := &config.Config{}
63+
err := ygconfig.LoadYamlLocalFile(workerConfigPath, cfg)
25364
if err != nil {
254-
return nil, err
65+
return nil, fmt.Errorf("failed to load config from %s: %w", workerConfigPath, err)
25566
}
256-
if err := router.Register(name, runner); err != nil {
257-
return nil, err
258-
}
259-
registered++
260-
logs.Infof("Registering external agent CLI runtime: %s", name)
67+
return cfg, nil
26168
}
26269

263-
if registered == 0 {
264-
return nil, fmt.Errorf("no agent runtime is available")
265-
}
266-
if cfg.LLM == nil || cfg.LLM.APIKey == "" {
267-
if cfg.CLI != nil && cfg.CLI.Default != "" {
268-
router.SetDefault(cfg.CLI.Default)
269-
} else if len(cliNames) > 0 {
270-
router.SetDefault(cliNames[0])
271-
}
272-
} else {
273-
router.SetDefault(agent.RuntimeKindSingerOS)
274-
}
275-
return router, nil
70+
return &config.Config{}, nil
27671
}

backend/internal/api/router.go

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import (
1616
githubprovider "github.com/insmtx/SingerOS/backend/internal/infra/providers/github"
1717
"github.com/insmtx/SingerOS/backend/internal/infra/websocket"
1818
"github.com/insmtx/SingerOS/backend/internal/service"
19+
"github.com/insmtx/SingerOS/backend/internal/worker/scheduler"
1920
workerserver "github.com/insmtx/SingerOS/backend/internal/worker/server"
2021
singerMCP "github.com/insmtx/SingerOS/backend/mcp"
2122
ygmiddleware "github.com/ygpkg/yg-go/apis/runtime/middleware"
@@ -61,14 +62,18 @@ func SetupRouter(cfg config.Config, publisher eventbus.Publisher, db *gorm.DB) *
6162
websocket.RegisterWebSocketRoutes(v1, publisher)
6263
logs.Info("WebSocket connector registered successfully")
6364

64-
digitalAssistantService := service.NewDigitalAssistantService(db)
65-
handler.RegisterDigitalAssistantRoutes(v1, digitalAssistantService)
66-
logs.Info("Digital assistant routes registered successfully")
65+
workerScheduler := scheduler.NewProcessScheduler(&scheduler.ProcessConfig{
66+
ServerAddr: ":8080",
67+
})
6768

68-
workerServer := workerserver.NewServer()
69+
workerServer := workerserver.NewServer(workerScheduler, db)
6970
workerServer.RegisterRoutes(v1)
7071
logs.Info("Worker server routes registered successfully")
7172

73+
digitalAssistantService := service.NewDigitalAssistantService(db, workerScheduler)
74+
handler.RegisterDigitalAssistantRoutes(v1, digitalAssistantService)
75+
logs.Info("Digital assistant routes registered successfully")
76+
7277
singerMCP.RegisterRoutes(v1, singerMCP.NewServer())
7378
logs.Info("MCP routes registered successfully")
7479

backend/internal/service/digital_assistant_service.go

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,18 +10,22 @@ import (
1010
"github.com/insmtx/SingerOS/backend/internal/api/auth"
1111
"github.com/insmtx/SingerOS/backend/internal/api/contract"
1212
"github.com/insmtx/SingerOS/backend/internal/infra/db"
13+
"github.com/insmtx/SingerOS/backend/internal/worker"
1314
"github.com/insmtx/SingerOS/backend/types"
15+
"github.com/ygpkg/yg-go/logs"
1416
)
1517

1618
var _ contract.DigitalAssistantService = (*digitalAssistantService)(nil)
1719

1820
type digitalAssistantService struct {
19-
db *gorm.DB
21+
db *gorm.DB
22+
workerScheduler worker.WorkerScheduler
2023
}
2124

22-
func NewDigitalAssistantService(db *gorm.DB) contract.DigitalAssistantService {
25+
func NewDigitalAssistantService(db *gorm.DB, workerScheduler worker.WorkerScheduler) contract.DigitalAssistantService {
2326
return &digitalAssistantService{
24-
db: db,
27+
db: db,
28+
workerScheduler: workerScheduler,
2529
}
2630
}
2731

@@ -70,6 +74,17 @@ func (s *digitalAssistantService) CreateDigitalAssistant(ctx context.Context, re
7074
return nil, err
7175
}
7276

77+
if s.workerScheduler != nil && da.Status == string(contract.DigitalAssistantStatusActive) {
78+
spec := &worker.WorkerSpec{
79+
ID: da.Code,
80+
Name: da.Name,
81+
EnvType: worker.WorkerEnvProcess,
82+
}
83+
if _, err := s.workerScheduler.Start(ctx, spec); err != nil {
84+
logs.Warnf("Failed to start worker for assistant %s: %v", da.Code, err)
85+
}
86+
}
87+
7388
return convertToContractDigitalAssistant(da), nil
7489
}
7590

0 commit comments

Comments
 (0)