返回 DeepSeek-Reasonix
bot_runtime_app.go
根目录 / desktop / bot_runtime_app.go
1 package main
2
3 import (
4 "context"
5 "fmt"
6 "log/slog"
7 "os"
8 "strings"
9 "sync"
10 "time"
11
12 "reasonix/internal/bot"
13 "reasonix/internal/botruntime"
14 "reasonix/internal/config"
15 "reasonix/internal/session"
16 )
17
18 type BotRuntimeStatusView struct {
19 Running bool `json:"running"`
20 Status string `json:"status"`
21 Message string `json:"message"`
22 Connections int `json:"connections"`
23 StartedAt string `json:"startedAt"`
24 Platforms map[string]string `json:"platforms,omitempty"`
25 }
26
27 type desktopBotRuntime struct {
28 // lifecycleMu serializes start/stop transitions so two apply/stop calls
29 // can't race a gateway into existence. The slow work (gw.Stop teardown,
30 // gw.Start dials) runs while holding it but NOT r.mu, so status/send reads
31 // never block on a restart.
32 lifecycleMu sync.Mutex
33 mu sync.Mutex
34 cancel context.CancelFunc
35 gw *bot.BotGateway
36 status BotRuntimeStatusView
37 }
38
39 func newDesktopBotRuntime() *desktopBotRuntime {
40 return &desktopBotRuntime{status: BotRuntimeStatusView{Status: "stopped", Message: "bot runtime is not started"}}
41 }
42
43 func desktopBotChannelsWithLegacyQQ(qq config.QQBotConfig, channels map[bot.Platform]bot.ChannelConfig, connectionChannels map[string]bot.ChannelConfig) (map[bot.Platform]bot.ChannelConfig, map[string]bot.ChannelConfig) {
44 channel := bot.ChannelConfig{
45 Model: strings.TrimSpace(qq.Model),
46 ToolApprovalMode: normalizeBotConnectionToolApprovalMode(qq.ToolApprovalMode),
47 WorkspaceRoot: strings.TrimSpace(qq.WorkspaceRoot),
48 }
49 if channel.Model == "" && channel.ToolApprovalMode == "" && channel.WorkspaceRoot == "" {
50 return channels, connectionChannels
51 }
52 if channels == nil {
53 channels = make(map[bot.Platform]bot.ChannelConfig)
54 }
55 if _, ok := channels[bot.PlatformQQ]; !ok {
56 channels[bot.PlatformQQ] = channel
57 }
58 if connectionChannels == nil {
59 connectionChannels = make(map[string]bot.ChannelConfig)
60 }
61 if _, ok := connectionChannels[string(bot.PlatformQQ)]; !ok {
62 connectionChannels[string(bot.PlatformQQ)] = channel
63 }
64 return channels, connectionChannels
65 }
66
67 // desktopBotChannelsWithLegacyDingtalk 把 legacy [bot.dingtalk] 的模型/权限/
68 // 工作目录合成进 Channels 与 ConnectionChannels,使直配(无 connection)的
69 // 钉钉 bot 也能从设置面板配置这些运行选项(与 legacy QQ 同路径)。
70 func desktopBotChannelsWithLegacyDingtalk(dt config.DingtalkBotConfig, channels map[bot.Platform]bot.ChannelConfig, connectionChannels map[string]bot.ChannelConfig) (map[bot.Platform]bot.ChannelConfig, map[string]bot.ChannelConfig) {
71 channel := bot.ChannelConfig{
72 Model: strings.TrimSpace(dt.Model),
73 ToolApprovalMode: normalizeBotConnectionToolApprovalMode(dt.ToolApprovalMode),
74 WorkspaceRoot: strings.TrimSpace(dt.WorkspaceRoot),
75 SessionMappings: botruntime.SessionMappings(dt.SessionMappings),
76 }
77 if channel.Model == "" && channel.ToolApprovalMode == "" && channel.WorkspaceRoot == "" && len(channel.SessionMappings) == 0 {
78 return channels, connectionChannels
79 }
80 if channels == nil {
81 channels = make(map[bot.Platform]bot.ChannelConfig)
82 }
83 if _, ok := channels[bot.PlatformDingtalk]; !ok {
84 channels[bot.PlatformDingtalk] = channel
85 }
86 if connectionChannels == nil {
87 connectionChannels = make(map[string]bot.ChannelConfig)
88 }
89 if _, ok := connectionChannels[string(bot.PlatformDingtalk)]; !ok {
90 connectionChannels[string(bot.PlatformDingtalk)] = channel
91 }
92 return channels, connectionChannels
93 }
94
95 func (a *App) refreshBotRuntimeAsync() {
96 if a.ctx == nil {
97 return
98 }
99 a.goSafe("refreshBotRuntime", a.refreshBotRuntime)
100 }
101
102 func (a *App) refreshBotRuntime() {
103 // NewApp always pre-fills botRuntime; a nil here means a test-constructed
104 // App with no bot runtime, which must not lazily create one from a
105 // background goroutine (that would race a concurrent refresh).
106 if a.botRuntime == nil {
107 return
108 }
109 var watcherVersion uint64
110 if a.botBridge != nil {
111 watcherVersion = a.botBridge.watcherVersion()
112 }
113 cfg, err := a.loadDesktopBotConfig()
114 if err != nil {
115 a.botRuntime.stop("error", err.Error())
116 return
117 }
118 // Assign through a typed local so a nil *botBridgeHub never becomes a
119 // non-nil bot.DesktopBridge interface inside the gateway config.
120 var bridge bot.DesktopBridge
121 if a.botBridge != nil {
122 // 配置是订阅的持久化事实源:每次运行时重算前重新种子,桌面重启后
123 // /desktop watch 的订阅继续生效。
124 a.botBridge.seedWatchers(bridgeRoutesFromConfig(cfg.Bot.DesktopWatchers), watcherVersion)
125 bridge = a.botBridge
126 }
127 _ = a.botRuntime.apply(a.bootContext(), cfg, globalTabWorkspaceRoot(), a.persistRemoteBotToolApprovalMode, bridge, a.historicalSessionService)
128 }
129
130 func (a *App) loadDesktopBotConfig() (*config.Config, error) {
131 // Read-only load feeding the bot runtime and connection diagnostics. It
132 // must load credentials: the runtime resolves app secrets and control
133 // tokens from the process env (AppSecretEnv, Control.TokenEnv), which the
134 // credential-free view load would leave unset on a fresh process.
135 cfg, _, err := a.loadDesktopUserConfigForViewWithCredentials()
136 if err != nil {
137 return nil, err
138 }
139 return cfg, nil
140 }
141
142 func (a *App) stopBotRuntime() {
143 if a.botRuntime != nil {
144 a.botRuntime.stop("stopped", "bot runtime stopped")
145 }
146 }
147
148 func (a *App) BotRuntimeStatus() BotRuntimeStatusView {
149 if a.botRuntime == nil {
150 return BotRuntimeStatusView{Status: "stopped", Message: "bot runtime is not started"}
151 }
152 return a.botRuntime.snapshot()
153 }
154
155 func (r *desktopBotRuntime) apply(parent context.Context, cfg *config.Config, workspaceRoot string, onToolApprovalModeChange func(bot.InboundMessage, string) error, bridge bot.DesktopBridge, sessionServiceForRoot func(string) (*session.Service, error)) error {
156 if r == nil {
157 return nil
158 }
159 if parent == nil {
160 parent = context.Background()
161 }
162 plan := desktopBotRuntimePlan(cfg)
163 r.lifecycleMu.Lock()
164 defer r.lifecycleMu.Unlock()
165 r.stopCurrent()
166 if !plan.Start {
167 r.setStatus(BotRuntimeStatusView{Status: plan.Status, Message: plan.Message})
168 return nil
169 }
170
171 logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo}))
172 ctx, cancel := context.WithCancel(parent)
173 modelName := botruntime.ModelName(cfg, "")
174 channels := botruntime.ChannelConfigs(cfg.Bot.Connections, true, true)
175 connectionChannels := botruntime.ConnectionChannelConfigs(cfg.Bot.Connections, true, true)
176 channels, connectionChannels = desktopBotChannelsWithLegacyQQ(cfg.Bot.QQ, channels, connectionChannels)
177 channels, connectionChannels = desktopBotChannelsWithLegacyDingtalk(cfg.Bot.Dingtalk, channels, connectionChannels)
178 gwCfg := bot.GatewayConfig{
179 Model: modelName,
180 ToolApprovalMode: cfg.Bot.ToolApprovalMode,
181 MaxSteps: cfg.Bot.MaxSteps,
182 QueueMode: cfg.Bot.QueueMode,
183 QueueCap: cfg.Bot.QueueCap,
184 QueueDrop: cfg.Bot.QueueDrop,
185 PairingEnabled: cfg.Bot.Pairing.Enabled,
186 PairingTTL: time.Duration(cfg.Bot.Pairing.RequestTTLMinutes) * time.Minute,
187 PairingMaxPending: cfg.Bot.Pairing.MaxPendingPerPlatform,
188 IgnoreSelfMessages: cfg.Bot.IgnoreSelfMessages,
189 SelfUserIDs: map[bot.Platform][]string{
190 bot.PlatformQQ: cfg.Bot.SelfUserIDs.QQ,
191 bot.PlatformFeishu: cfg.Bot.SelfUserIDs.Feishu,
192 bot.PlatformWeixin: cfg.Bot.SelfUserIDs.Weixin,
193 bot.PlatformDingtalk: cfg.Bot.SelfUserIDs.Dingtalk,
194 },
195 ControlEnabled: cfg.Bot.Control.Enabled,
196 ControlAddr: cfg.Bot.Control.Addr,
197 ControlToken: os.Getenv(strings.TrimSpace(cfg.Bot.Control.TokenEnv)),
198 WorkspaceRoot: workspaceRoot,
199 Channels: channels,
200 ConnectionChannels: connectionChannels,
201 Routes: botruntime.RouteConfigs(cfg.Bot.Routes, true, true),
202 ConnectionAccess: botruntime.ConnectionAccessConfigs(cfg),
203 Enabled: plan.Enabled,
204 Allowlist: bot.AllowlistConfig{
205 Enabled: cfg.Bot.Allowlist.Enabled,
206 AllowAll: cfg.Bot.Allowlist.AllowAll,
207 Users: map[bot.Platform][]string{
208 bot.PlatformQQ: cfg.Bot.Allowlist.QQUsers,
209 bot.PlatformFeishu: cfg.Bot.Allowlist.FeishuUsers,
210 bot.PlatformWeixin: cfg.Bot.Allowlist.WeixinUsers,
211 bot.PlatformDingtalk: cfg.Bot.Allowlist.DingtalkUsers,
212 },
213 Approvers: map[bot.Platform][]string{
214 bot.PlatformQQ: cfg.Bot.Allowlist.QQApprovers,
215 bot.PlatformFeishu: cfg.Bot.Allowlist.FeishuApprovers,
216 bot.PlatformWeixin: cfg.Bot.Allowlist.WeixinApprovers,
217 bot.PlatformDingtalk: cfg.Bot.Allowlist.DingtalkApprovers,
218 },
219 Admins: map[bot.Platform][]string{
220 bot.PlatformQQ: cfg.Bot.Allowlist.QQAdmins,
221 bot.PlatformFeishu: cfg.Bot.Allowlist.FeishuAdmins,
222 bot.PlatformWeixin: cfg.Bot.Allowlist.WeixinAdmins,
223 bot.PlatformDingtalk: cfg.Bot.Allowlist.DingtalkAdmins,
224 },
225 Groups: map[bot.Platform][]string{
226 bot.PlatformQQ: cfg.Bot.Allowlist.QQGroups,
227 bot.PlatformFeishu: cfg.Bot.Allowlist.FeishuGroups,
228 bot.PlatformWeixin: cfg.Bot.Allowlist.WeixinGroups,
229 bot.PlatformDingtalk: cfg.Bot.Allowlist.DingtalkGroups,
230 },
231 },
232 Debounce: time.Duration(cfg.Bot.DebounceMs) * time.Millisecond,
233 ModelResolver: botruntime.ModelResolver(cfg),
234 OnInbound: botruntime.NewRemoteRememberer(logger),
235 OnSessionReady: botruntime.NewSessionRemembererWithWorkspace(logger, workspaceRoot),
236 OnToolApprovalModeChange: onToolApprovalModeChange,
237 Desktop: bridge,
238 SessionServiceForRoot: sessionServiceForRoot,
239 }
240 bindings := botruntime.AdapterBindings(cfg, plan.Enabled, nil, logger)
241 if len(bindings) == 0 {
242 cancel()
243 r.setStatus(BotRuntimeStatusView{Status: "stopped", Message: "no bot adapters configured"})
244 return nil
245 }
246 gw := bot.NewGatewayWithAdapterBindings(gwCfg, bindings, logger)
247 if err := gw.Start(ctx); err != nil {
248 cancel()
249 gw.Stop()
250 r.setStatus(BotRuntimeStatusView{Status: "error", Message: err.Error(), Connections: gw.AdapterCount()})
251 return err
252 }
253 runningConnections := gw.AdapterCount()
254 startErrors := gw.StartErrors()
255 status := "running"
256 message := fmt.Sprintf("%d bot connection(s) running", runningConnections)
257 if len(startErrors) > 0 {
258 status = "degraded"
259 message = fmt.Sprintf("%d bot connection(s) running; %d failed to start: %s", runningConnections, len(startErrors), summarizeBotRuntimeErrors(startErrors))
260 }
261 r.mu.Lock()
262 r.cancel = cancel
263 r.gw = gw
264 r.status = BotRuntimeStatusView{
265 Running: true,
266 Status: status,
267 Message: message,
268 Connections: runningConnections,
269 StartedAt: time.Now().UTC().Format(time.RFC3339),
270 }
271 r.mu.Unlock()
272 return nil
273 }
274
275 func (a *App) persistRemoteBotToolApprovalMode(msg bot.InboundMessage, mode string) error {
276 mode = normalizeBotConnectionToolApprovalMode(mode)
277 if mode == "" {
278 return nil
279 }
280 return a.applyConfigOnly(func(c *config.Config) error {
281 id := strings.TrimSpace(msg.ConnectionID)
282 now := time.Now().UTC().Format(time.RFC3339)
283 if id != "" {
284 for i := range c.Bot.Connections {
285 if c.Bot.Connections[i].ID == id || botruntime.ConnectionRuntimeID(c.Bot.Connections[i]) == id {
286 c.Bot.Connections[i].ToolApprovalMode = mode
287 c.Bot.Connections[i].UpdatedAt = now
288 return nil
289 }
290 }
291 }
292 c.Bot.ToolApprovalMode = mode
293 return nil
294 })
295 }
296
297 func summarizeBotRuntimeErrors(errs []error) string {
298 parts := make([]string, 0, len(errs))
299 for _, err := range errs {
300 if err == nil {
301 continue
302 }
303 parts = append(parts, err.Error())
304 }
305 if len(parts) == 0 {
306 return ""
307 }
308 if len(parts) > 3 {
309 hidden := len(parts) - 3
310 parts = append(parts[:3], fmt.Sprintf("%d more", hidden))
311 }
312 return strings.Join(parts, "; ")
313 }
314
315 type botRuntimePlan struct {
316 Start bool
317 Status string
318 Message string
319 Enabled map[bot.Platform]bool
320 }
321
322 func desktopBotRuntimePlan(cfg *config.Config) botRuntimePlan {
323 if cfg == nil {
324 return botRuntimePlan{Status: "error", Message: "config is unavailable"}
325 }
326 if !cfg.Bot.Enabled {
327 return botRuntimePlan{Status: "stopped", Message: "bot is disabled"}
328 }
329 if !botruntime.BotConfigHasAccessControl(cfg.Bot) {
330 return botRuntimePlan{Status: "blocked", Message: "bot requires an allowlist, pairing, per-bot access, or allow_all=true"}
331 }
332 enabled, unknown := botruntime.EnabledPlatforms(cfg, nil)
333 if len(unknown) > 0 {
334 return botRuntimePlan{Status: "error", Message: "unknown bot channel: " + strings.Join(unknown, ", ")}
335 }
336 if !botruntime.HasEnabledPlatform(enabled) {
337 return botRuntimePlan{Status: "stopped", Message: "no bot channels enabled"}
338 }
339 return botRuntimePlan{Start: true, Status: "running", Message: "bot runtime can start", Enabled: enabled}
340 }
341
342 func (r *desktopBotRuntime) stop(status, message string) {
343 r.lifecycleMu.Lock()
344 defer r.lifecycleMu.Unlock()
345 r.stopCurrent()
346 r.setStatus(BotRuntimeStatusView{Status: status, Message: message})
347 }
348
349 // stopCurrent detaches the running gateway under r.mu, then tears it down
350 // off-lock: gw.Stop() closes every session controller (up to the jobs teardown
351 // grace each) and must not stall status/send readers. Callers hold lifecycleMu.
352 func (r *desktopBotRuntime) stopCurrent() {
353 r.mu.Lock()
354 cancel := r.cancel
355 gw := r.gw
356 r.cancel = nil
357 r.gw = nil
358 r.mu.Unlock()
359 if cancel != nil {
360 cancel()
361 }
362 if gw != nil {
363 gw.Stop()
364 }
365 }
366
367 func (r *desktopBotRuntime) setStatus(status BotRuntimeStatusView) {
368 r.mu.Lock()
369 r.status = status
370 r.mu.Unlock()
371 }
372
373 func (r *desktopBotRuntime) snapshot() BotRuntimeStatusView {
374 r.mu.Lock()
375 defer r.mu.Unlock()
376 s := r.status
377 if r.gw != nil {
378 s.Platforms = botAdapterPlatformStatuses(r.gw.AdapterHealth())
379 }
380 return s
381 }
382
383 // botAdapterPlatformStatuses 把 gateway 的适配器健康快照收敛为
384 // platform → status 映射(如 dingtalk → running),供设置面板显示在线状态。
385 func botAdapterPlatformStatuses(health []bot.AdapterHealthSnapshot) map[string]string {
386 out := make(map[string]string, len(health))
387 for _, h := range health {
388 if strings.TrimSpace(string(h.Platform)) != "" {
389 out[string(h.Platform)] = h.Status
390 }
391 }
392 return out
393 }
394
395 // updateConnectionToolApprovalMode updates a connection's tool approval mode
396 // on the running gateway without restarting. Returns true if updated, false if
397 // the gateway is not running or the connection is unknown.
398 func (r *desktopBotRuntime) updateConnectionToolApprovalMode(connID, mode string) bool {
399 r.mu.Lock()
400 defer r.mu.Unlock()
401 if r.gw == nil {
402 return false
403 }
404 mode = normalizeBotConnectionToolApprovalMode(mode)
405 // Update ConnectionChannels in the internal GatewayConfig so new sessions
406 // pick up the mode. Existing sessions are updated by the gateway directly.
407 r.gw.UpdateConnectionToolApprovalMode(connID, mode)
408 return true
409 }
410
411 // SendToAdapter sends a message through the running gateway's adapter
412 // identified by connID. Returns an error if the gateway is not running
413 // or no matching adapter is found.
414 func (r *desktopBotRuntime) SendToAdapter(ctx context.Context, connID, domain string, msg bot.OutboundMessage) (bot.SendResult, error) {
415 r.mu.Lock()
416 gw := r.gw
417 r.mu.Unlock()
418 if gw == nil {
419 return bot.SendResult{}, nil // gateway not running — silent no-op
420 }
421 return gw.SendToAdapter(ctx, connID, domain, msg)
422 }
423
424 // TestSendToAdapter sends a test message through the running gateway's adapter
425 // identified by connID. The adapter must implement bot.TestSender (dingtalk).
426 func (r *desktopBotRuntime) TestSendToAdapter(ctx context.Context, connID, domain, text string) (bot.SendResult, error) {
427 r.mu.Lock()
428 gw := r.gw
429 r.mu.Unlock()
430 if gw == nil {
431 return bot.SendResult{}, fmt.Errorf("bot runtime is not running")
432 }
433 return gw.TestSendToAdapter(ctx, connID, domain, text)
434 }
435
436 // Running returns true if the bot gateway is currently active.
437 func (r *desktopBotRuntime) Running() bool {
438 r.mu.Lock()
439 defer r.mu.Unlock()
440 return r.gw != nil
441 }
442
443 // ForwardTargets returns the list of bot forward targets derived from the
444 // current config's bot connections and their session mappings. Each mapping
445 // produces one target (connID + chatID + chatType) for event forwarding.
446 func (r *desktopBotRuntime) ForwardTargets(cfg *config.Config) []botForwardTarget {
447 if cfg == nil {
448 return nil
449 }
450 var targets []botForwardTarget
451 seen := make(map[botForwardTarget]bool)
452 for _, conn := range cfg.Bot.Connections {
453 if !conn.Enabled {
454 continue
455 }
456 connID := botruntime.ConnectionRuntimeID(conn)
457 domain := strings.TrimSpace(conn.Domain)
458 for _, sm := range conn.SessionMappings {
459 remoteID := strings.TrimSpace(sm.RemoteID)
460 if remoteID == "" {
461 continue
462 }
463 chatType := bot.ChatDM
464 if sm.ChatType != "" {
465 chatType = bot.ChatType(sm.ChatType)
466 }
467 target := botForwardTarget{
468 ConnID: connID,
469 Domain: domain,
470 ChatID: remoteID,
471 ChatType: chatType,
472 }
473 if seen[target] {
474 continue
475 }
476 seen[target] = true
477 targets = append(targets, target)
478 }
479 }
480 return targets
481 }
482
482 lines GO