返回 DeepSeek-Reasonix
gateway.go
根目录 / internal / bot / gateway.go
1 package bot
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "os"
9 "path/filepath"
10 "sort"
11 "strconv"
12 "strings"
13 "sync"
14 "time"
15
16 "reasonix/internal/agent"
17 "reasonix/internal/boot"
18 "reasonix/internal/config"
19 "reasonix/internal/control"
20 "reasonix/internal/event"
21 "reasonix/internal/secrets"
22 "reasonix/internal/session"
23 "reasonix/internal/sessioninbox"
24 )
25
26 // GatewayConfig 是 BotGateway 的配置。
27 type GatewayConfig struct {
28 Model string
29 ToolApprovalMode string
30 MaxSteps int
31 QueueMode string
32 QueueCap int
33 QueueDrop string
34 PairingEnabled bool
35 PairingTTL time.Duration
36 PairingMaxPending int
37 // ModelResolver 校验模型引用是否可解析且已配置(provider 存在、模型
38 // 存在、API key 已配)。/model 切换前调用:校验失败则不写入覆盖,
39 // 保留当前 controller 可继续聊天(失败原子性)。nil 时跳过预校验。
40 ModelResolver func(ref string) error
41 // IgnoreSelfMessages drops messages that are clearly sent by this bot. It
42 // uses configured SelfUserIDs plus recently returned outbound message IDs.
43 IgnoreSelfMessages bool
44 SelfUserIDs map[Platform][]string
45 ControlEnabled bool
46 ControlAddr string
47 ControlToken string
48 // ApprovalTimeout bounds how long a tool-approval/ask prompt blocks a bot
49 // session waiting for a remote user's reply. Zero falls back to
50 // defaultBotApprovalTimeout so an abandoned prompt can't wedge the bot forever
51 // (#4626, #4402). A negative value disables the timeout (wait indefinitely).
52 ApprovalTimeout time.Duration
53 WorkspaceRoot string
54 Channels map[Platform]ChannelConfig
55 ConnectionChannels map[string]ChannelConfig
56 Routes []RouteConfig
57 ConnectionAccess map[string]AccessConfig
58 Allowlist AllowlistConfig
59 Enabled map[Platform]bool
60 Debounce time.Duration
61 // OnInbound observes every allowlisted inbound message before dispatch.
62 //
63 // Reentrancy contract for all GatewayConfig callbacks (OnInbound,
64 // OnSessionReady, OnToolApprovalModeChange): they run synchronously on
65 // gateway-owned dispatch/turn goroutines; OnSessionReady can also run on a
66 // controller recovery/autosave goroutine. Stop drains all of those paths
67 // before returning. A callback must therefore never call Stop, nor block
68 // until a goroutine that does so completes — Stop would wait on the very
69 // goroutine running the callback, a guaranteed deadlock. Hosts that want to
70 // shut the gateway down in reaction to a callback must trigger the shutdown
71 // asynchronously.
72 OnInbound func(InboundMessage)
73 // OnSessionReady notifies the host after the bot has created, reused, or
74 // recovered the controller for an inbound remote. Hosts may persist the
75 // concrete session ID or keep the remote as a read-only channel.
76 OnSessionReady func(InboundMessage, string) error
77 // OnToolApprovalModeChange persists a remote IM /mode request.
78 // The gateway updates the live session and in-memory defaults first; this
79 // callback lets desktop save the chosen connection mode to user config.
80 OnToolApprovalModeChange func(InboundMessage, string) error
81 // Desktop, when the gateway is embedded in the desktop app, gives bot
82 // chats a god view over desktop sessions (/desktop commands): global
83 // status, event subscriptions, and remote approvals for any live desktop
84 // session. Nil when the gateway runs standalone (reasonix bot start).
85 Desktop DesktopBridge
86 // SessionServiceForRoot lets an embedded host share ownership of a session
87 // store. The host keeps the returned service alive until after Stop drains
88 // all gateway controllers; the gateway never shuts that service down.
89 SessionServiceForRoot func(string) (*session.Service, error)
90 }
91
92 // ChannelConfig overrides gateway defaults for one IM channel.
93 type ChannelConfig struct {
94 Model string
95 ToolApprovalMode string
96 WorkspaceRoot string
97 SessionMappings []SessionMapping
98 }
99
100 // SessionMapping is the runtime subset of a saved bot connection mapping used
101 // to route a remote chat/user/thread back to its intended workspace.
102 type SessionMapping struct {
103 RemoteID string
104 SessionID string
105 SessionSource string
106 ChatType string
107 UserID string
108 ThreadID string
109 Scope string
110 WorkspaceRoot string
111 UpdatedAt string
112 }
113
114 // RouteConfig applies per-remote overrides. Empty match fields are wildcards;
115 // the first matching route wins.
116 type RouteConfig struct {
117 ConnectionID string
118 Platform Platform
119 ChatType ChatType
120 ChatID string
121 UserID string
122 ThreadID string
123 Channel ChannelConfig
124 }
125
126 // AdapterBinding attaches an adapter instance to one saved bot connection.
127 // Feishu and Lark share PlatformFeishu, so ID/Domain keep their sessions,
128 // replies, and per-connection settings separated at runtime.
129 type AdapterBinding struct {
130 ID string
131 Domain string
132 Platform Platform
133 Adapter Adapter
134 }
135
136 // AllowlistConfig 控制哪些用户/群可以使用 bot。
137 type AllowlistConfig struct {
138 Enabled bool
139 AllowAll bool
140 Users map[Platform][]string
141 Approvers map[Platform][]string
142 Admins map[Platform][]string
143 Groups map[Platform][]string
144 }
145
146 // AccessConfig controls who may use one concrete bot connection.
147 type AccessConfig struct {
148 Enabled bool
149 AllowAll bool
150 PairingEnabled bool
151 Users []string
152 Groups []string
153 Approvers []string
154 Admins []string
155 }
156
157 // AdapterHealthSnapshot describes the gateway's current view of one adapter.
158 type AdapterHealthSnapshot struct {
159 ID string `json:"id"`
160 Platform Platform `json:"platform"`
161 Domain string `json:"domain,omitempty"`
162 Name string `json:"name,omitempty"`
163 Status string `json:"status"`
164 StartedAt time.Time `json:"started_at,omitempty"`
165 LastMessageAt time.Time `json:"last_message_at,omitempty"`
166 LastSendAt time.Time `json:"last_send_at,omitempty"`
167 LastErrorAt time.Time `json:"last_error_at,omitempty"`
168 LastError string `json:"last_error,omitempty"`
169 Messages int64 `json:"messages"`
170 Sends int64 `json:"sends"`
171 SendErrors int64 `json:"send_errors"`
172 Closed bool `json:"closed"`
173 }
174
175 // BotGateway 是 reasonix bot 消息网关,管理 Controller 生命周期、session 并发、
176 // 事件渲染和平台适配器。
177 type BotGateway struct {
178 cfg GatewayConfig
179 adapters []AdapterBinding
180 sessions *SessionManager
181 startErr []error
182
183 lifecycleMu sync.Mutex
184 started bool
185 stopped bool
186 runCancel context.CancelFunc
187 startDone chan struct{}
188 stopDone chan struct{}
189 gatewayWG sync.WaitGroup
190 turnWG sync.WaitGroup
191
192 mu sync.Mutex
193 controllers map[string]*sessionState // session key -> active state
194 pendingReactionCleanups map[string][]func()
195 allowlist map[Platform]map[string]bool
196 groupAllowlist map[Platform]map[string]bool
197 selfUserIDs map[Platform]map[string]bool
198 outboundMessageIDs map[string]time.Time
199 adapterHealth map[string]*AdapterHealthSnapshot
200 controlServer *controlHTTPServer
201 sessionOverrides map[string]sessionRuntimeOverride
202 buildController func(context.Context, boot.Options) (*control.Controller, error)
203
204 sessionServicesMu sync.Mutex
205 sessionServices map[string]*session.Service
206
207 logger *slog.Logger
208 }
209
210 // botController is the slice of the controller's driving port the gateway needs:
211 // session lifecycle, turn execution, and approval/ask handling. The bot never
212 // touches goals, checkpoints, or memory, so it depends on those sub-ports only —
213 // not the concrete *control.Controller and its ~99 methods.
214 type botController interface {
215 control.Lifecycle
216 control.TurnControl
217 control.Approvals
218 }
219
220 type sessionState struct {
221 lifecycleMu sync.Mutex
222 retired bool
223 ctrl botController
224 sink *sessionEventSink
225 leases *control.SessionLeaseKeeper
226 platform Platform
227 connectionID string
228 model string
229 workspaceRoot string
230 toolApprovalMode string
231 sessionPath string
232 sessionRef session.SessionRef
233 releaseRuntimeOnly bool
234 onSessionTransition func(control.SessionTransitionInfo) error
235 // mappingDegraded records that this state intentionally runs on a fresh
236 // session because its session_mappings target could not be used at build
237 // time. It keeps later messages (whose profile re-resolves the mapping)
238 // from tearing the state down every turn; convergence back onto the
239 // mapped file happens on the next gateway restart.
240 mappingDegraded bool
241 cancel context.CancelFunc
242 pendingAsks map[string][]event.AskQuestion
243 pendingApprovals map[string]event.Approval
244 lastApprovalID string
245 lastAskID string
246 createdAt time.Time
247 lastActive time.Time
248 }
249
250 var errBotSessionRetired = errors.New("bot session retired during recovery")
251
252 type sessionRuntimeProfile struct {
253 model string
254 workspaceRoot string
255 toolApprovalMode string
256 sessionPath string
257 sessionRef session.SessionRef
258 // sessionPathOptional marks sessionPath as a persisted session_mappings
259 // binding rather than an explicit /attach: when the mapped file cannot be
260 // loaded or leased, the session degrades to a fresh path instead of
261 // dropping the message (#6917).
262 sessionPathOptional bool
263 sessionRefOptional bool
264 }
265
266 type sessionRuntimeOverride struct {
267 channel ChannelConfig
268 sessionPath string
269 label string
270 }
271
272 type sessionEventSink struct {
273 mu sync.RWMutex
274 target event.Sink
275 }
276
277 type pendingReactionAdapter interface {
278 AddPendingReaction(ctx context.Context, messageID string) (func(), error)
279 }
280
281 const outboundEchoTTL = 10 * time.Minute
282
283 func (s *sessionEventSink) setTarget(target event.Sink) {
284 s.mu.Lock()
285 defer s.mu.Unlock()
286 s.target = target
287 }
288
289 func (s *sessionEventSink) Emit(e event.Event) {
290 s.mu.RLock()
291 target := s.target
292 s.mu.RUnlock()
293 if target != nil {
294 target.Emit(e)
295 }
296 }
297
298 // NewGateway 创建一个新的 BotGateway。
299 func NewGateway(cfg GatewayConfig, adapters map[Platform]Adapter, logger *slog.Logger) *BotGateway {
300 bindings := make([]AdapterBinding, 0, len(adapters))
301 for plat, adapter := range adapters {
302 bindings = append(bindings, AdapterBinding{ID: string(plat), Platform: plat, Adapter: adapter})
303 }
304 return NewGatewayWithAdapterBindings(cfg, bindings, logger)
305 }
306
307 // NewGatewayWithAdapterBindings creates a gateway with one or more adapter
308 // instances per platform.
309 func NewGatewayWithAdapterBindings(cfg GatewayConfig, adapters []AdapterBinding, logger *slog.Logger) *BotGateway {
310 if logger == nil {
311 logger = slog.Default()
312 }
313 if cfg.Debounce <= 0 {
314 cfg.Debounce = 1500 * time.Millisecond
315 }
316 cfg.QueueMode = NormalizeQueueMode(cfg.QueueMode)
317 if cfg.QueueCap <= 0 {
318 cfg.QueueCap = DefaultQueueCap
319 }
320 cfg.QueueDrop = NormalizeQueueDrop(cfg.QueueDrop)
321 if cfg.PairingTTL <= 0 {
322 cfg.PairingTTL = defaultPairingTTL
323 }
324 if cfg.PairingMaxPending <= 0 {
325 cfg.PairingMaxPending = defaultPairingMaxPending
326 }
327 gw := &BotGateway{
328 cfg: cfg,
329 adapters: normalizeAdapterBindings(adapters),
330 sessions: NewSessionManager(cfg.Debounce),
331 controllers: make(map[string]*sessionState),
332 pendingReactionCleanups: make(map[string][]func()),
333 allowlist: make(map[Platform]map[string]bool),
334 groupAllowlist: make(map[Platform]map[string]bool),
335 selfUserIDs: make(map[Platform]map[string]bool),
336 outboundMessageIDs: make(map[string]time.Time),
337 adapterHealth: make(map[string]*AdapterHealthSnapshot),
338 sessionOverrides: make(map[string]sessionRuntimeOverride),
339 sessionServices: make(map[string]*session.Service),
340 buildController: boot.Build,
341 logger: logger.With("component", "bot_gateway"),
342 }
343 gw.buildAllowlist()
344 gw.buildSelfUserIDs()
345 for _, binding := range gw.adapters {
346 gw.setAdapterConfigured(binding)
347 }
348 return gw
349 }
350
351 func normalizeAdapterBindings(adapters []AdapterBinding) []AdapterBinding {
352 out := make([]AdapterBinding, 0, len(adapters))
353 for _, binding := range adapters {
354 if binding.Adapter == nil {
355 continue
356 }
357 if binding.Platform == "" {
358 binding.Platform = binding.Adapter.Platform()
359 }
360 if strings.TrimSpace(binding.ID) == "" {
361 binding.ID = string(binding.Platform)
362 }
363 binding.ID = strings.TrimSpace(binding.ID)
364 binding.Domain = strings.TrimSpace(binding.Domain)
365 out = append(out, binding)
366 }
367 return out
368 }
369
370 func (gw *BotGateway) buildAllowlist() {
371 for _, plat := range []Platform{PlatformQQ, PlatformFeishu, PlatformWeixin, PlatformDingtalk} {
372 gw.allowlist[plat] = make(map[string]bool)
373 if !gw.cfg.Allowlist.Enabled {
374 continue
375 }
376 addAllowlistUsers(gw.allowlist[plat], gw.cfg.Allowlist.Users[plat])
377 addAllowlistUsers(gw.allowlist[plat], gw.cfg.Allowlist.Admins[plat])
378 addAllowlistUsers(gw.allowlist[plat], gw.cfg.Allowlist.Approvers[plat])
379 gw.groupAllowlist[plat] = make(map[string]bool)
380 for _, gid := range gw.cfg.Allowlist.Groups[plat] {
381 gw.groupAllowlist[plat][gid] = true
382 }
383 }
384 }
385
386 func addAllowlistUsers(dst map[string]bool, users []string) {
387 for _, uid := range users {
388 uid = strings.TrimSpace(uid)
389 if uid != "" {
390 dst[uid] = true
391 }
392 }
393 }
394
395 func (gw *BotGateway) buildSelfUserIDs() {
396 for _, plat := range []Platform{PlatformQQ, PlatformFeishu, PlatformWeixin, PlatformDingtalk} {
397 gw.selfUserIDs[plat] = stringSet(gw.cfg.SelfUserIDs[plat])
398 }
399 }
400
401 // Start 启动所有已启用的平台适配器并开始处理消息。
402 func (gw *BotGateway) Start(ctx context.Context) (err error) {
403 gw.lifecycleMu.Lock()
404 if gw.stopped {
405 gw.lifecycleMu.Unlock()
406 return errors.New("bot gateway already stopped")
407 }
408 if gw.started {
409 gw.lifecycleMu.Unlock()
410 return errors.New("bot gateway already started")
411 }
412 gw.started = true
413 runCtx, cancel := context.WithCancel(ctx)
414 gw.runCancel = cancel
415 startDone := make(chan struct{})
416 gw.startDone = startDone
417 gw.lifecycleMu.Unlock()
418 defer func() {
419 if err != nil {
420 cancel()
421 }
422 gw.lifecycleMu.Lock()
423 if err != nil {
424 gw.runCancel = nil
425 }
426 close(startDone)
427 gw.lifecycleMu.Unlock()
428 }()
429
430 started := make([]AdapterBinding, 0, len(gw.adapters))
431 var startErr []error
432 for _, binding := range gw.adapters {
433 if !gw.cfg.Enabled[binding.Platform] {
434 gw.logger.Info("platform disabled, skipping", "platform", binding.Platform, "connection", binding.ID)
435 gw.markAdapterDisabled(binding)
436 continue
437 }
438 gw.logger.Info("starting adapter", "platform", binding.Platform, "connection", binding.ID, "domain", binding.Domain)
439 if err := binding.Adapter.Start(runCtx); err != nil {
440 wrapped := fmt.Errorf("start adapter %s: %w", binding.ID, err)
441 startErr = append(startErr, wrapped)
442 gw.markAdapterStartFailed(binding, err)
443 gw.logger.Warn("adapter start failed", "platform", binding.Platform, "connection", binding.ID, "domain", binding.Domain, "err", err)
444 continue
445 }
446 gw.markAdapterStarted(binding)
447 started = append(started, binding)
448 }
449 // SendToAdapter reads gw.adapters under gw.mu; publish the started set under
450 // the same lock.
451 gw.mu.Lock()
452 gw.adapters = started
453 gw.startErr = startErr
454 gw.mu.Unlock()
455 if len(started) == 0 && len(startErr) > 0 {
456 return errors.Join(startErr...)
457 }
458 if err := gw.startControlServer(runCtx); err != nil {
459 for _, binding := range started {
460 _ = binding.Adapter.Stop()
461 }
462 return err
463 }
464
465 // 合并所有适配器的消息通道
466 for _, binding := range gw.adapters {
467 gw.gatewayWG.Go(func() {
468 gw.dispatchLoop(runCtx, binding)
469 })
470 }
471
472 return nil
473 }
474
475 func (gw *BotGateway) AdapterCount() int {
476 gw.mu.Lock()
477 defer gw.mu.Unlock()
478 return len(gw.adapters)
479 }
480
481 func (gw *BotGateway) StartErrors() []error {
482 gw.mu.Lock()
483 defer gw.mu.Unlock()
484 out := make([]error, len(gw.startErr))
485 copy(out, gw.startErr)
486 return out
487 }
488
489 // AdapterHealth returns a stable snapshot of all configured adapter instances.
490 func (gw *BotGateway) AdapterHealth() []AdapterHealthSnapshot {
491 gw.mu.Lock()
492 defer gw.mu.Unlock()
493 out := make([]AdapterHealthSnapshot, 0, len(gw.adapterHealth))
494 for _, health := range gw.adapterHealth {
495 if health == nil {
496 continue
497 }
498 out = append(out, *health)
499 }
500 sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
501 return out
502 }
503
504 func (gw *BotGateway) setAdapterConfigured(binding AdapterBinding) {
505 gw.mu.Lock()
506 defer gw.mu.Unlock()
507 gw.ensureAdapterHealthLocked(binding).Status = "configured"
508 }
509
510 func (gw *BotGateway) markAdapterDisabled(binding AdapterBinding) {
511 gw.mu.Lock()
512 defer gw.mu.Unlock()
513 health := gw.ensureAdapterHealthLocked(binding)
514 health.Status = "disabled"
515 health.Closed = true
516 }
517
518 func (gw *BotGateway) markAdapterStarted(binding AdapterBinding) {
519 now := time.Now()
520 gw.mu.Lock()
521 defer gw.mu.Unlock()
522 health := gw.ensureAdapterHealthLocked(binding)
523 health.Status = "running"
524 health.StartedAt = now
525 health.LastError = ""
526 health.Closed = false
527 }
528
529 func (gw *BotGateway) markAdapterStartFailed(binding AdapterBinding, err error) {
530 gw.mu.Lock()
531 defer gw.mu.Unlock()
532 health := gw.ensureAdapterHealthLocked(binding)
533 health.Status = "error"
534 health.Closed = true
535 health.LastErrorAt = time.Now()
536 if err != nil {
537 health.LastError = err.Error()
538 }
539 }
540
541 func (gw *BotGateway) markAdapterMessage(binding AdapterBinding) {
542 now := time.Now()
543 gw.mu.Lock()
544 defer gw.mu.Unlock()
545 health := gw.ensureAdapterHealthLocked(binding)
546 health.Status = "running"
547 health.LastMessageAt = now
548 health.Messages++
549 health.Closed = false
550 }
551
552 func (gw *BotGateway) markAdapterClosed(binding AdapterBinding) {
553 gw.mu.Lock()
554 defer gw.mu.Unlock()
555 health := gw.ensureAdapterHealthLocked(binding)
556 if health.Status == "running" {
557 health.Status = "closed"
558 }
559 health.Closed = true
560 }
561
562 func (gw *BotGateway) markAdapterSend(binding AdapterBinding, err error) {
563 now := time.Now()
564 gw.mu.Lock()
565 defer gw.mu.Unlock()
566 health := gw.ensureAdapterHealthLocked(binding)
567 if err != nil {
568 health.SendErrors++
569 health.LastErrorAt = now
570 health.LastError = err.Error()
571 if health.Status == "running" {
572 health.Status = "degraded"
573 }
574 return
575 }
576 health.Sends++
577 health.LastSendAt = now
578 if health.Status == "degraded" {
579 health.Status = "running"
580 }
581 }
582
583 func (gw *BotGateway) ensureAdapterHealthLocked(binding AdapterBinding) *AdapterHealthSnapshot {
584 id := strings.TrimSpace(binding.ID)
585 if id == "" && binding.Adapter != nil {
586 id = binding.Adapter.Name()
587 }
588 if id == "" {
589 id = string(binding.Platform)
590 }
591 health := gw.adapterHealth[id]
592 if health == nil {
593 health = &AdapterHealthSnapshot{ID: id}
594 gw.adapterHealth[id] = health
595 }
596 health.Platform = binding.Platform
597 health.Domain = strings.TrimSpace(binding.Domain)
598 if binding.Adapter != nil {
599 health.Name = binding.Adapter.Name()
600 }
601 if strings.TrimSpace(health.Status) == "" {
602 health.Status = "configured"
603 }
604 return health
605 }
606
607 // Stop 停止所有适配器并关闭所有 session。它会等待 dispatch 与 turn goroutine
608 // 全部退出,所以绝不能在 GatewayConfig 回调里同步调用(见 OnInbound 的
609 // reentrancy contract),否则 Stop 会等待正在运行该回调的 goroutine 自己。
610 func (gw *BotGateway) Stop() {
611 gw.lifecycleMu.Lock()
612 if gw.stopped {
613 stopDone := gw.stopDone
614 gw.lifecycleMu.Unlock()
615 if stopDone != nil {
616 <-stopDone
617 }
618 return
619 }
620 gw.stopped = true
621 stopDone := make(chan struct{})
622 gw.stopDone = stopDone
623 cancel := gw.runCancel
624 gw.runCancel = nil
625 startDone := gw.startDone
626 gw.lifecycleMu.Unlock()
627 defer close(stopDone)
628
629 if cancel != nil {
630 cancel()
631 }
632 if startDone != nil {
633 <-startDone
634 }
635
636 // Cancel sessions that already exist before waiting for dispatch to drain.
637 // A dispatch already inside handleMessage may still publish a late session,
638 // so closeSessions is repeated after gatewayWG and turnWG reach zero.
639 gw.closeSessions()
640 for _, binding := range gw.adapters {
641 if err := binding.Adapter.Stop(); err != nil {
642 gw.logger.Warn("error stopping adapter", "platform", binding.Platform, "connection", binding.ID, "err", err)
643 }
644 gw.markAdapterClosed(binding)
645 }
646 gw.stopControlServer()
647 gw.gatewayWG.Wait()
648 gw.closeSessions()
649 gw.turnWG.Wait()
650 gw.finishSessionTeardown()
651 }
652
653 func (gw *BotGateway) closeSessions() {
654 var states []*sessionState
655 gw.mu.Lock()
656 for key, state := range gw.controllers {
657 states = append(states, state)
658 delete(gw.controllers, key)
659 }
660 gw.mu.Unlock()
661 for _, state := range states {
662 gw.closeSessionState(state)
663 }
664 }
665
666 // closeSessionState tears down a session state that has been unlinked from
667 // gw.controllers. runTurn publishes state.cancel under gw.mu on every turn —
668 // possibly after the state was already unlinked — so snapshot and clear the
669 // field inside the lock and invoke it outside (the same discipline as
670 // cancelActiveSession).
671 func (gw *BotGateway) closeSessionState(state *sessionState) {
672 if state == nil {
673 return
674 }
675 // Serialize retirement with recovery ownership handoffs. Stop unlinks
676 // sessions before turn goroutines drain, so a recovery callback captured by
677 // the controller can still arrive here. Marking the state retired under the
678 // same lock prevents that callback from reacquiring a lease after teardown;
679 // an already-running handoff completes before the lease is released below.
680 state.lifecycleMu.Lock()
681 if state.retired {
682 state.lifecycleMu.Unlock()
683 return
684 }
685 state.retired = true
686 state.lifecycleMu.Unlock()
687
688 gw.mu.Lock()
689 cancel := state.cancel
690 state.cancel = nil
691 gw.mu.Unlock()
692 if cancel != nil {
693 cancel()
694 }
695 if state.ctrl != nil {
696 if state.releaseRuntimeOnly {
697 if releaser, ok := state.ctrl.(interface{ ReleaseResources() }); ok {
698 releaser.ReleaseResources()
699 } else {
700 state.ctrl.Close()
701 }
702 } else {
703 state.ctrl.Close()
704 }
705 }
706 if state.leases != nil {
707 state.leases.Release()
708 }
709 }
710
711 // unlinkAndCloseSessionState removes state from the live gateway before closing
712 // it. It is used when a controller has already rotated its transcript but the
713 // replacement lease could not be acquired: retaining that state would let the
714 // next message reuse a controller that no longer owns its active session path.
715 func (gw *BotGateway) unlinkAndCloseSessionState(key string, state *sessionState) {
716 if state == nil {
717 return
718 }
719 gw.mu.Lock()
720 if gw.controllers[key] == state {
721 delete(gw.controllers, key)
722 }
723 gw.mu.Unlock()
724 gw.closeSessionState(state)
725 }
726
727 func (gw *BotGateway) dispatchLoop(ctx context.Context, binding AdapterBinding) {
728 for {
729 select {
730 case <-ctx.Done():
731 gw.markAdapterClosed(binding)
732 return
733 case msg, ok := <-binding.Adapter.Messages():
734 if !ok {
735 gw.markAdapterClosed(binding)
736 return
737 }
738 gw.markAdapterMessage(binding)
739 gw.handleMessage(ctx, binding, msg)
740 }
741 }
742 }
743
744 func (gw *BotGateway) handleMessage(ctx context.Context, binding AdapterBinding, msg InboundMessage) {
745 msg.Platform = binding.Platform
746 if msg.ConnectionID == "" {
747 msg.ConnectionID = binding.ID
748 }
749 if msg.Domain == "" {
750 msg.Domain = binding.Domain
751 }
752 if gw.isSelfMessage(msg) {
753 gw.logger.Debug("bot ignored self message", "platform", binding.Platform, "connection", msg.ConnectionID, "chat", hashID(msg.ChatID), "message", hashID(msg.MessageID), "user", hashID(msg.UserID))
754 return
755 }
756 src := msg.Session()
757 key := BuildSessionKey(src)
758 logFields := []any{
759 "platform", binding.Platform,
760 "connection", msg.ConnectionID,
761 "domain", msg.Domain,
762 "chat_type", msg.ChatType,
763 "chat", hashID(msg.ChatID),
764 "user", hashID(msg.UserID),
765 "operator", hashID(msg.OperatorID),
766 "thread", hashID(msg.ThreadID),
767 "message", hashID(msg.MessageID),
768 "text_chars", len([]rune(msg.Text)),
769 "session", key[:8],
770 }
771 gw.logger.Info("bot inbound message", logFields...)
772
773 // allowlist 检查
774 if !gw.checkAllowlist(binding.Platform, msg) {
775 gw.logger.Info("user not in allowlist", "platform", binding.Platform, "connection", msg.ConnectionID, "user", hashID(msg.UserID))
776 if gw.offerPairing(ctx, binding.Adapter, msg) {
777 return
778 }
779 _ = gw.sendText(ctx, binding.Adapter, msg, "抱歉,您没有使用此 bot 的权限。")
780 return
781 }
782 if gw.cfg.OnInbound != nil {
783 gw.cfg.OnInbound(msg)
784 }
785
786 if normalized, ok := gw.normalizeApprovalShortcut(key, msg.Text); ok {
787 msg.Text = normalized
788 } else if normalized, ok := gw.normalizeAskShortcut(key, msg.Text); ok {
789 msg.Text = normalized
790 } else if _, ok := decisionShortcutCommand(msg.Text); ok && gw.sessions.IsActive(key) {
791 _ = gw.sendText(ctx, binding.Adapter, msg, "没有找到可匹配的待处理操作。请重新触发一次操作后回复编号,或按消息中的 ID 使用 /approve、/deny 或 /answer。")
792 return
793 }
794
795 // 斜杠命令处理
796 if IsSlashBypass(msg.Text) {
797 gw.logger.Info("bot slash command", logFields...)
798 gw.handleSlashCommand(ctx, binding.Adapter, key, msg)
799 return
800 }
801
802 // 已接管桌面会话的聊天:普通消息直接驱动那个桌面会话,不进 bot 自己的
803 // 会话机器(斜杠命令仍走上面的分支,/desktop release 永远可达)。
804 if gw.divertToDesktopTakeover(ctx, binding.Adapter, msg) {
805 gw.logger.Info("bot message diverted to desktop takeover", logFields...)
806 return
807 }
808
809 cleanup := gw.addPendingReaction(ctx, binding.Platform, binding.Adapter, msg)
810
811 queueMode := gw.queueMode(key, msg)
812 warnDeprecatedQueueDrop(gw.cfg.QueueDrop)
813 if gw.sessions.IsActive(key) {
814 // Busy session: durable inbox is the authority (not SessionManager.pending).
815 if IsSlashBypass(msg.Text) {
816 // Slash commands still acquire through the session lock below.
817 } else {
818 switch queueMode {
819 case QueueModeSteer:
820 if rec, ok := gw.steerActiveSessionDurable(ctx, binding.Adapter, key, msg); ok {
821 gw.logger.Info("bot message steered into active turn", "session", key[:8], "item", rec.ItemID)
822 if cleanup != nil {
823 cleanup()
824 }
825 _ = gw.sendText(ctx, binding.Adapter, msg, formatQueuedReceipt(rec)+"(已并入当前任务)")
826 return
827 }
828 case QueueModeInterrupt:
829 gw.cancelActiveSession(key)
830 runReactionCleanups(gw.takeReactionCleanups(key))
831 rec, err := gw.interruptActiveSessionDurable(ctx, binding.Adapter, key, msg)
832 gw.storeReactionCleanup(key, cleanup)
833 if err != nil {
834 gw.logger.Warn("bot interrupt enqueue failed", "session", key[:8], "err", err)
835 _ = gw.sendText(ctx, binding.Adapter, msg, "排队失败:"+err.Error())
836 return
837 }
838 gw.logger.Info("bot active turn interrupted; newest message durable-queued", "session", key[:8], "item", rec.ItemID)
839 _ = gw.sendText(ctx, binding.Adapter, msg, "已停止当前任务。"+formatQueuedReceipt(rec))
840 return
841 case QueueModeCollect:
842 if rec, err := gw.collectActiveSessionDurable(ctx, binding.Adapter, key, msg); err == nil {
843 gw.storeReactionCleanup(key, cleanup)
844 _ = gw.sendText(ctx, binding.Adapter, msg, formatQueuedReceipt(rec))
845 return
846 } else if errors.Is(err, sessioninbox.ErrCapacityItems) || errors.Is(err, sessioninbox.ErrCapacityBytes) || errors.Is(err, sessioninbox.ErrItemTooLarge) {
847 if cleanup != nil {
848 cleanup()
849 }
850 _ = gw.sendText(ctx, binding.Adapter, msg, "当前会话排队已满,请稍后再发,或使用 /queue pause 后清理。")
851 return
852 }
853 default: // followup
854 if rec, err := gw.followupActiveSessionDurable(ctx, binding.Adapter, key, msg); err == nil {
855 gw.storeReactionCleanup(key, cleanup)
856 _ = gw.sendText(ctx, binding.Adapter, msg, formatQueuedReceipt(rec))
857 return
858 } else if errors.Is(err, sessioninbox.ErrCapacityItems) || errors.Is(err, sessioninbox.ErrCapacityBytes) || errors.Is(err, sessioninbox.ErrItemTooLarge) {
859 if cleanup != nil {
860 cleanup()
861 }
862 _ = gw.sendText(ctx, binding.Adapter, msg, "当前会话排队已满,请稍后再发。")
863 return
864 }
865 }
866 }
867 }
868
869 // session 并发控制 — only the active-turn lock remains here; bodies live in inbox.
870 result := gw.sessions.TryAcquireWithQueue(key, msg, QueueOptions{
871 Mode: QueueModeFollowup, // never drop_old; capacity enforced by inbox
872 Cap: sessioninbox.DefaultMaxItems,
873 Drop: QueueDropNew,
874 })
875 if result.Rejected {
876 gw.logger.Warn("bot queue rejected message", "session", key[:8], "pending", result.Pending, "mode", result.Mode)
877 if cleanup != nil {
878 cleanup()
879 }
880 _ = gw.sendText(ctx, binding.Adapter, msg, "当前会话排队已满,请稍后再发,或使用 /queue 管理队列。")
881 return
882 }
883 gw.dispatchQueueResult(ctx, binding.Adapter, key, msg, cleanup, result)
884 }
885
886 func (gw *BotGateway) queueMode(key string, msg InboundMessage) string {
887 return gw.sessions.QueueMode(key, gw.cfg.QueueMode)
888 }
889
890 func (gw *BotGateway) sessionAPI(key string) control.SessionAPI {
891 gw.mu.Lock()
892 state, ok := gw.controllers[key]
893 gw.mu.Unlock()
894 if !ok || state == nil || state.ctrl == nil {
895 return nil
896 }
897 if api, ok := state.ctrl.(control.SessionAPI); ok {
898 return api
899 }
900 return nil
901 }
902
903 func (gw *BotGateway) steerActiveSessionDurable(ctx context.Context, adapter Adapter, key string, msg InboundMessage) (sessioninbox.InboxReceipt, bool) {
904 text := strings.TrimSpace(msg.Text)
905 if text == "" && len(msg.MediaURLs) == 0 && len(msg.Media) == 0 {
906 return sessioninbox.InboxReceipt{}, false
907 }
908 gw.mu.Lock()
909 state, ok := gw.controllers[key]
910 gw.mu.Unlock()
911 if !ok || state.ctrl == nil {
912 return sessioninbox.InboxReceipt{}, false
913 }
914 msg = gw.prepareDurableInboxMessage(ctx, adapter, msg, state)
915 text = msg.Text
916 if strings.TrimSpace(text) == "" {
917 return sessioninbox.InboxReceipt{}, false
918 }
919 msg.Text = text
920 api, ok := state.ctrl.(control.SessionAPI)
921 if !ok {
922 // Legacy fallback.
923 if steerer, ok := state.ctrl.(interface{ TrySteer(string) bool }); ok && steerer.TrySteer(text) {
924 return sessioninbox.InboxReceipt{Disposition: sessioninbox.DispositionSteerAccepted}, true
925 }
926 return sessioninbox.InboxReceipt{}, false
927 }
928 rec, err := enqueueViaInbox(api, msg, sessioninbox.IntentSteer)
929 if err != nil {
930 return sessioninbox.InboxReceipt{}, false
931 }
932 return rec, true
933 }
934
935 func (gw *BotGateway) cancelActiveSession(key string) {
936 // state.cancel is rewritten under gw.mu on every turn (runTurn), so copy it
937 // inside the lock and invoke it outside.
938 var cancel context.CancelFunc
939 gw.mu.Lock()
940 state, ok := gw.controllers[key]
941 if ok && state != nil {
942 cancel = state.cancel
943 }
944 gw.mu.Unlock()
945 if !ok || state == nil {
946 return
947 }
948 if cancel != nil {
949 cancel()
950 return
951 }
952 if state.ctrl != nil {
953 state.ctrl.Cancel()
954 }
955 }
956
957 func (gw *BotGateway) storeReactionCleanup(key string, cleanup func()) {
958 if cleanup == nil {
959 return
960 }
961 gw.mu.Lock()
962 defer gw.mu.Unlock()
963 gw.pendingReactionCleanups[key] = append(gw.pendingReactionCleanups[key], cleanup)
964 }
965
966 func (gw *BotGateway) flushReactionCleanups(key string, cleanup func()) {
967 stored := gw.takeReactionCleanups(key)
968 runReactionCleanups(stored)
969 if cleanup != nil {
970 cleanup()
971 }
972 }
973
974 func (gw *BotGateway) takeReactionCleanups(key string) []func() {
975 gw.mu.Lock()
976 defer gw.mu.Unlock()
977 stored := gw.pendingReactionCleanups[key]
978 delete(gw.pendingReactionCleanups, key)
979 return stored
980 }
981
982 func runReactionCleanups(cleanups []func()) {
983 for _, cleanup := range cleanups {
984 if cleanup != nil {
985 cleanup()
986 }
987 }
988 }
989
990 func makeReactionCleanup(cleanups []func()) func() {
991 if len(cleanups) == 0 {
992 return nil
993 }
994 return func() {
995 runReactionCleanups(cleanups)
996 }
997 }
998
999 func (gw *BotGateway) addPendingReaction(ctx context.Context, plat Platform, adapter Adapter, msg InboundMessage) func() {
1000 if strings.TrimSpace(msg.MessageID) == "" {
1001 return nil
1002 }
1003 reactor, ok := adapter.(pendingReactionAdapter)
1004 if !ok {
1005 return nil
1006 }
1007 cleanup, err := reactor.AddPendingReaction(ctx, msg.MessageID)
1008 if err != nil {
1009 gw.logger.Warn("pending reaction failed", "platform", plat, "err", err)
1010 return nil
1011 }
1012 return cleanup
1013 }
1014
1015 func (gw *BotGateway) isSelfMessage(msg InboundMessage) bool {
1016 if !gw.cfg.IgnoreSelfMessages {
1017 return false
1018 }
1019 actor := strings.TrimSpace(msg.UserID)
1020 if strings.TrimSpace(msg.OperatorID) != "" {
1021 actor = strings.TrimSpace(msg.OperatorID)
1022 }
1023 if actor != "" && gw.selfUserIDs[msg.Platform][actor] {
1024 return true
1025 }
1026 messageID := strings.TrimSpace(msg.MessageID)
1027 if messageID == "" {
1028 return false
1029 }
1030 key := outboundMessageKey(msg.Platform, msg.ConnectionID, msg.Domain, msg.ChatID, messageID)
1031 now := time.Now()
1032 gw.mu.Lock()
1033 defer gw.mu.Unlock()
1034 gw.pruneOutboundMessagesLocked(now)
1035 _, ok := gw.outboundMessageIDs[key]
1036 return ok
1037 }
1038
1039 func (gw *BotGateway) rememberOutboundMessage(platform Platform, connID, domain, chatID, messageID string) {
1040 messageID = strings.TrimSpace(messageID)
1041 if !gw.cfg.IgnoreSelfMessages || messageID == "" {
1042 return
1043 }
1044 now := time.Now()
1045 key := outboundMessageKey(platform, connID, domain, chatID, messageID)
1046 gw.mu.Lock()
1047 defer gw.mu.Unlock()
1048 gw.pruneOutboundMessagesLocked(now)
1049 gw.outboundMessageIDs[key] = now.Add(outboundEchoTTL)
1050 }
1051
1052 func (gw *BotGateway) pruneOutboundMessagesLocked(now time.Time) {
1053 for key, expiresAt := range gw.outboundMessageIDs {
1054 if !expiresAt.After(now) {
1055 delete(gw.outboundMessageIDs, key)
1056 }
1057 }
1058 }
1059
1060 func outboundMessageKey(platform Platform, connID, domain, chatID, messageID string) string {
1061 return strings.Join([]string{
1062 string(platform),
1063 strings.TrimSpace(connID),
1064 strings.TrimSpace(domain),
1065 strings.TrimSpace(chatID),
1066 strings.TrimSpace(messageID),
1067 }, "\x00")
1068 }
1069
1070 func (gw *BotGateway) connectionAccess(msg InboundMessage) (AccessConfig, bool) {
1071 if gw.cfg.ConnectionAccess == nil {
1072 return AccessConfig{}, false
1073 }
1074 id := strings.TrimSpace(msg.ConnectionID)
1075 if id == "" {
1076 return AccessConfig{}, false
1077 }
1078 access, ok := gw.cfg.ConnectionAccess[id]
1079 if !ok {
1080 return AccessConfig{}, false
1081 }
1082 if !accessConfigActive(access) {
1083 return AccessConfig{}, false
1084 }
1085 return access, true
1086 }
1087
1088 func accessConfigActive(access AccessConfig) bool {
1089 return access.Enabled ||
1090 access.AllowAll ||
1091 access.PairingEnabled ||
1092 len(access.Users) > 0 ||
1093 len(access.Groups) > 0 ||
1094 len(access.Approvers) > 0 ||
1095 len(access.Admins) > 0
1096 }
1097
1098 func (gw *BotGateway) checkAllowlist(plat Platform, msg InboundMessage) bool {
1099 if access, ok := gw.connectionAccess(msg); ok {
1100 return checkConnectionAllowlist(access, msg)
1101 }
1102 if gw.cfg.Allowlist.AllowAll {
1103 return true
1104 }
1105 if !gw.cfg.Allowlist.Enabled {
1106 return false
1107 }
1108 actor := msg.UserID
1109 if msg.OperatorID != "" {
1110 actor = msg.OperatorID
1111 }
1112 if !gw.allowlist[plat][actor] {
1113 return false
1114 }
1115 groups := gw.groupAllowlist[plat]
1116 if chatUsesGroupAllowlist(msg.ChatType) && len(groups) > 0 && !groups[msg.ChatID] {
1117 return false
1118 }
1119 return true
1120 }
1121
1122 func checkConnectionAllowlist(access AccessConfig, msg InboundMessage) bool {
1123 if access.AllowAll {
1124 return true
1125 }
1126 if !access.Enabled {
1127 return false
1128 }
1129 actor := msg.UserID
1130 if msg.OperatorID != "" {
1131 actor = msg.OperatorID
1132 }
1133 users := stringSet(append(append(append([]string{}, access.Users...), access.Admins...), access.Approvers...))
1134 groups := stringSet(access.Groups)
1135 actorAllowed := users[actor]
1136 groupAllowed := chatUsesGroupAllowlist(msg.ChatType) && groups[msg.ChatID]
1137 if len(users) == 0 && len(groups) == 0 {
1138 return false
1139 }
1140 return actorAllowed || groupAllowed
1141 }
1142
1143 func (gw *BotGateway) requireCommandRole(ctx context.Context, adapter Adapter, msg InboundMessage, role string) bool {
1144 if gw.checkCommandRole(msg.Platform, msg, role) {
1145 return true
1146 }
1147 _ = gw.sendText(ctx, adapter, msg, "抱歉,你没有执行此 bot 命令的权限。")
1148 return false
1149 }
1150
1151 func (gw *BotGateway) checkCommandRole(plat Platform, msg InboundMessage, role string) bool {
1152 actor := msg.UserID
1153 if msg.OperatorID != "" {
1154 actor = msg.OperatorID
1155 }
1156 if strings.TrimSpace(actor) == "" {
1157 return false
1158 }
1159 if access, ok := gw.connectionAccess(msg); ok {
1160 admins := stringSet(access.Admins)
1161 approvers := stringSet(access.Approvers)
1162 if len(admins) == 0 && len(approvers) == 0 {
1163 return true
1164 }
1165 if admins[actor] {
1166 return true
1167 }
1168 return role == "approver" && approvers[actor]
1169 }
1170 admins := stringSet(gw.cfg.Allowlist.Admins[plat])
1171 approvers := stringSet(gw.cfg.Allowlist.Approvers[plat])
1172 if len(admins) == 0 && len(approvers) == 0 {
1173 return true
1174 }
1175 if admins[actor] {
1176 return true
1177 }
1178 if role == "approver" && approvers[actor] {
1179 return true
1180 }
1181 return false
1182 }
1183
1184 func stringSet(values []string) map[string]bool {
1185 out := make(map[string]bool, len(values))
1186 for _, value := range values {
1187 value = strings.TrimSpace(value)
1188 if value != "" {
1189 out[value] = true
1190 }
1191 }
1192 return out
1193 }
1194
1195 func (gw *BotGateway) offerPairing(ctx context.Context, adapter Adapter, msg InboundMessage) bool {
1196 if access, ok := gw.connectionAccess(msg); ok {
1197 if !access.PairingEnabled {
1198 return false
1199 }
1200 } else if !gw.cfg.PairingEnabled {
1201 return false
1202 }
1203 req, created, err := CreateOrRefreshPairingRequest(msg, PairingConfig{
1204 Enabled: true,
1205 RequestTTL: gw.cfg.PairingTTL,
1206 MaxPendingPerPlatform: gw.cfg.PairingMaxPending,
1207 })
1208 if err != nil {
1209 gw.logger.Warn("bot pairing request failed", "platform", msg.Platform, "chat_type", msg.ChatType, "err", err)
1210 return false
1211 }
1212 prefix := "需要先完成配对。"
1213 if !created {
1214 prefix = "你已有待批准的配对请求。"
1215 }
1216 text := fmt.Sprintf("%s\n配对码: %s\n请在本机运行: reasonix bot pairing approve %s\n此码将在 %s 过期。",
1217 prefix, req.Code, req.Code, req.ExpiresAt.Local().Format("2006-01-02 15:04"))
1218 _ = gw.sendText(ctx, adapter, msg, text)
1219 return true
1220 }
1221
1222 func chatUsesGroupAllowlist(chatType ChatType) bool {
1223 switch chatType {
1224 case ChatGroup, ChatGuild, ChatThread:
1225 return true
1226 default:
1227 return false
1228 }
1229 }
1230
1231 func approvalShortcutCommand(text string) (string, bool) {
1232 switch strings.ToLower(strings.TrimSpace(text)) {
1233 case "1", "y", "yes", "ok", "同意", "批准", "允许", "允许一次":
1234 return "/approve", true
1235 case "2", "0", "n", "no", "deny", "拒绝":
1236 return "/deny", true
1237 default:
1238 return "", false
1239 }
1240 }
1241
1242 func recoveryShortcutCommand(text string, canGrantTask bool) (string, bool) {
1243 switch strings.ToLower(strings.TrimSpace(text)) {
1244 case "1", "y", "yes", "ok", "继续", "继续此变更", "continue":
1245 return "/recovery-continue", true
1246 case "2", "a", "同类", "本任务允许", "allow similar":
1247 if canGrantTask {
1248 return "/recovery-continue-task", true
1249 }
1250 return "/recovery-revise", true
1251 case "3":
1252 if canGrantTask {
1253 return "/recovery-revise", true
1254 }
1255 return "", false
1256 case "修改", "修改方案", "换个办法", "revise":
1257 return "/recovery-revise", true
1258 default:
1259 return "", false
1260 }
1261 }
1262
1263 func (gw *BotGateway) pendingRecoveryCanGrantTask(key, id string) bool {
1264 gw.mu.Lock()
1265 defer gw.mu.Unlock()
1266 state, ok := gw.controllers[key]
1267 if !ok || state.pendingApprovals == nil {
1268 return false
1269 }
1270 a, ok := state.pendingApprovals[id]
1271 return ok && a.Recovery != nil && a.Recovery.CanGrantTask
1272 }
1273
1274 func (gw *BotGateway) pendingApprovalIsRecovery(key, id string) bool {
1275 gw.mu.Lock()
1276 defer gw.mu.Unlock()
1277 state, ok := gw.controllers[key]
1278 if !ok || state.pendingApprovals == nil {
1279 return false
1280 }
1281 a, ok := state.pendingApprovals[id]
1282 if !ok {
1283 return false
1284 }
1285 return strings.EqualFold(strings.TrimSpace(a.Kind), "recovery") || a.Recovery != nil
1286 }
1287
1288 func decisionShortcutCommand(text string) (string, bool) {
1289 if command, ok := approvalShortcutCommand(text); ok {
1290 return command, true
1291 }
1292 if _, ok := askShortcutAnswer(text); ok {
1293 return "/answer", true
1294 }
1295 return "", false
1296 }
1297
1298 func (gw *BotGateway) currentPendingApprovalID(key string) string {
1299 gw.mu.Lock()
1300 defer gw.mu.Unlock()
1301 state, ok := gw.controllers[key]
1302 if !ok || len(state.pendingApprovals) == 0 {
1303 return ""
1304 }
1305 if state.lastApprovalID != "" {
1306 if _, ok := state.pendingApprovals[state.lastApprovalID]; ok {
1307 return state.lastApprovalID
1308 }
1309 }
1310 for id := range state.pendingApprovals {
1311 return id
1312 }
1313 return ""
1314 }
1315
1316 func (gw *BotGateway) forgetPendingApproval(key, id string) {
1317 gw.mu.Lock()
1318 defer gw.mu.Unlock()
1319 state, ok := gw.controllers[key]
1320 if !ok || state.pendingApprovals == nil {
1321 return
1322 }
1323 delete(state.pendingApprovals, id)
1324 if state.lastApprovalID == id {
1325 state.lastApprovalID = ""
1326 for nextID := range state.pendingApprovals {
1327 state.lastApprovalID = nextID
1328 break
1329 }
1330 }
1331 }
1332
1333 func (gw *BotGateway) normalizeAskShortcut(key, text string) (string, bool) {
1334 raw := strings.TrimSpace(text)
1335 if raw == "" || strings.HasPrefix(raw, "/") {
1336 return "", false
1337 }
1338 askID := gw.currentPendingAskIDForReply(key)
1339 if askID == "" {
1340 return "", false
1341 }
1342 return "/answer " + askID + " " + raw, true
1343 }
1344
1345 func askShortcutAnswer(text string) (string, bool) {
1346 raw := strings.TrimSpace(text)
1347 if raw == "" {
1348 return "", false
1349 }
1350 if strings.ContainsAny(raw, " \t\n;=") {
1351 return "", false
1352 }
1353 if _, err := strconv.Atoi(raw); err == nil {
1354 return raw, true
1355 }
1356 return "", false
1357 }
1358
1359 func (gw *BotGateway) currentPendingAskIDForReply(key string) string {
1360 gw.mu.Lock()
1361 defer gw.mu.Unlock()
1362 state, ok := gw.controllers[key]
1363 if !ok || len(state.pendingAsks) == 0 {
1364 return ""
1365 }
1366 if state.lastAskID != "" {
1367 if _, ok := state.pendingAsks[state.lastAskID]; ok {
1368 return state.lastAskID
1369 }
1370 }
1371 if len(state.pendingAsks) != 1 {
1372 return ""
1373 }
1374 for id := range state.pendingAsks {
1375 return id
1376 }
1377 return ""
1378 }
1379
1380 func (gw *BotGateway) handleSlashCommandCore(ctx context.Context, adapter Adapter, key string, msg InboundMessage) {
1381 switch {
1382 case strings.HasPrefix(msg.Text, "/stop"):
1383 var cancel context.CancelFunc
1384 gw.mu.Lock()
1385 if state, ok := gw.controllers[key]; ok {
1386 cancel = state.cancel
1387 }
1388 gw.mu.Unlock()
1389 if cancel != nil {
1390 cancel()
1391 }
1392 gw.sessions.ForceRelease(key)
1393 _ = gw.sendText(ctx, adapter, msg, "已停止当前任务。")
1394
1395 case strings.HasPrefix(msg.Text, "/new") || strings.HasPrefix(msg.Text, "/reset"):
1396 var cancel context.CancelFunc
1397 gw.mu.Lock()
1398 state, ok := gw.controllers[key]
1399 if ok {
1400 cancel = state.cancel
1401 }
1402 gw.mu.Unlock()
1403 if ok {
1404 if cancel != nil {
1405 cancel()
1406 }
1407 // NewSession refuses to rotate while a turn is running; the cancel
1408 // above is asynchronous, so give the turn a bounded window to
1409 // unwind before rotating.
1410 deadline := time.Now().Add(5 * time.Second)
1411 for state.ctrl.Running() && time.Now().Before(deadline) {
1412 time.Sleep(10 * time.Millisecond)
1413 }
1414 if err := state.ctrl.NewSession(); err != nil {
1415 gw.logger.Warn("new session failed", "err", err)
1416 gw.sessions.ForceRelease(key)
1417 _ = gw.sendText(ctx, adapter, msg, "新会话创建失败,请稍后重试。")
1418 return
1419 }
1420 if state.leases != nil && strings.TrimSpace(state.ctrl.SessionPath()) != "" {
1421 if err := rebindBotSessionWriteAuthority(state, state.ctrl.SessionPath()); err != nil {
1422 gw.logger.Warn("new session lease failed", "err", control.SessionInUseMessage(err))
1423 gw.unlinkAndCloseSessionState(key, state)
1424 gw.sessions.ForceRelease(key)
1425 _ = gw.sendText(ctx, adapter, msg, "新会话创建失败:无法取得写入权限。请关闭其他 Reasonix 窗口或进程后重试。")
1426 return
1427 }
1428 }
1429 // /new 后把旋转出的不可变身份钉为会话覆盖,避免下一条消息重新解析回旧绑定。
1430 gw.mu.Lock()
1431 if gw.controllers[key] == state {
1432 rotated := state.ctrl.SessionPath()
1433 override, exists := gw.sessionOverrides[key]
1434 if !exists {
1435 override = sessionRuntimeOverride{}
1436 }
1437 if identity, ok := state.ctrl.(control.IdentityLifecycle); ok && identity.UsesExclusiveSession() {
1438 if ref, bound := identity.SessionRef(); bound {
1439 state.sessionRef = ref
1440 state.sessionPath = ""
1441 override.sessionPath = botSessionRefTarget(ref)
1442 }
1443 } else {
1444 state.sessionPath = rotated
1445 override.sessionPath = rotated
1446 }
1447 gw.sessionOverrides[key] = override
1448 }
1449 gw.mu.Unlock()
1450 gw.rememberSessionReady(msg, state.ctrl)
1451 }
1452 gw.sessions.ForceRelease(key)
1453 _ = gw.sendText(ctx, adapter, msg, "已开始新会话。")
1454
1455 case strings.HasPrefix(msg.Text, "/approve"):
1456 if !gw.requireCommandRole(ctx, adapter, msg, "approver") {
1457 return
1458 }
1459 // 从消息中解析 approval ID
1460 parts := strings.Fields(msg.Text)
1461 if len(parts) < 2 {
1462 _ = gw.sendText(ctx, adapter, msg, "用法: /approve <id>")
1463 return
1464 }
1465 gw.mu.Lock()
1466 state, ok := gw.controllers[key]
1467 gw.mu.Unlock()
1468 if ok && state.ctrl != nil {
1469 // Recovery cards map allow → continue for older clients that only know Approve.
1470 if gw.pendingApprovalIsRecovery(key, parts[1]) {
1471 _ = state.ctrl.ResolveRecovery(parts[1], agent.RecoveryActionContinue, "")
1472 } else {
1473 state.ctrl.Approve(parts[1], true, false, false)
1474 }
1475 gw.forgetPendingApproval(key, parts[1])
1476 _ = gw.sendText(ctx, adapter, msg, "已批准。")
1477 } else {
1478 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话中的待审批操作,请重新触发一次操作。")
1479 }
1480
1481 case strings.HasPrefix(msg.Text, "/deny"):
1482 if !gw.requireCommandRole(ctx, adapter, msg, "approver") {
1483 return
1484 }
1485 parts := strings.Fields(msg.Text)
1486 if len(parts) < 2 {
1487 _ = gw.sendText(ctx, adapter, msg, "用法: /deny <id>")
1488 return
1489 }
1490 gw.mu.Lock()
1491 state, ok := gw.controllers[key]
1492 gw.mu.Unlock()
1493 if ok && state.ctrl != nil {
1494 if gw.pendingApprovalIsRecovery(key, parts[1]) {
1495 _ = state.ctrl.ResolveRecovery(parts[1], agent.RecoveryActionRevise, "")
1496 } else {
1497 state.ctrl.Approve(parts[1], false, false, false)
1498 }
1499 gw.forgetPendingApproval(key, parts[1])
1500 _ = gw.sendText(ctx, adapter, msg, "已拒绝。")
1501 } else {
1502 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话中的待审批操作,请重新触发一次操作。")
1503 }
1504
1505 case strings.HasPrefix(msg.Text, "/recovery-continue-task"):
1506 if !gw.requireCommandRole(ctx, adapter, msg, "approver") {
1507 return
1508 }
1509 parts := strings.Fields(msg.Text)
1510 if len(parts) < 2 {
1511 _ = gw.sendText(ctx, adapter, msg, "用法: /recovery-continue-task <id>")
1512 return
1513 }
1514 gw.mu.Lock()
1515 state, ok := gw.controllers[key]
1516 gw.mu.Unlock()
1517 if ok && state.ctrl != nil {
1518 if err := state.ctrl.ResolveRecovery(parts[1], agent.RecoveryActionContinueTask, ""); err != nil {
1519 _ = gw.sendText(ctx, adapter, msg, "确认失败: "+err.Error())
1520 return
1521 }
1522 gw.forgetPendingApproval(key, parts[1])
1523 _ = gw.sendText(ctx, adapter, msg, "已继续;本任务内同类操作将自动执行,范围扩大或风险升级仍会确认。")
1524 } else {
1525 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话中的待确认操作。")
1526 }
1527
1528 case strings.HasPrefix(msg.Text, "/recovery-continue"):
1529 if !gw.requireCommandRole(ctx, adapter, msg, "approver") {
1530 return
1531 }
1532 parts := strings.Fields(msg.Text)
1533 if len(parts) < 2 {
1534 _ = gw.sendText(ctx, adapter, msg, "用法: /recovery-continue <id>")
1535 return
1536 }
1537 gw.mu.Lock()
1538 state, ok := gw.controllers[key]
1539 gw.mu.Unlock()
1540 if ok && state.ctrl != nil {
1541 if err := state.ctrl.ResolveRecovery(parts[1], agent.RecoveryActionContinue, ""); err != nil {
1542 _ = gw.sendText(ctx, adapter, msg, "确认失败: "+err.Error())
1543 return
1544 }
1545 gw.forgetPendingApproval(key, parts[1])
1546 _ = gw.sendText(ctx, adapter, msg, "已继续。")
1547 } else {
1548 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话中的待确认操作。")
1549 }
1550
1551 case strings.HasPrefix(msg.Text, "/recovery-revise"):
1552 if !gw.requireCommandRole(ctx, adapter, msg, "approver") {
1553 return
1554 }
1555 parts := strings.Fields(msg.Text)
1556 if len(parts) < 2 {
1557 _ = gw.sendText(ctx, adapter, msg, "用法: /recovery-revise <id> [补充要求]")
1558 return
1559 }
1560 feedback := strings.TrimSpace(strings.Join(parts[2:], " "))
1561 gw.mu.Lock()
1562 state, ok := gw.controllers[key]
1563 gw.mu.Unlock()
1564 if ok && state.ctrl != nil {
1565 if err := state.ctrl.ResolveRecovery(parts[1], agent.RecoveryActionRevise, feedback); err != nil {
1566 _ = gw.sendText(ctx, adapter, msg, "修改方案失败: "+err.Error())
1567 return
1568 }
1569 gw.forgetPendingApproval(key, parts[1])
1570 _ = gw.sendText(ctx, adapter, msg, "已拒绝当前变更并注入修改要求。")
1571 } else {
1572 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话中的恢复检查点。")
1573 }
1574
1575 case strings.HasPrefix(msg.Text, "/recovery-stop"):
1576 // Backward compatibility for cards rendered by an older client: reject
1577 // the proposed mutation but leave task cancellation to ordinary /stop.
1578 if !gw.requireCommandRole(ctx, adapter, msg, "approver") {
1579 return
1580 }
1581 parts := strings.Fields(msg.Text)
1582 if len(parts) < 2 {
1583 _ = gw.sendText(ctx, adapter, msg, "用法: /recovery-stop <id>")
1584 return
1585 }
1586 gw.mu.Lock()
1587 state, ok := gw.controllers[key]
1588 gw.mu.Unlock()
1589 if ok && state.ctrl != nil {
1590 if err := state.ctrl.ResolveRecovery(parts[1], agent.RecoveryActionRevise, "cancel this proposed action"); err != nil {
1591 _ = gw.sendText(ctx, adapter, msg, "取消变更失败: "+err.Error())
1592 return
1593 }
1594 gw.forgetPendingApproval(key, parts[1])
1595 _ = gw.sendText(ctx, adapter, msg, "已取消当前变更;如需停止整个任务,请使用 /stop。")
1596 } else {
1597 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话中的恢复检查点。")
1598 }
1599
1600 case strings.HasPrefix(msg.Text, "/answer"):
1601 parts := strings.Fields(msg.Text)
1602 if len(parts) < 3 {
1603 _ = gw.sendText(ctx, adapter, msg, "用法: /answer <id> <选项或 q1=选项;q2=选项>")
1604 return
1605 }
1606 askID := parts[1]
1607 rawAnswer := strings.TrimSpace(strings.Join(parts[2:], " "))
1608 gw.mu.Lock()
1609 state, ok := gw.controllers[key]
1610 var questions []event.AskQuestion
1611 if ok {
1612 questions = state.pendingAsks[askID]
1613 delete(state.pendingAsks, askID)
1614 if state.lastAskID == askID {
1615 state.lastAskID = ""
1616 for nextID := range state.pendingAsks {
1617 state.lastAskID = nextID
1618 break
1619 }
1620 }
1621 }
1622 gw.mu.Unlock()
1623 if !ok || state.ctrl == nil {
1624 _ = gw.sendText(ctx, adapter, msg, "没有找到当前会话。")
1625 return
1626 }
1627 answers := parseAskAnswers(questions, rawAnswer)
1628 state.ctrl.AnswerQuestion(askID, answers)
1629 _ = gw.sendText(ctx, adapter, msg, "已提交回答。")
1630
1631 case strings.HasPrefix(msg.Text, "/mode"):
1632 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1633 return
1634 }
1635 mode, statusOnly, ok := parseToolApprovalModeCommand(msg.Text)
1636 if !ok {
1637 _ = gw.sendText(ctx, adapter, msg, "用法: /mode read-only|workspace-write|danger-full-access|status")
1638 return
1639 }
1640 if statusOnly {
1641 _ = gw.sendText(ctx, adapter, msg, gw.toolApprovalModeStatusText(key, msg))
1642 return
1643 }
1644 persistErr := gw.setToolApprovalModeForMessage(key, msg, mode)
1645 text := toolApprovalModeChangedText(mode)
1646 if persistErr != nil {
1647 text += "\n当前会话已生效,但保存到设置失败:" + persistErr.Error()
1648 }
1649 _ = gw.sendText(ctx, adapter, msg, text)
1650
1651 case strings.HasPrefix(msg.Text, "/queue"):
1652 if reply, handled, kick := gw.handleQueueInboxCommand(ctx, key, msg); handled {
1653 _ = gw.sendText(ctx, adapter, msg, reply)
1654 if kick {
1655 gw.kickInbox(ctx, adapter, key, msg)
1656 }
1657 return
1658 }
1659 mode, clear, statusOnly, ok := parseQueueCommand(msg.Text)
1660 if !ok {
1661 _ = gw.sendText(ctx, adapter, msg, "用法: /queue steer|followup|collect|interrupt|status|list|show|delete|move|pause|resume|retry|default")
1662 return
1663 }
1664 if statusOnly {
1665 _ = gw.sendText(ctx, adapter, msg, gw.queueStatusText(key, msg))
1666 return
1667 }
1668 if clear {
1669 gw.sessions.ClearQueueMode(key)
1670 _ = gw.sendText(ctx, adapter, msg, "已恢复默认队列模式:"+queueModeLabel(gw.queueMode(key, msg))+"。")
1671 return
1672 }
1673 gw.sessions.SetQueueMode(key, mode)
1674 _ = gw.sendText(ctx, adapter, msg, "已切换队列模式:"+queueModeLabel(mode)+"。")
1675
1676 case slashCommandVerb(msg.Text) == "/projects":
1677 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1678 return
1679 }
1680 query := strings.TrimSpace(strings.TrimPrefix(msg.Text, "/projects"))
1681 _ = gw.sendText(ctx, adapter, msg, formatBotProjects(gw.buildProjectIndex(), query, botProjectListLimit))
1682
1683 case slashCommandVerb(msg.Text) == "/use":
1684 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1685 return
1686 }
1687 _ = gw.sendText(ctx, adapter, msg, gw.handleUseProjectCommand(ctx, msg, msg.Text))
1688
1689 case slashCommandVerb(msg.Text) == "/model":
1690 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1691 return
1692 }
1693 _ = gw.sendText(ctx, adapter, msg, gw.handleModelCommand(ctx, msg, msg.Text))
1694
1695 case slashCommandVerb(msg.Text) == "/sessions":
1696 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1697 return
1698 }
1699 _ = gw.sendText(ctx, adapter, msg, gw.handleSessionsCommand(msg.Text))
1700
1701 case slashCommandVerb(msg.Text) == "/attach":
1702 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1703 return
1704 }
1705 _ = gw.sendText(ctx, adapter, msg, gw.handleAttachSessionCommand(ctx, msg, msg.Text))
1706
1707 case slashCommandVerb(msg.Text) == "/search":
1708 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1709 return
1710 }
1711 _ = gw.sendText(ctx, adapter, msg, gw.handleProjectSearchCommand(ctx, msg.Text))
1712
1713 case strings.HasPrefix(msg.Text, "/desktop"):
1714 // God view over the embedding desktop app: listing every live desktop
1715 // session and answering its approvals is strictly more power than the
1716 // per-session approver role, so gate on admin.
1717 if !gw.requireCommandRole(ctx, adapter, msg, "admin") {
1718 return
1719 }
1720 _ = gw.sendText(ctx, adapter, msg, gw.handleDesktopCommand(msg))
1721
1722 case strings.HasPrefix(msg.Text, "/status"):
1723 active := gw.sessions.ActiveCount()
1724 pending := gw.sessions.PendingCount(key)
1725 gw.mu.Lock()
1726 sessions := len(gw.controllers)
1727 gw.mu.Unlock()
1728 mode := gw.currentToolApprovalMode(key, msg)
1729 _ = gw.sendText(ctx, adapter, msg, fmt.Sprintf("活跃任务数: %d\n保留会话数: %d\n工具审批模式: %s\n队列模式: %s\n当前会话排队: %d\n连接健康: %s", active, sessions, toolApprovalModeLabel(mode), queueModeLabel(gw.queueMode(key, msg)), pending, gw.adapterHealthSummaryText()))
1730
1731 case strings.HasPrefix(msg.Text, "/help"):
1732 _ = gw.sendText(ctx, adapter, msg, botHelpText())
1733 }
1734 }
1735
1736 func (gw *BotGateway) kickInbox(ctx context.Context, adapter Adapter, key string, fallback InboundMessage) {
1737 if gw.sessions.IsActive(key) {
1738 return
1739 }
1740 next := gw.nextInboxTurn(key, fallback)
1741 if next == nil {
1742 return
1743 }
1744 if !gw.sessions.TryAcquireIdle(key) {
1745 return
1746 }
1747 gw.turnWG.Go(func() {
1748 gw.runTurnItem(ctx, adapter, key, next.msg, next.itemID, nil)
1749 })
1750 }
1751
1752 func slashCommandVerb(text string) string {
1753 parts := strings.Fields(strings.TrimSpace(text))
1754 if len(parts) == 0 {
1755 return ""
1756 }
1757 return strings.ToLower(parts[0])
1758 }
1759
1760 func (gw *BotGateway) handleUseProjectCommand(ctx context.Context, msg InboundMessage, text string) string {
1761 key := BuildSessionKey(msg.Session())
1762 selector := parseUseProjectSelector(text)
1763 if selector == "" {
1764 return "用法: /use project <项目 id|名称|路径>,或 /use project default 恢复默认路由。"
1765 }
1766 if isDefaultBotSelector(selector) {
1767 switched, err := gw.setSessionRuntimeOverride(ctx, key, msg, sessionRuntimeOverride{}, false)
1768 if err != nil {
1769 return botRuntimeSwitchFailedText("切换项目")
1770 }
1771 if !switched {
1772 return botRuntimeSwitchBusyText()
1773 }
1774 return "已恢复当前远端会话的默认项目路由。下一条消息会按 bot 配置重新选择 workspace。"
1775 }
1776 projects := gw.buildProjectIndex()
1777 project, matches := resolveBotProject(projects, selector)
1778 if project.Root == "" {
1779 if len(matches) > 0 {
1780 return "匹配到多个项目,请使用项目 id:\n" + formatBotProjects(matches, "", botProjectListLimit)
1781 }
1782 return "没有匹配的项目。可先用 /projects 查看当前索引。"
1783 }
1784 switched, err := gw.setSessionRuntimeOverride(ctx, key, msg, sessionRuntimeOverride{
1785 channel: ChannelConfig{WorkspaceRoot: project.Root},
1786 label: "project:" + project.ID,
1787 }, true)
1788 if err != nil {
1789 return botRuntimeSwitchFailedText("切换项目")
1790 }
1791 if !switched {
1792 return botRuntimeSwitchBusyText()
1793 }
1794 return fmt.Sprintf("已将当前远端会话切到项目 %s %s。\n下一条消息将在 %s 中运行。", project.ID, project.Name, displayBotPath(project.Root))
1795 }
1796
1797 // handleModelCommand 处理 /model:无参查询当前会话生效模型,带参切换当前
1798 // 远端会话的模型(可选 --provider <name> 一并切换供应商)。模型以
1799 // provider/model 写入会话运行时覆盖,仅影响当前会话。
1800 func (gw *BotGateway) handleModelCommand(ctx context.Context, msg InboundMessage, text string) string {
1801 model, provider, statusOnly, ok := parseModelSelector(text)
1802 if !ok {
1803 return "用法: /model <模型名> [--provider <供应商>],或 /model 查看当前模型。"
1804 }
1805 if statusOnly {
1806 // 查询会话生效模型:走完整解析(覆盖 → 通道/路由 → 全局默认),否则
1807 // per-channel/per-connection 的模型设置(如钉钉直配 model)会被漏报。
1808 effective, _, _ := gw.sessionOptionsForMessage(msg)
1809 if strings.TrimSpace(effective) == "" {
1810 return "当前会话未指定模型,使用 bot 默认模型。"
1811 }
1812 return fmt.Sprintf("当前会话模型:%s", effective)
1813 }
1814 if strings.TrimSpace(model) == "" && strings.TrimSpace(provider) != "" {
1815 // 仅 provider 无模型名会存成无法解析的 "provider/" 空模型,下一条消息
1816 // 构建会话失败;要求显式模型名。
1817 return "用法: /model <模型名> [--provider <供应商>],或 /model 查看当前模型。"
1818 }
1819 key := BuildSessionKey(msg.Session())
1820 ref := strings.TrimSpace(model)
1821 if provider != "" {
1822 ref = strings.TrimSpace(provider) + "/" + ref
1823 }
1824 // 失败原子性:先校验模型可解析且已配置,无效则直接拒绝并保留当前
1825 // controller,不写入覆盖、不销毁旧会话(否则下一条消息构建失败)。
1826 if gw.cfg.ModelResolver != nil {
1827 if err := gw.cfg.ModelResolver(ref); err != nil {
1828 return fmt.Sprintf("模型 %s 不可用:%v", ref, err)
1829 }
1830 }
1831 // 复用 /use 的会话覆盖机制:只改 model,保留现有 workspace/tool 覆盖。
1832 var existing sessionRuntimeOverride
1833 gw.mu.Lock()
1834 existing = gw.sessionOverrides[key]
1835 gw.mu.Unlock()
1836 existing.channel.Model = ref
1837 switched, err := gw.setSessionRuntimeOverride(ctx, key, msg, existing, true)
1838 if err != nil {
1839 return botRuntimeSwitchFailedText("切换模型")
1840 }
1841 if !switched {
1842 return botRuntimeSwitchBusyText()
1843 }
1844 if provider != "" {
1845 return fmt.Sprintf("已将当前会话模型切换到 %s(供应商 %s)。", model, provider)
1846 }
1847 return fmt.Sprintf("已将当前会话模型切换到 %s。", ref)
1848 }
1849
1850 func parseModelSelector(text string) (model, provider string, statusOnly, ok bool) {
1851 parts := strings.Fields(text)
1852 if len(parts) == 0 || strings.ToLower(strings.TrimSpace(parts[0])) != "/model" {
1853 return "", "", false, false
1854 }
1855 if len(parts) == 1 {
1856 return "", "", true, true
1857 }
1858 rest := parts[1:]
1859 var models, providers []string
1860 for i := 0; i < len(rest); i++ {
1861 tok := rest[i]
1862 if strings.EqualFold(tok, "--provider") || strings.EqualFold(tok, "-p") {
1863 if i+1 < len(rest) {
1864 providers = append(providers, rest[i+1])
1865 i++
1866 }
1867 continue
1868 }
1869 if strings.HasPrefix(tok, "-") {
1870 continue
1871 }
1872 models = append(models, tok)
1873 }
1874 if len(models) == 0 {
1875 return "", strings.Join(providers, " "), false, true
1876 }
1877 return strings.Join(models, " "), strings.Join(providers, " "), false, true
1878 }
1879
1880 func parseUseProjectSelector(text string) string {
1881 parts := strings.Fields(text)
1882 if len(parts) < 2 || strings.ToLower(parts[0]) != "/use" {
1883 return ""
1884 }
1885 if len(parts) >= 3 && strings.EqualFold(parts[1], "project") {
1886 return strings.TrimSpace(strings.Join(parts[2:], " "))
1887 }
1888 return strings.TrimSpace(strings.Join(parts[1:], " "))
1889 }
1890
1891 func (gw *BotGateway) handleSessionsCommand(text string) string {
1892 query := parseSessionsQuery(text)
1893 projects := gw.buildProjectIndex()
1894 sessions := gw.buildSessionIndex(projects)
1895 return formatBotSessions(sessions, query, botSessionListLimit)
1896 }
1897
1898 func parseSessionsQuery(text string) string {
1899 parts := strings.Fields(text)
1900 if len(parts) <= 1 {
1901 return ""
1902 }
1903 if strings.EqualFold(parts[1], "search") {
1904 return strings.TrimSpace(strings.Join(parts[2:], " "))
1905 }
1906 return strings.TrimSpace(strings.Join(parts[1:], " "))
1907 }
1908
1909 func (gw *BotGateway) handleAttachSessionCommand(ctx context.Context, msg InboundMessage, text string) string {
1910 key := BuildSessionKey(msg.Session())
1911 selector := parseAttachSessionSelector(text)
1912 if selector == "" {
1913 return "用法: /attach session <会话 id|关键词|path:...>"
1914 }
1915 projects := gw.buildProjectIndex()
1916 sessions := gw.buildSessionIndex(projects)
1917 session, matches := resolveBotSession(sessions, selector)
1918 if session.ID == "" {
1919 if len(matches) > 0 {
1920 return "匹配到多个会话,请使用会话 id:\n" + formatBotSessions(matches, "", botSessionListLimit)
1921 }
1922 return "没有匹配的会话。可先用 /sessions search <关键词> 查看当前索引。"
1923 }
1924 if session.SessionPath == "" {
1925 return "这个会话没有可恢复的 path: transcript,暂时不能 attach。"
1926 }
1927 if info, err := os.Stat(session.SessionPath); err != nil || info.IsDir() {
1928 return "会话文件不可用或已被移动:" + displayBotPath(session.SessionPath)
1929 }
1930 workspaceRoot := session.WorkspaceRoot
1931 if workspaceRoot == "" {
1932 project := botProjectForPath(projects, session.SessionPath)
1933 workspaceRoot = project.Root
1934 }
1935 switched, err := gw.setSessionRuntimeOverride(ctx, key, msg, sessionRuntimeOverride{
1936 channel: ChannelConfig{WorkspaceRoot: workspaceRoot},
1937 sessionPath: session.SessionPath,
1938 label: "session:" + session.ID,
1939 }, true)
1940 if err != nil {
1941 return botRuntimeSwitchFailedText("attach")
1942 }
1943 if !switched {
1944 return botRuntimeSwitchBusyText()
1945 }
1946 projectName := firstNonEmptyString(session.ProjectName, botProjectName(workspaceRoot), "global")
1947 return fmt.Sprintf("已 attach 到会话 %s(%s)。\n下一条消息会从 %s 继续。", session.ID, projectName, displayBotPath(session.SessionPath))
1948 }
1949
1950 func parseAttachSessionSelector(text string) string {
1951 parts := strings.Fields(text)
1952 if len(parts) < 3 || !strings.EqualFold(parts[0], "/attach") || !strings.EqualFold(parts[1], "session") {
1953 return ""
1954 }
1955 return strings.TrimSpace(strings.Join(parts[2:], " "))
1956 }
1957
1958 func (gw *BotGateway) handleProjectSearchCommand(ctx context.Context, text string) string {
1959 parts := strings.Fields(text)
1960 if len(parts) < 3 || !strings.EqualFold(parts[1], "all") {
1961 return "用法: /search all <关键词>"
1962 }
1963 query := strings.TrimSpace(strings.Join(parts[2:], " "))
1964 searchCtx, cancel := context.WithTimeout(ctx, 8*time.Second)
1965 defer cancel()
1966 results, err := searchBotProjects(searchCtx, gw.buildProjectIndex(), query, botSearchListLimit)
1967 if err != nil {
1968 return "检索失败:" + err.Error()
1969 }
1970 return formatBotProjectSearchResults(results, botSearchListLimit)
1971 }
1972
1973 func botSessionHasActiveWork(state *sessionState) bool {
1974 if state == nil || state.ctrl == nil {
1975 return false
1976 }
1977 status, ok := safeBotControllerRuntimeStatus(state.ctrl)
1978 if !ok {
1979 return true
1980 }
1981 return status.Running || status.PendingPrompt || status.BackgroundJobs > 0
1982 }
1983
1984 func safeBotControllerRuntimeStatus(ctrl botController) (control.RuntimeStatus, bool) {
1985 if ctrl == nil {
1986 return control.RuntimeStatus{}, false
1987 }
1988 return ctrl.RuntimeStatus(), true
1989 }
1990
1991 func (gw *BotGateway) sessionRuntimeOverrideForMessage(msg InboundMessage) (sessionRuntimeOverride, bool) {
1992 key := BuildSessionKey(msg.Session())
1993 gw.mu.Lock()
1994 defer gw.mu.Unlock()
1995 override, ok := gw.sessionOverrides[key]
1996 return override, ok
1997 }
1998
1999 func isDefaultBotSelector(selector string) bool {
2000 switch strings.ToLower(strings.TrimSpace(selector)) {
2001 case "default", "reset", "inherit", "global", "none", "默认", "重置":
2002 return true
2003 default:
2004 return false
2005 }
2006 }
2007
2008 func parseQueueCommand(text string) (mode string, clear bool, statusOnly bool, ok bool) {
2009 parts := strings.Fields(text)
2010 if len(parts) == 0 || strings.ToLower(strings.TrimSpace(parts[0])) != "/queue" {
2011 return "", false, false, false
2012 }
2013 if len(parts) == 1 {
2014 return "", false, true, true
2015 }
2016 switch strings.ToLower(strings.TrimSpace(parts[1])) {
2017 case "status", "state", "show", "状态", "查看":
2018 return "", false, true, true
2019 case "default", "reset", "inherit", "默认", "重置":
2020 return "", true, false, true
2021 default:
2022 if normalized := NormalizeOptionalQueueMode(parts[1]); normalized != "" {
2023 return normalized, false, false, true
2024 }
2025 return "", false, false, false
2026 }
2027 }
2028
2029 func (gw *BotGateway) queueStatusText(key string, msg InboundMessage) string {
2030 inboxN := 0
2031 paused := false
2032 if api := gw.sessionAPI(key); api != nil {
2033 snap := api.InboxSnapshot()
2034 inboxN = len(snap.Items)
2035 paused = snap.Paused
2036 }
2037 return fmt.Sprintf("当前队列模式:%s\n持久化 Inbox: %d%s\n全局上限: %d\n溢出策略: 拒绝新消息(queue_drop 已弃用)\n用法:/queue steer|followup|collect|interrupt|status|list|show|delete|move|pause|resume|retry|default",
2038 queueModeLabel(gw.queueMode(key, msg)),
2039 inboxN,
2040 map[bool]string{true: " (paused)", false: ""}[paused],
2041 sessioninbox.DefaultMaxItems,
2042 )
2043 }
2044
2045 func queueModeLabel(mode string) string {
2046 switch NormalizeQueueMode(mode) {
2047 case QueueModeFollowup:
2048 return "逐条跟进"
2049 case QueueModeCollect:
2050 return "合并收集"
2051 case QueueModeInterrupt:
2052 return "打断重跑"
2053 default:
2054 return "即时补充"
2055 }
2056 }
2057
2058 func (gw *BotGateway) adapterHealthSummaryText() string {
2059 snapshots := gw.AdapterHealth()
2060 if len(snapshots) == 0 {
2061 return "未启动"
2062 }
2063 parts := make([]string, 0, len(snapshots))
2064 for _, h := range snapshots {
2065 label := strings.TrimSpace(h.ID)
2066 if label == "" {
2067 label = string(h.Platform)
2068 }
2069 status := strings.TrimSpace(h.Status)
2070 if status == "" {
2071 status = "unknown"
2072 }
2073 parts = append(parts, fmt.Sprintf("%s=%s", label, status))
2074 }
2075 return strings.Join(parts, ", ")
2076 }
2077
2078 func parseToolApprovalModeCommand(text string) (mode string, statusOnly bool, ok bool) {
2079 parts := strings.Fields(text)
2080 if len(parts) == 0 {
2081 return "", false, false
2082 }
2083 cmd := strings.ToLower(strings.TrimSpace(parts[0]))
2084 switch cmd {
2085 case "/mode":
2086 if len(parts) == 1 {
2087 return "", true, true
2088 }
2089 return parseToolApprovalModeArg(parts[1])
2090 default:
2091 return "", false, false
2092 }
2093 }
2094
2095 func parseToolApprovalModeArg(arg string) (mode string, statusOnly bool, ok bool) {
2096 switch strings.ToLower(strings.TrimSpace(arg)) {
2097 case "status", "state", "show", "状态", "查看":
2098 return "", true, true
2099 case "danger-full-access", "full", "full-access", "完全权限":
2100 return control.ToolApprovalDangerFullAccess, false, true
2101 case "read-only", "readonly", "ask", "仅可查看":
2102 return control.ToolApprovalReadOnly, false, true
2103 case "workspace-write", "workspace", "auto", "yolo", "工作区内修改":
2104 return control.ToolApprovalWorkspaceWrite, false, true
2105 default:
2106 return "", false, false
2107 }
2108 }
2109
2110 func (gw *BotGateway) setToolApprovalModeForMessage(key string, msg InboundMessage, mode string) error {
2111 mode = normalizeBotToolApprovalMode(mode)
2112 var ctrl botController
2113
2114 gw.mu.Lock()
2115 if state, ok := gw.controllers[key]; ok {
2116 ctrl = state.ctrl
2117 }
2118 gw.updateToolApprovalModeDefaultLocked(msg, mode)
2119 gw.mu.Unlock()
2120
2121 if ctrl != nil {
2122 ctrl.SetToolApprovalMode(mode)
2123 }
2124 if gw.cfg.OnToolApprovalModeChange != nil {
2125 return gw.cfg.OnToolApprovalModeChange(msg, mode)
2126 }
2127 return nil
2128 }
2129
2130 func (gw *BotGateway) updateToolApprovalModeDefaultLocked(msg InboundMessage, mode string) {
2131 if id := strings.TrimSpace(msg.ConnectionID); id != "" {
2132 if gw.cfg.ConnectionChannels == nil {
2133 gw.cfg.ConnectionChannels = make(map[string]ChannelConfig)
2134 }
2135 channel := gw.cfg.ConnectionChannels[id]
2136 channel.ToolApprovalMode = mode
2137 gw.cfg.ConnectionChannels[id] = channel
2138 return
2139 }
2140 if msg.Platform != "" {
2141 if gw.cfg.Channels == nil {
2142 gw.cfg.Channels = make(map[Platform]ChannelConfig)
2143 }
2144 channel := gw.cfg.Channels[msg.Platform]
2145 channel.ToolApprovalMode = mode
2146 gw.cfg.Channels[msg.Platform] = channel
2147 return
2148 }
2149 gw.cfg.ToolApprovalMode = mode
2150 }
2151
2152 func (gw *BotGateway) currentToolApprovalMode(key string, msg InboundMessage) string {
2153 var ctrl botController
2154 gw.mu.Lock()
2155 if state, ok := gw.controllers[key]; ok {
2156 ctrl = state.ctrl
2157 }
2158 gw.mu.Unlock()
2159 if ctrl != nil {
2160 return ctrl.ToolApprovalMode()
2161 }
2162 _, _, mode := gw.sessionOptionsForMessage(msg)
2163 return mode
2164 }
2165
2166 func (gw *BotGateway) toolApprovalModeStatusText(key string, msg InboundMessage) string {
2167 mode := gw.currentToolApprovalMode(key, msg)
2168 return fmt.Sprintf("当前权限:%s\n用法:/mode read-only|workspace-write|danger-full-access|status", toolApprovalModeLabel(mode))
2169 }
2170
2171 func toolApprovalModeChangedText(mode string) string {
2172 switch normalizeBotToolApprovalMode(mode) {
2173 case control.ToolApprovalDangerFullAccess:
2174 return "已切换为完全权限:普通工具审批将自动放行,显式禁止仍然生效。"
2175 case control.ToolApprovalWorkspaceWrite:
2176 return "已切换为工作区内修改:工作区与会话临时目录可写,越界操作需要授权。"
2177 default:
2178 return "已切换为仅可查看:读取可直接执行,写入与外部副作用需要授权。"
2179 }
2180 }
2181
2182 func toolApprovalModeLabel(mode string) string {
2183 switch normalizeBotToolApprovalMode(mode) {
2184 case control.ToolApprovalDangerFullAccess:
2185 return "完全权限"
2186 case control.ToolApprovalWorkspaceWrite:
2187 return "工作区内修改"
2188 default:
2189 return "仅可查看"
2190 }
2191 }
2192
2193 func (gw *BotGateway) runTurn(ctx context.Context, adapter Adapter, key string, msg InboundMessage, cleanup func()) {
2194 gw.runTurnItem(ctx, adapter, key, msg, "", cleanup)
2195 }
2196
2197 func (gw *BotGateway) runTurnItem(ctx context.Context, adapter Adapter, key string, msg InboundMessage, inboxItemID string, cleanup func()) {
2198 gw.logger.Info("bot turn started", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8])
2199 defer gw.finishTurnItem(ctx, adapter, key, msg, cleanup)
2200
2201 state := gw.sessionForNewTurn(ctx, adapter, key, msg)
2202 if state == nil {
2203 return
2204 }
2205 gw.rememberSessionReady(msg, state.ctrl)
2206
2207 // 构建输入文本:群聊中在消息前加上发送者名,并把 IM 媒体保存为 @附件引用。
2208 input := msg.Text
2209 if inboxItemID == "" {
2210 input = gw.inputTextWithMedia(ctx, adapter, msg, state)
2211 }
2212 if inboxItemID == "" && msg.ChatType == ChatGroup {
2213 userName := strings.TrimSpace(msg.UserName)
2214 if msg.ResolveUserName != nil {
2215 if resolved := strings.TrimSpace(msg.ResolveUserName(ctx)); resolved != "" {
2216 userName = resolved
2217 }
2218 }
2219 input = fmt.Sprintf("[%s] %s", userName, input)
2220 }
2221
2222 // 发送"正在输入"状态
2223 _ = adapter.SendTyping(ctx, msg.ChatID)
2224
2225 // 创建事件渲染 sink
2226 sink := newRenderSink(
2227 ctx,
2228 adapter,
2229 msg.ConnectionID,
2230 msg.Domain,
2231 msg.ChatID,
2232 msg.ChatType,
2233 msg.UserID,
2234 msg.MessageID,
2235 gw.logger,
2236 func(approval event.Approval) {
2237 gw.mu.Lock()
2238 if state.pendingApprovals == nil {
2239 state.pendingApprovals = make(map[string]event.Approval)
2240 }
2241 state.pendingApprovals[approval.ID] = approval
2242 state.lastApprovalID = approval.ID
2243 gw.mu.Unlock()
2244 },
2245 func(ask event.Ask) {
2246 gw.mu.Lock()
2247 if state.pendingAsks == nil {
2248 state.pendingAsks = make(map[string][]event.AskQuestion)
2249 }
2250 state.pendingAsks[ask.ID] = ask.Questions
2251 state.lastAskID = ask.ID
2252 gw.mu.Unlock()
2253 },
2254 )
2255 // Finish initializing the sink before publishing it as the live target: once
2256 // setTarget runs, other goroutines can reach this sink via state.sink.Emit.
2257 sink.ctrl = state.ctrl
2258 state.sink.setTarget(sink)
2259 defer state.sink.setTarget(nil)
2260
2261 // 创建带取消的 context
2262 turnCtx, cancel := context.WithCancel(ctx)
2263 defer cancel()
2264
2265 gw.mu.Lock()
2266 live := gw.controllers[key] == state
2267 if live {
2268 state.cancel = cancel
2269 }
2270 state.lastActive = time.Now()
2271 gw.mu.Unlock()
2272 if !live {
2273 // The session was closed (gateway stop or runtime rebuild) after this
2274 // turn picked it up; a cancel published now would never be consumed, so
2275 // abort the turn instead of running it uncancellable.
2276 cancel()
2277 }
2278
2279 // 运行一轮对话
2280 var err error
2281 if inboxItemID == "" {
2282 err = state.ctrl.RunTurn(turnCtx, input)
2283 } else if api, ok := state.ctrl.(interface {
2284 RunInboxTurn(context.Context, string) error
2285 }); ok {
2286 err = api.RunInboxTurn(turnCtx, inboxItemID)
2287 } else {
2288 err = fmt.Errorf("controller cannot run durable inbox item")
2289 }
2290 sink.Emit(event.Event{Kind: event.TurnDone, Err: err})
2291 if err != nil {
2292 gw.logger.Warn("turn error", "session", key[:8], "err", err)
2293 return
2294 }
2295 gw.logger.Info("bot turn completed", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8])
2296 }
2297
2298 func (gw *BotGateway) inputTextWithMedia(ctx context.Context, adapter Adapter, msg InboundMessage, state *sessionState) string {
2299 input := msg.Text
2300 if len(msg.MediaURLs) == 0 && len(msg.Media) == 0 {
2301 return input
2302 }
2303 workspaceRoot := ""
2304 if state != nil && state.ctrl != nil {
2305 workspaceRoot = state.ctrl.WorkspaceRoot()
2306 }
2307 if strings.TrimSpace(workspaceRoot) == "" {
2308 _, workspaceRoot, _ = gw.sessionOptionsForMessage(msg)
2309 }
2310 refs, errs := saveInboundMedia(ctx, workspaceRoot, msg.MediaURLs)
2311 itemRefs, fallbacks, itemErrs := saveInboundMediaItems(ctx, workspaceRoot, msg.Media)
2312 refs = append(refs, itemRefs...)
2313 errs = append(errs, itemErrs...)
2314 if len(errs) > 0 {
2315 gw.logger.Warn("bot media attachment failed", "platform", msg.Platform, "chat", hashID(msg.ChatID), "errors", len(errs))
2316 _ = gw.sendText(ctx, adapter, msg, fmt.Sprintf("有 %d 个附件保存失败;我会先处理可用内容。", len(errs)))
2317 }
2318 return appendMediaRefs(appendMediaFallbacks(input, fallbacks), refs)
2319 }
2320
2321 func (gw *BotGateway) getOrCreateSession(ctx context.Context, key string, msg InboundMessage) *sessionState {
2322 profile := gw.sessionProfileForMessage(msg)
2323 switch state, claim := gw.claimSession(key, msg, profile); claim {
2324 case sessionReused:
2325 safeBotSetToolApprovalMode(state.ctrl, profile.toolApprovalMode)
2326 gw.logger.Info("bot session reused", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8])
2327 return state
2328 case sessionChangeDeferred:
2329 safeBotSetToolApprovalMode(state.ctrl, profile.toolApprovalMode)
2330 gw.logger.Warn("bot session runtime change deferred while work is active", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8])
2331 return state
2332 case sessionRetired:
2333 gw.closeSessionState(state)
2334 gw.logger.Warn("bot session runtime changed; rebuilding", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8], "old_workspace_set", strings.TrimSpace(state.workspaceRoot) != "", "new_workspace_set", profile.workspaceRoot != "", "old_model", state.model, "new_model", profile.model)
2335 }
2336
2337 // Create the lease owner first so recovery or intentional transitions can
2338 // move ownership before the controller commits to the target path.
2339 sessionSink := &sessionEventSink{}
2340 leases := control.NewSessionLeaseKeeper()
2341 state := &sessionState{
2342 sink: sessionSink,
2343 leases: leases,
2344 platform: msg.Platform,
2345 connectionID: strings.TrimSpace(msg.ConnectionID),
2346 model: profile.model,
2347 workspaceRoot: profile.workspaceRoot,
2348 toolApprovalMode: profile.toolApprovalMode,
2349 sessionPath: profile.sessionPath,
2350 pendingAsks: make(map[string][]event.AskQuestion),
2351 createdAt: time.Now(),
2352 lastActive: time.Now(),
2353 }
2354 state.onSessionTransition = gw.botSessionTransitionHandler(key, msg, state)
2355 gw.logger.Info("bot session creating", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8], "model", profile.model, "workspace_set", profile.workspaceRoot != "", "tool_approval_mode", profile.toolApprovalMode)
2356 ctrl, err := gw.buildBotController(ctx, boot.Options{
2357 Model: profile.model,
2358 MaxSteps: gw.cfg.MaxSteps,
2359 MaxStepsKey: "bot.max_steps",
2360 RequireKey: true,
2361 Sink: sessionSink,
2362 StatsSource: "bot",
2363 WorkspaceRoot: profile.workspaceRoot,
2364 SessionDir: botSessionDir(profile.workspaceRoot),
2365 ApprovalTimeout: gw.approvalTimeout(),
2366 OnSessionRecovered: gw.botSessionRecoveredHandler(key, msg, state),
2367 OnSessionTransition: state.onSessionTransition,
2368 })
2369 if err != nil {
2370 leases.Release()
2371 gw.logger.Error("build controller failed", "err", secrets.RedactError(err))
2372 return nil
2373 }
2374 state.ctrl = ctrl
2375 if identity, ok := any(ctrl).(control.IdentityLifecycle); ok && identity.UsesExclusiveSession() {
2376 ref, bindErr := bindBotSessionIdentity(ctx, identity, profile, msg)
2377 if bindErr != nil && (profile.sessionRefOptional || profile.sessionPathOptional) {
2378 gw.logger.Warn("mapped bot session unavailable; starting fresh", "err", bindErr)
2379 ref, bindErr = identity.BindFreshSession(ctx, "")
2380 state.mappingDegraded = bindErr == nil
2381 }
2382 if bindErr != nil {
2383 ctrl.Close()
2384 leases.Release()
2385 gw.logger.Error("bind bot v3 session failed", "err", secrets.RedactError(bindErr))
2386 return nil
2387 }
2388 state.sessionRef = ref
2389 state.sessionPath = ""
2390 } else if profile.sessionPath != "" {
2391 // A mapped binding degrades to a fresh session on failure; only an
2392 // explicit /attach is allowed to hard-fail the message, because the
2393 // user named that exact session.
2394 degrade := func(reason string, err error) bool {
2395 if !profile.sessionPathOptional {
2396 return false
2397 }
2398 gw.logger.Warn("mapped bot session unavailable; starting fresh", "reason", reason, "session_path", profile.sessionPath, "err", err)
2399 profile.sessionPath = ""
2400 state.sessionPath = ""
2401 state.mappingDegraded = true
2402 return true
2403 }
2404 if err := leases.Rebind(profile.sessionPath); err != nil {
2405 if !degrade("lease held elsewhere", err) {
2406 ctrl.Close()
2407 leases.Release()
2408 gw.logger.Error("attached bot session is in use", "err", control.SessionInUseMessage(err))
2409 return nil
2410 }
2411 } else if loaded, err := agent.LoadSession(profile.sessionPath); err != nil {
2412 if os.IsNotExist(err) && profile.sessionPathOptional {
2413 // First message on a deterministic chat→file path: pin the new
2414 // conversation there instead of orphaning a timestamp file.
2415 ctrl.SetSessionPath(profile.sessionPath)
2416 } else if !degrade("load failed", err) {
2417 ctrl.Close()
2418 leases.Release()
2419 if os.IsNotExist(err) {
2420 gw.logger.Error("attached bot session missing", "session_path", profile.sessionPath)
2421 } else {
2422 gw.logger.Error("attached bot session load failed", "session_path", profile.sessionPath, "err", err)
2423 }
2424 return nil
2425 }
2426 } else {
2427 ctrl.Resume(loaded, profile.sessionPath)
2428 }
2429 }
2430 ctrl.EnableInteractiveApproval()
2431 ctrl.SetToolApprovalMode(profile.toolApprovalMode)
2432 if identity, ok := any(ctrl).(control.IdentityLifecycle); !ok || !identity.UsesExclusiveSession() {
2433 ctrl.EnsureSessionPath()
2434 if err := rebindBotSessionWriteAuthority(state, ctrl.SessionPath()); err != nil {
2435 ctrl.Close()
2436 leases.Release()
2437 gw.logger.Error("bot session lease failed", "err", control.SessionInUseMessage(err))
2438 return nil
2439 }
2440 }
2441 var replace *sessionState
2442 gw.mu.Lock()
2443 // Re-check under the lock: while we were off-lock in boot.Build, a second
2444 // message for the same key may have built and registered its own session.
2445 // Reuse it only when it still targets this message's runtime profile.
2446 if existing, ok := gw.controllers[key]; ok {
2447 if sessionStateMatchesRuntime(existing, profile) {
2448 updateSessionStateRuntime(existing, msg, profile)
2449 gw.mu.Unlock()
2450 ctrl.Close()
2451 leases.Release()
2452 safeBotSetToolApprovalMode(existing.ctrl, profile.toolApprovalMode)
2453 gw.logger.Info("bot session built concurrently; discarding duplicate", "platform", msg.Platform, "chat", hashID(msg.ChatID), "session", key[:8])
2454 return existing
2455 }
2456 delete(gw.controllers, key)
2457 replace = existing
2458 }
2459 gw.controllers[key] = state
2460 gw.mu.Unlock()
2461 gw.closeSessionState(replace)
2462
2463 gw.logger.Info("bot session created", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "session", key[:8])
2464 return state
2465 }
2466
2467 func updateSessionStateRuntime(state *sessionState, msg InboundMessage, profile sessionRuntimeProfile) {
2468 if state == nil {
2469 return
2470 }
2471 if state.connectionID == "" {
2472 state.connectionID = strings.TrimSpace(msg.ConnectionID)
2473 }
2474 if state.platform == "" {
2475 state.platform = msg.Platform
2476 }
2477 state.model = profile.model
2478 state.workspaceRoot = profile.workspaceRoot
2479 state.toolApprovalMode = profile.toolApprovalMode
2480 state.sessionPath = profile.sessionPath
2481 if profile.sessionRef.SessionID != "" {
2482 if profile.sessionRef.HostID == "" {
2483 profile.sessionRef.HostID = state.sessionRef.HostID
2484 }
2485 state.sessionRef = profile.sessionRef
2486 }
2487 state.lastActive = time.Now()
2488 }
2489
2490 func (gw *BotGateway) sessionProfileForMessage(msg InboundMessage) sessionRuntimeProfile {
2491 override, enabled := gw.sessionRuntimeOverrideForMessage(msg)
2492 return gw.sessionProfileForResolvedOverride(msg, override, enabled)
2493 }
2494
2495 func (gw *BotGateway) sessionProfileForResolvedOverride(msg InboundMessage, override sessionRuntimeOverride, enabled bool) sessionRuntimeProfile {
2496 model, workspaceRoot, toolApprovalMode := gw.sessionOptionsForResolvedOverride(msg, override, enabled)
2497 var sessionPath string
2498 var sessionRef session.SessionRef
2499 sessionPathOptional := false
2500 sessionRefOptional := false
2501 if enabled {
2502 if ref, ok := parseBotSessionRefTarget(override.sessionPath); ok {
2503 sessionRef = ref
2504 } else {
2505 sessionPath = override.sessionPath
2506 }
2507 }
2508 // A persisted session_mappings binding is the durable chat→session link
2509 // the desktop writes into the connection config. Without consuming it
2510 // here, every gateway restart or runtime rebuild opened a brand-new
2511 // session file for the chat and the configured binding was display-only
2512 // (#6917, #6934).
2513 if sessionPath == "" && sessionRef.SessionID == "" {
2514 if mapped := gw.sessionMappingTargetForMessage(msg); mapped != "" {
2515 if ref, ok := parseBotSessionRefTarget(mapped); ok {
2516 sessionRef = ref
2517 sessionRefOptional = true
2518 } else if path := botSessionPathFromTarget(mapped); path != "" {
2519 sessionPath = path
2520 sessionPathOptional = true
2521 }
2522 }
2523 }
2524 // No explicit binding: pin a deterministic per-chat file so the chat reuses
2525 // one conversation across restarts (dsh-dingtalk-channel's `ding-<chatId>`
2526 // analogue). Optional, mirroring mapping degrade semantics.
2527 if sessionPath == "" && sessionRef.SessionID == "" && strings.TrimSpace(msg.ChatID) != "" {
2528 sessionRef = session.SessionRef{SessionID: "bot-" + BuildSessionKey(msg.Session())}
2529 sessionRefOptional = true
2530 }
2531 return sessionRuntimeProfile{
2532 model: strings.TrimSpace(model),
2533 workspaceRoot: strings.TrimSpace(workspaceRoot),
2534 toolApprovalMode: normalizeBotToolApprovalMode(toolApprovalMode),
2535 sessionPath: canonicalBotPath(sessionPath),
2536 sessionPathOptional: sessionPathOptional,
2537 sessionRef: sessionRef,
2538 sessionRefOptional: sessionRefOptional,
2539 }
2540 }
2541
2542 // sessionMappingPathForMessage resolves the persisted session_mappings entry
2543 // for a message to an existing session file. Only bindings that resolve to a
2544 // present, readable file participate — a moved or deleted target quietly
2545 // degrades to normal session creation rather than blocking the chat.
2546 func (gw *BotGateway) sessionMappingTargetForMessage(msg InboundMessage) string {
2547 gw.mu.Lock()
2548 var mappings []SessionMapping
2549 if msg.ConnectionID != "" {
2550 if channel, ok := gw.cfg.ConnectionChannels[msg.ConnectionID]; ok {
2551 mappings = channel.SessionMappings
2552 }
2553 }
2554 if len(mappings) == 0 {
2555 if channel, ok := gw.cfg.Channels[msg.Platform]; ok {
2556 mappings = channel.SessionMappings
2557 }
2558 }
2559 gw.mu.Unlock()
2560 mapping, ok := matchingSessionMapping(mappings, msg)
2561 if !ok {
2562 return ""
2563 }
2564 target := strings.TrimSpace(mapping.SessionID)
2565 if target == "" {
2566 target = strings.TrimSpace(mapping.SessionSource)
2567 }
2568 if _, ok := parseBotSessionRefTarget(target); ok {
2569 return target
2570 }
2571 path := botSessionPathFromTarget(target)
2572 if path == "" {
2573 return ""
2574 }
2575 if info, err := os.Stat(path); err != nil || info.IsDir() {
2576 return ""
2577 }
2578 return path
2579 }
2580
2581 func sessionStateMatchesRuntime(state *sessionState, profile sessionRuntimeProfile) bool {
2582 if state == nil || state.ctrl == nil {
2583 return false
2584 }
2585 if stateModel := strings.TrimSpace(state.model); stateModel != "" && profile.model != "" && stateModel != profile.model {
2586 return false
2587 }
2588 stateRoot := strings.TrimSpace(state.workspaceRoot)
2589 wantRoot := strings.TrimSpace(profile.workspaceRoot)
2590 if stateRoot == "" {
2591 root, ok := safeBotControllerWorkspaceRoot(state.ctrl)
2592 if ok {
2593 stateRoot = strings.TrimSpace(root)
2594 } else if wantRoot != "" {
2595 return false
2596 }
2597 }
2598 if stateRoot != wantRoot {
2599 return false
2600 }
2601 // A state that already degraded off its mapped session keeps running on
2602 // its fresh path even though the profile re-resolves the mapping each
2603 // message; rebuilding here would spawn a new session per message while the
2604 // mapped file stays unavailable.
2605 if profile.sessionPathOptional && state.mappingDegraded {
2606 return true
2607 }
2608 if profile.sessionRefOptional && state.mappingDegraded {
2609 return true
2610 }
2611 if profile.sessionRef.SessionID != "" {
2612 if state.sessionRef.SessionID != profile.sessionRef.SessionID {
2613 return false
2614 }
2615 if profile.sessionRef.HostID != "" && state.sessionRef.HostID != profile.sessionRef.HostID {
2616 return false
2617 }
2618 identity, ok := state.ctrl.(control.IdentityLifecycle)
2619 if !ok {
2620 return false
2621 }
2622 ref, bound := identity.SessionRef()
2623 return bound && ref == state.sessionRef
2624 }
2625 if canonicalBotPath(state.sessionPath) != canonicalBotPath(profile.sessionPath) {
2626 return false
2627 }
2628 if profile.sessionPath != "" && canonicalBotPath(state.ctrl.SessionPath()) != canonicalBotPath(profile.sessionPath) {
2629 return false
2630 }
2631 return true
2632 }
2633
2634 func safeBotControllerWorkspaceRoot(ctrl botController) (string, bool) {
2635 if ctrl == nil {
2636 return "", false
2637 }
2638 return ctrl.WorkspaceRoot(), true
2639 }
2640
2641 func safeBotSetToolApprovalMode(ctrl botController, mode string) {
2642 if ctrl == nil {
2643 return
2644 }
2645 ctrl.SetToolApprovalMode(mode)
2646 }
2647
2648 // defaultBotApprovalTimeout caps how long a bot session waits for a remote
2649 // user's approval/ask reply before treating it as denied, so an abandoned
2650 // prompt (or a dropped IM event) can't leave the session wedged forever
2651 // (#4626, #4402). 30 minutes is generous for a human reply yet bounded.
2652 const defaultBotApprovalTimeout = 30 * time.Minute
2653
2654 // approvalTimeout resolves the configured bot approval wait: zero uses the
2655 // bounded default; a negative value opts out (wait indefinitely).
2656 func (gw *BotGateway) approvalTimeout() time.Duration {
2657 switch {
2658 case gw.cfg.ApprovalTimeout < 0:
2659 return 0
2660 case gw.cfg.ApprovalTimeout == 0:
2661 return defaultBotApprovalTimeout
2662 default:
2663 return gw.cfg.ApprovalTimeout
2664 }
2665 }
2666
2667 func botSessionDir(workspaceRoot string) string {
2668 if strings.TrimSpace(workspaceRoot) == "" {
2669 return config.SessionDir()
2670 }
2671 if dir := config.ProjectSessionDir(workspaceRoot); dir != "" {
2672 return dir
2673 }
2674 return config.SessionDir()
2675 }
2676
2677 // BotSessionPathForChat derives the deterministic per-chat session file for a
2678 // message with no persisted mapping or /attach binding. Reusing BuildSessionKey
2679 // (already a stable chat-identity hash) makes the same chat hit one file across
2680 // restarts — the chat-side analogue of dsh-dingtalk-channel's `ding-<chatId>`
2681 // scheme. Empty when the message has no stable chat identity.
2682 func BotSessionPathForChat(sessionDir string, src SessionSource) string {
2683 if strings.TrimSpace(sessionDir) == "" || strings.TrimSpace(src.ChatID) == "" {
2684 return ""
2685 }
2686 return filepath.Join(sessionDir, "bot-"+BuildSessionKey(src)+".jsonl")
2687 }
2688
2689 func (gw *BotGateway) rememberSessionReady(msg InboundMessage, ctrl botController) {
2690 if gw.cfg.OnSessionReady == nil || ctrl == nil {
2691 return
2692 }
2693 if identity, ok := ctrl.(control.IdentityLifecycle); ok && identity.UsesExclusiveSession() {
2694 if ref, bound := identity.SessionRef(); bound {
2695 gw.rememberSessionTarget(msg, botSessionRefTarget(ref))
2696 return
2697 }
2698 }
2699 gw.rememberSessionPath(msg, ctrl.SessionPath())
2700 }
2701
2702 func (gw *BotGateway) rememberSessionPath(msg InboundMessage, sessionPath string) {
2703 gw.rememberSessionTarget(msg, botSessionTarget(sessionPath))
2704 }
2705
2706 func (gw *BotGateway) rememberSessionTarget(msg InboundMessage, sessionID string) {
2707 if gw.cfg.OnSessionReady == nil {
2708 return
2709 }
2710 if sessionID == "" {
2711 return
2712 }
2713 if err := gw.cfg.OnSessionReady(msg, sessionID); err != nil {
2714 gw.logger.Warn("remember bot session failed", "platform", msg.Platform, "connection", msg.ConnectionID, "err", err)
2715 }
2716 }
2717
2718 func botSessionRefTarget(ref session.SessionRef) string {
2719 if strings.TrimSpace(ref.HostID) == "" || strings.TrimSpace(ref.SessionID) == "" {
2720 return ""
2721 }
2722 return "session:" + ref.HostID + ":" + ref.SessionID
2723 }
2724
2725 func parseBotSessionRefTarget(target string) (session.SessionRef, bool) {
2726 target = strings.TrimSpace(target)
2727 if !strings.HasPrefix(target, "session:") {
2728 return session.SessionRef{}, false
2729 }
2730 parts := strings.SplitN(strings.TrimPrefix(target, "session:"), ":", 2)
2731 if len(parts) != 2 || strings.TrimSpace(parts[0]) == "" || strings.TrimSpace(parts[1]) == "" {
2732 return session.SessionRef{}, false
2733 }
2734 return session.SessionRef{HostID: strings.TrimSpace(parts[0]), SessionID: strings.TrimSpace(parts[1])}, true
2735 }
2736
2737 // botSessionRecoveredHandler keeps the controller path, its writer lease, and
2738 // the remote-to-session mapping on the same recovery generation. The lease
2739 // handoff runs first and is failure-atomic: if the recovery path is already
2740 // owned, the controller stays on the original path and the old lease remains
2741 // held. Mapping updates are limited to this exact sessionState so a late
2742 // callback from a retired controller cannot overwrite its replacement.
2743 func (gw *BotGateway) botSessionRecoveredHandler(key string, msg InboundMessage, state *sessionState) func(control.SessionRecoveryInfo) error {
2744 return func(info control.SessionRecoveryInfo) error {
2745 if state == nil || state.leases == nil {
2746 return nil
2747 }
2748 // Keep the lease handoff and mapping publication atomic with respect to
2749 // state retirement. In particular, never let a callback that outlives
2750 // Stop reacquire a lease after closeSessionState has released it.
2751 state.lifecycleMu.Lock()
2752 defer state.lifecycleMu.Unlock()
2753 if state.retired {
2754 return errBotSessionRetired
2755 }
2756 if err := state.leases.HandleSessionRecovered(info); err != nil {
2757 return err
2758 }
2759
2760 originalPath := canonicalBotPath(info.OriginalPath)
2761 recoveryPath := canonicalBotPath(info.RecoveryPath)
2762 live := false
2763 gw.mu.Lock()
2764 if gw.controllers[key] == state {
2765 live = true
2766 if canonicalBotPath(state.sessionPath) == originalPath {
2767 state.sessionPath = recoveryPath
2768 }
2769 if override, ok := gw.sessionOverrides[key]; ok && canonicalBotPath(override.sessionPath) == originalPath {
2770 override.sessionPath = recoveryPath
2771 gw.sessionOverrides[key] = override
2772 }
2773 }
2774 gw.mu.Unlock()
2775
2776 if live {
2777 gw.rememberSessionPath(msg, recoveryPath)
2778 }
2779 return nil
2780 }
2781 }
2782
2783 func botSessionTarget(sessionPath string) string {
2784 sessionPath = strings.TrimSpace(sessionPath)
2785 if sessionPath == "" {
2786 return ""
2787 }
2788 return "path:" + sessionPath
2789 }
2790
2791 func (gw *BotGateway) sessionOptionsForMessage(msg InboundMessage) (model string, workspaceRoot string, toolApprovalMode string) {
2792 override, enabled := gw.sessionRuntimeOverrideForMessage(msg)
2793 return gw.sessionOptionsForResolvedOverride(msg, override, enabled)
2794 }
2795
2796 func (gw *BotGateway) sessionOptionsForResolvedOverride(msg InboundMessage, override sessionRuntimeOverride, enabled bool) (model string, workspaceRoot string, toolApprovalMode string) {
2797 // cfg.ToolApprovalMode / Channels / ConnectionChannels are rewritten under
2798 // gw.mu at runtime (/mode, UpdateConnectionToolApprovalMode), so snapshot them
2799 // under a short lock and resolve outside it. Copying the ChannelConfig value is enough: writers
2800 // replace whole map entries and never mutate SessionMappings in place.
2801 gw.mu.Lock()
2802 model = gw.cfg.Model
2803 workspaceRoot = gw.cfg.WorkspaceRoot
2804 toolApprovalMode = normalizeBotToolApprovalMode(gw.cfg.ToolApprovalMode)
2805 var connChannel ChannelConfig
2806 connOK := false
2807 if msg.ConnectionID != "" {
2808 connChannel, connOK = gw.cfg.ConnectionChannels[msg.ConnectionID]
2809 }
2810 platChannel, platOK := gw.cfg.Channels[msg.Platform]
2811 gw.mu.Unlock()
2812
2813 var mappings []SessionMapping
2814 if connOK {
2815 applyBotChannelOptions(connChannel, &model, &workspaceRoot, &toolApprovalMode)
2816 mappings = connChannel.SessionMappings
2817 if mapping, ok := matchingSessionMapping(mappings, msg); ok {
2818 workspaceRoot = workspaceRootForSessionMapping(mapping, workspaceRoot)
2819 }
2820 model, workspaceRoot, toolApprovalMode = gw.applyRouteOptions(msg, model, workspaceRoot, toolApprovalMode)
2821 if enabled {
2822 applyBotChannelOptions(override.channel, &model, &workspaceRoot, &toolApprovalMode)
2823 }
2824 return model, workspaceRoot, toolApprovalMode
2825 }
2826 if platOK {
2827 applyBotChannelOptions(platChannel, &model, &workspaceRoot, &toolApprovalMode)
2828 mappings = platChannel.SessionMappings
2829 }
2830 if mapping, ok := matchingSessionMapping(mappings, msg); ok {
2831 workspaceRoot = workspaceRootForSessionMapping(mapping, workspaceRoot)
2832 }
2833 model, workspaceRoot, toolApprovalMode = gw.applyRouteOptions(msg, model, workspaceRoot, toolApprovalMode)
2834 if enabled {
2835 applyBotChannelOptions(override.channel, &model, &workspaceRoot, &toolApprovalMode)
2836 }
2837 return model, workspaceRoot, toolApprovalMode
2838 }
2839
2840 func (gw *BotGateway) applyRouteOptions(msg InboundMessage, model, workspaceRoot, toolApprovalMode string) (string, string, string) {
2841 for _, route := range gw.cfg.Routes {
2842 if routeMatchesMessage(route, msg) {
2843 applyBotChannelOptions(route.Channel, &model, &workspaceRoot, &toolApprovalMode)
2844 break
2845 }
2846 }
2847 return model, workspaceRoot, toolApprovalMode
2848 }
2849
2850 func applyBotChannelOptions(channel ChannelConfig, model *string, workspaceRoot *string, toolApprovalMode *string) {
2851 if value := strings.TrimSpace(channel.Model); value != "" {
2852 *model = value
2853 }
2854 if value := strings.TrimSpace(channel.WorkspaceRoot); value != "" {
2855 *workspaceRoot = value
2856 }
2857 if value := normalizeOptionalBotToolApprovalMode(channel.ToolApprovalMode); value != "" {
2858 *toolApprovalMode = value
2859 }
2860 }
2861
2862 func matchingSessionMapping(mappings []SessionMapping, msg InboundMessage) (SessionMapping, bool) {
2863 for i := range mappings {
2864 if sessionMappingMatches(mappings[i], msg) {
2865 return mappings[i], true
2866 }
2867 }
2868 return SessionMapping{}, false
2869 }
2870
2871 func sessionMappingMatches(mapping SessionMapping, msg InboundMessage) bool {
2872 if strings.TrimSpace(mapping.RemoteID) != strings.TrimSpace(msg.ChatID) {
2873 return false
2874 }
2875 chatType, userID, threadID := sessionMappingIdentity(msg)
2876 mappingChatType := strings.TrimSpace(mapping.ChatType)
2877 if mappingChatType == "" {
2878 return chatType == ""
2879 }
2880 if mappingChatType != chatType {
2881 return false
2882 }
2883 if strings.TrimSpace(mapping.UserID) != userID {
2884 return false
2885 }
2886 return strings.TrimSpace(mapping.ThreadID) == threadID
2887 }
2888
2889 func sessionMappingIdentity(msg InboundMessage) (chatType string, userID string, threadID string) {
2890 switch msg.ChatType {
2891 case ChatGroup, ChatGuild:
2892 chatType = string(msg.ChatType)
2893 userID = strings.TrimSpace(msg.UserID)
2894 case ChatThread:
2895 chatType = string(msg.ChatType)
2896 threadID = strings.TrimSpace(msg.ThreadID)
2897 if threadID == "" {
2898 threadID = strings.TrimSpace(msg.ChatID)
2899 }
2900 }
2901 return chatType, userID, threadID
2902 }
2903
2904 func workspaceRootForSessionMapping(mapping SessionMapping, fallback string) string {
2905 if root := strings.TrimSpace(mapping.WorkspaceRoot); root != "" {
2906 return root
2907 }
2908 if strings.EqualFold(strings.TrimSpace(mapping.Scope), "global") {
2909 return ""
2910 }
2911 return fallback
2912 }
2913
2914 func routeMatchesMessage(route RouteConfig, msg InboundMessage) bool {
2915 if value := strings.TrimSpace(route.ConnectionID); value != "" && value != strings.TrimSpace(msg.ConnectionID) {
2916 return false
2917 }
2918 if route.Platform != "" && route.Platform != msg.Platform {
2919 return false
2920 }
2921 if route.ChatType != "" && route.ChatType != msg.ChatType {
2922 return false
2923 }
2924 if value := strings.TrimSpace(route.ChatID); value != "" && value != strings.TrimSpace(msg.ChatID) {
2925 return false
2926 }
2927 if value := strings.TrimSpace(route.UserID); value != "" && value != strings.TrimSpace(msg.UserID) {
2928 return false
2929 }
2930 if value := strings.TrimSpace(route.ThreadID); value != "" && value != strings.TrimSpace(msg.ThreadID) {
2931 return false
2932 }
2933 return true
2934 }
2935
2936 func normalizeBotToolApprovalMode(mode string) string {
2937 if value := normalizeOptionalBotToolApprovalMode(mode); value != "" {
2938 return value
2939 }
2940 return control.ToolApprovalWorkspaceWrite
2941 }
2942
2943 func normalizeOptionalBotToolApprovalMode(mode string) string {
2944 if strings.TrimSpace(mode) == "" {
2945 return ""
2946 }
2947 return config.NormalizeToolApprovalMode(mode)
2948 }
2949
2950 func (gw *BotGateway) sendText(ctx context.Context, adapter Adapter, msg InboundMessage, text string) error {
2951 out := OutboundMessage{
2952 ConnectionID: msg.ConnectionID,
2953 Domain: msg.Domain,
2954 ChatID: msg.ChatID,
2955 ChatType: msg.ChatType,
2956 Text: text,
2957 ReplyToMsgID: msg.MessageID,
2958 SessionWebhook: msg.SessionWebhook,
2959 }
2960 binding := AdapterBinding{
2961 ID: strings.TrimSpace(msg.ConnectionID),
2962 Domain: strings.TrimSpace(msg.Domain),
2963 Platform: msg.Platform,
2964 Adapter: adapter,
2965 }
2966 if binding.Platform == "" && adapter != nil {
2967 binding.Platform = adapter.Platform()
2968 }
2969 if binding.ID == "" && adapter != nil {
2970 binding.ID = adapter.Name()
2971 }
2972 result, err := gw.sendViaAdapter(ctx, binding, out)
2973 if err != nil {
2974 gw.logger.Warn("bot send failed", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "reply_to", hashID(msg.MessageID), "err", err)
2975 return err
2976 }
2977 gw.logger.Info("bot send completed", "platform", msg.Platform, "chat_type", msg.ChatType, "chat", hashID(msg.ChatID), "reply_to", hashID(msg.MessageID), "message", hashID(result.MessageID))
2978 return err
2979 }
2980
2981 func (gw *BotGateway) sendViaAdapter(ctx context.Context, binding AdapterBinding, msg OutboundMessage) (SendResult, error) {
2982 if binding.Adapter == nil {
2983 return SendResult{}, errors.New("bot send: adapter is nil")
2984 }
2985 if strings.TrimSpace(msg.ConnectionID) == "" {
2986 msg.ConnectionID = binding.ID
2987 }
2988 if strings.TrimSpace(msg.Domain) == "" {
2989 msg.Domain = binding.Domain
2990 }
2991 result, err := binding.Adapter.Send(ctx, msg)
2992 gw.markAdapterSend(binding, err)
2993 for _, messageID := range result.DeliveredMessageIDs() {
2994 gw.rememberOutboundMessage(binding.Platform, binding.ID, binding.Domain, msg.ChatID, messageID)
2995 }
2996 return result, err
2997 }
2998
2999 func parseAskAnswers(questions []event.AskQuestion, raw string) []event.AskAnswer {
3000 raw = strings.TrimSpace(raw)
3001 if len(questions) == 0 {
3002 return []event.AskAnswer{{Selected: []string{raw}}}
3003 }
3004 byID := make(map[string]*event.AskQuestion, len(questions))
3005 for i := range questions {
3006 q := &questions[i]
3007 byID[q.ID] = q
3008 byID[fmt.Sprintf("%d", i+1)] = q
3009 }
3010 answerMap := make(map[string][]string, len(questions))
3011 if strings.Contains(raw, "=") {
3012 for part := range strings.SplitSeq(raw, ";") {
3013 k, v, ok := strings.Cut(part, "=")
3014 if !ok {
3015 continue
3016 }
3017 q := byID[strings.TrimSpace(k)]
3018 if q == nil {
3019 continue
3020 }
3021 answerMap[q.ID] = normalizeAskSelection(*q, strings.TrimSpace(v))
3022 }
3023 } else if len(questions) == 1 {
3024 answerMap[questions[0].ID] = normalizeAskSelection(questions[0], raw)
3025 }
3026 out := make([]event.AskAnswer, 0, len(questions))
3027 for _, q := range questions {
3028 out = append(out, event.AskAnswer{QuestionID: q.ID, Selected: answerMap[q.ID]})
3029 }
3030 return out
3031 }
3032
3033 func normalizeAskSelection(q event.AskQuestion, raw string) []string {
3034 parts := []string{raw}
3035 if q.Multi && strings.Contains(raw, ",") {
3036 parts = strings.Split(raw, ",")
3037 }
3038 out := make([]string, 0, len(parts))
3039 for _, part := range parts {
3040 part = strings.TrimSpace(part)
3041 if part == "" {
3042 continue
3043 }
3044 if idx, err := strconv.Atoi(part); err == nil && idx >= 1 && idx <= len(q.Options) {
3045 out = append(out, q.Options[idx-1].Label)
3046 continue
3047 }
3048 out = append(out, part)
3049 }
3050 return out
3051 }
3052
3053 // UpdateConnectionToolApprovalMode updates the in-memory tool approval mode for
3054 // a single bot connection without restarting the gateway. Empty mode clears the
3055 // connection override, so existing sessions inherit the current gateway default.
3056 func (gw *BotGateway) UpdateConnectionToolApprovalMode(connID, mode string) {
3057 connID = strings.TrimSpace(connID)
3058 if connID == "" {
3059 return
3060 }
3061 mode = normalizeOptionalBotToolApprovalMode(mode)
3062 type controllerMode struct {
3063 ctrl botController
3064 mode string
3065 }
3066 var updates []controllerMode
3067
3068 gw.mu.Lock()
3069 if gw.cfg.ConnectionChannels == nil {
3070 gw.cfg.ConnectionChannels = make(map[string]ChannelConfig)
3071 }
3072 ch := gw.cfg.ConnectionChannels[connID]
3073 ch.ToolApprovalMode = mode
3074 gw.cfg.ConnectionChannels[connID] = ch
3075 // Update every active session that belongs to this connection.
3076 for _, state := range gw.controllers {
3077 if state == nil || state.ctrl == nil || strings.TrimSpace(state.connectionID) != connID {
3078 continue
3079 }
3080 effectiveMode := mode
3081 if effectiveMode == "" {
3082 effectiveMode = normalizeBotToolApprovalMode(gw.cfg.ToolApprovalMode)
3083 }
3084 updates = append(updates, controllerMode{ctrl: state.ctrl, mode: effectiveMode})
3085 }
3086 gw.mu.Unlock()
3087
3088 for _, update := range updates {
3089 update.ctrl.SetToolApprovalMode(update.mode)
3090 }
3091 }
3092
3093 // SendToAdapter sends a message through the adapter identified by connID.
3094 // Returns an error if no matching adapter is found.
3095 func (gw *BotGateway) SendToAdapter(ctx context.Context, connID, domain string, msg OutboundMessage) (SendResult, error) {
3096 connID = strings.TrimSpace(connID)
3097 domain = strings.TrimSpace(domain)
3098 var target AdapterBinding
3099 gw.mu.Lock()
3100 for _, binding := range gw.adapters {
3101 if strings.TrimSpace(binding.ID) == connID &&
3102 (domain == "" || strings.EqualFold(strings.TrimSpace(binding.Domain), domain)) {
3103 target = binding
3104 break
3105 }
3106 }
3107 gw.mu.Unlock()
3108 if target.Adapter != nil {
3109 return gw.sendViaAdapter(ctx, target, msg)
3110 }
3111 return SendResult{}, fmt.Errorf("SendToAdapter: no adapter found for connection %q (domain %q)", connID, domain)
3112 }
3113
3114 // SendTextToAdapter sends a plain text message through the adapter identified by connID.
3115 func (gw *BotGateway) SendTextToAdapter(ctx context.Context, connID, domain, chatID string, chatType ChatType, text string) (SendResult, error) {
3116 return gw.SendToAdapter(ctx, connID, domain, OutboundMessage{
3117 ChatID: chatID,
3118 ChatType: chatType,
3119 Text: text,
3120 })
3121 }
3122
3123 // TestSendToAdapter sends a test message through the adapter identified by
3124 // connID. The adapter must implement TestSender (currently dingtalk, which
3125 // replies to the most recent chat it learned a session webhook for). Returns
3126 // a readable error when the adapter is missing or does not support test sends.
3127 func (gw *BotGateway) TestSendToAdapter(ctx context.Context, connID, domain, text string) (SendResult, error) {
3128 connID = strings.TrimSpace(connID)
3129 domain = strings.TrimSpace(domain)
3130 var target AdapterBinding
3131 gw.mu.Lock()
3132 for _, binding := range gw.adapters {
3133 if strings.TrimSpace(binding.ID) == connID &&
3134 (domain == "" || strings.EqualFold(strings.TrimSpace(binding.Domain), domain)) {
3135 target = binding
3136 break
3137 }
3138 }
3139 gw.mu.Unlock()
3140 if target.Adapter == nil {
3141 return SendResult{}, fmt.Errorf("no bot adapter found for %q (domain %q)", connID, domain)
3142 }
3143 ts, ok := target.Adapter.(TestSender)
3144 if !ok {
3145 return SendResult{}, fmt.Errorf("bot adapter %q does not support test sends", connID)
3146 }
3147 return ts.TestSend(ctx, text)
3148 }
3149
3149 lines GO