返回 DeepSeek-Reasonix
bot_bridge.go
根目录 / desktop / bot_bridge.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "sort"
9 "strings"
10 "sync"
11 "time"
12
13 "reasonix/internal/bot"
14 "reasonix/internal/event"
15 )
16
17 // errDriveBusy signals that a takeover drive could not start because the target
18 // controller was already running a turn. The hub translates it into a
19 // user-facing "session busy" message rather than a generic drive failure.
20 var errDriveBusy = errors.New("desktop session busy")
21
22 // botBridgeHub 是 bot 网关对桌面端的"上帝视角"桥(bot.DesktopBridge 的实现)。
23 //
24 // 职责边界(刻意保持窄):
25 // - 观察:tabEventSink.Emit 把每个桌面会话的事件旁路到 observe;hub 记录
26 // 待审批/待回答项,并把审批请求、任务完成/出错推送给订阅聊天。
27 // - 遥控审批:/desktop approve|deny|answer 经 App.ApproveTab /
28 // AnswerQuestionForTab 按 tab 寻址回写。controller 侧幂等(先到者赢),
29 // 桌面 UI 与远端并发应答互不干扰。
30 // - 不做:向桌面会话注入输入、抢占写 lease。将来的显式接管在
31 // bot.DesktopBridge 上扩展,底层走 internal/control 的接管端口。
32 //
33 // observe 跑在 controller 的事件 goroutine 上,绝不能做网络调用——通知统一
34 // 进有界队列由 worker 异步发送,队列满时丢弃并告警。
35 type botBridgeHub struct {
36 sessions func() []bot.DesktopSessionInfo
37 approveTab func(tabID, id string, allow, session, persist bool)
38 answerTab func(tabID, id string, answers []QuestionAnswer)
39 notify func(ctx context.Context, connectionID, domain string, msg bot.OutboundMessage) (bot.SendResult, error)
40 // drive 把一条远程文本提交为 tab 的新 turn,并把该 turn 的输出转发回 route。
41 drive func(tabID, text string, route bot.DesktopWatchRoute) error
42 // announce 往会话 transcript 里发一条 Notice,让桌面用户看到接管状态变化。
43 announce func(tabID, text string)
44 // persistWatchers 把订阅全集回写用户配置(bot.desktop_watchers)。
45 persistWatchers func(routes []bot.DesktopWatchRoute) error
46 // takeoverChanged 通知桌面前端刷新(TabMeta.RemoteControlled 变化)。
47 takeoverChanged func()
48 logger *slog.Logger
49
50 mu sync.Mutex
51 watchers map[string]bot.DesktopWatchRoute
52 pending map[string]desktopPendingPrompt
53 // takeovers: tabID -> 驾驶该会话的聊天路由;takeoverTabs: routeKey -> tabID。
54 takeovers map[string]bot.DesktopWatchRoute
55 takeoverTabs map[string]string
56 // watchSeq 单调递增,标记订阅快照的新旧;persist 时用它丢弃过期写入。
57 watchSeq uint64
58 // watchPersistDirty keeps a failed local mutation authoritative in memory;
59 // a later runtime refresh must not silently restore the older disk snapshot.
60 watchPersistDirty bool
61
62 // persistMu 串行化订阅落盘,并保证只写最新快照(见 SetWatch)。
63 persistMu sync.Mutex
64 lastPersistSeq uint64
65
66 queue chan desktopBridgeNotification
67 // closed 终止 run worker;queue 本身不能 close——observe 仍可能在
68 // controller 事件 goroutine 上并发 enqueue,向已关闭 channel 发送会 panic。
69 closed chan struct{}
70 closeOnce sync.Once
71 }
72
73 type desktopPendingPrompt struct {
74 tabID string
75 kind string // "approval" | "ask"
76 tool string
77 subject string
78 questions []event.AskQuestion
79 }
80
81 // desktopBridgeNotification 是一条待推送的桌面事件。text 与 card 都按订阅路由
82 // 现做,因此群聊(多用户)与私聊可给出不同详略——命令行/错误详情只发私聊,群里
83 // 只给摘要。route 非 nil 时定向发给该聊天(不看 watch 订阅),用于接管收回等必达通知。
84 type desktopBridgeNotification struct {
85 text func(route bot.DesktopWatchRoute) string
86 card func(route bot.DesktopWatchRoute) *bot.InteractiveCard
87 route *bot.DesktopWatchRoute
88 }
89
90 // isSharedChat 判断一个聊天是否为多用户场景(群/话题群/服务器频道)。私聊/单聊
91 // 只有操作者本人,可安全展示命令行等敏感详情。
92 func isSharedChat(ct bot.ChatType) bool {
93 return ct != bot.ChatDM && ct != bot.ChatDirect
94 }
95
96 func constText(s string) func(bot.DesktopWatchRoute) string {
97 return func(bot.DesktopWatchRoute) string { return s }
98 }
99
100 const (
101 botBridgeQueueSize = 64
102 botBridgeSendTimeout = 15 * time.Second
103 botBridgeSubjectLimit = 200
104 botBridgePendingLimit = 200
105 botBridgeErrTextLimit = 300
106 botBridgePromptPreview = 500
107 )
108
109 // botBridgeDeps 打包 hub 对宿主(App)的全部依赖,便于测试注入。
110 type botBridgeDeps struct {
111 sessions func() []bot.DesktopSessionInfo
112 approveTab func(tabID, id string, allow, session, persist bool)
113 answerTab func(tabID, id string, answers []QuestionAnswer)
114 notify func(ctx context.Context, connectionID, domain string, msg bot.OutboundMessage) (bot.SendResult, error)
115 drive func(tabID, text string, route bot.DesktopWatchRoute) error
116 announce func(tabID, text string)
117 persistWatchers func(routes []bot.DesktopWatchRoute) error
118 takeoverChanged func()
119 logger *slog.Logger
120 }
121
122 func newBotBridgeHub(deps botBridgeDeps) *botBridgeHub {
123 logger := deps.logger
124 if logger == nil {
125 logger = slog.Default()
126 }
127 h := &botBridgeHub{
128 sessions: deps.sessions,
129 approveTab: deps.approveTab,
130 answerTab: deps.answerTab,
131 notify: deps.notify,
132 drive: deps.drive,
133 announce: deps.announce,
134 persistWatchers: deps.persistWatchers,
135 takeoverChanged: deps.takeoverChanged,
136 logger: logger.With("component", "bot_bridge"),
137 watchers: make(map[string]bot.DesktopWatchRoute),
138 pending: make(map[string]desktopPendingPrompt),
139 takeovers: make(map[string]bot.DesktopWatchRoute),
140 takeoverTabs: make(map[string]string),
141 queue: make(chan desktopBridgeNotification, botBridgeQueueSize),
142 closed: make(chan struct{}),
143 }
144 go h.run()
145 return h
146 }
147
148 // Close 停掉通知 worker goroutine。幂等;App 关停时调用。
149 func (h *botBridgeHub) Close() {
150 h.closeOnce.Do(func() { close(h.closed) })
151 }
152
153 // observe 接收某个桌面会话的一条事件。在 controller 事件 goroutine 上运行,
154 // 只做内存记账和入队,不做任何阻塞调用。
155 func (h *botBridgeHub) observe(tabID string, e event.Event) {
156 switch e.Kind {
157 case event.ApprovalRequest:
158 h.mu.Lock()
159 h.rememberPendingLocked(e.Approval.ID, desktopPendingPrompt{
160 tabID: tabID,
161 kind: "approval",
162 tool: e.Approval.Tool,
163 subject: truncateForBridge(e.Approval.Subject, botBridgeSubjectLimit),
164 })
165 watching := len(h.watchers) > 0
166 h.mu.Unlock()
167 if watching && !e.Replayed {
168 h.enqueue(h.approvalNotification(tabID, e.Approval))
169 }
170 case event.AskRequest:
171 h.mu.Lock()
172 h.rememberPendingLocked(e.Ask.ID, desktopPendingPrompt{
173 tabID: tabID,
174 kind: "ask",
175 questions: e.Ask.Questions,
176 })
177 watching := len(h.watchers) > 0
178 h.mu.Unlock()
179 if watching && !e.Replayed {
180 h.enqueue(h.askNotification(tabID, e.Ask))
181 }
182 case event.TurnDone:
183 h.mu.Lock()
184 for id, p := range h.pending {
185 if p.tabID == tabID {
186 delete(h.pending, id)
187 }
188 }
189 watching := len(h.watchers) > 0
190 h.mu.Unlock()
191 if !watching {
192 return
193 }
194 if e.Err != nil && strings.Contains(e.Err.Error(), "context canceled") {
195 // 桌面端主动停止的任务不推送,避免正常操作变成噪音。
196 return
197 }
198 h.enqueue(h.turnDoneNotification(tabID, e))
199 }
200 }
201
202 // rememberPendingLocked 记录待处理项;容量兜底防泄漏(正常路径 TurnDone 会清理)。
203 func (h *botBridgeHub) rememberPendingLocked(id string, p desktopPendingPrompt) {
204 if strings.TrimSpace(id) == "" {
205 return
206 }
207 if len(h.pending) >= botBridgePendingLimit {
208 h.pending = make(map[string]desktopPendingPrompt)
209 }
210 h.pending[id] = p
211 }
212
213 func (h *botBridgeHub) enqueue(n desktopBridgeNotification) {
214 select {
215 case h.queue <- n:
216 default:
217 h.logger.Warn("desktop bridge notification queue full; dropping")
218 }
219 }
220
221 func (h *botBridgeHub) run() {
222 for {
223 select {
224 case <-h.closed:
225 return
226 case n := <-h.queue:
227 h.deliver(n)
228 }
229 }
230 }
231
232 func (h *botBridgeHub) deliver(n desktopBridgeNotification) {
233 h.mu.Lock()
234 var routes []bot.DesktopWatchRoute
235 if n.route != nil {
236 routes = []bot.DesktopWatchRoute{*n.route}
237 } else {
238 routes = h.watcherRoutesLocked()
239 }
240 notify := h.notify
241 h.mu.Unlock()
242 if notify == nil || len(routes) == 0 {
243 return
244 }
245 // Fan out per route: a single slow/hung connection must not hold the queue
246 // worker for its full timeout and back-pressure everyone else's approvals.
247 var wg sync.WaitGroup
248 for _, route := range routes {
249 wg.Add(1)
250 go func(route bot.DesktopWatchRoute) {
251 defer wg.Done()
252 msg := bot.OutboundMessage{
253 ChatID: route.ChatID,
254 ChatType: route.ChatType,
255 }
256 if n.text != nil {
257 msg.Text = n.text(route)
258 }
259 if n.card != nil {
260 msg.Card = n.card(route)
261 }
262 ctx, cancel := context.WithTimeout(context.Background(), botBridgeSendTimeout)
263 defer cancel()
264 if _, err := notify(ctx, route.ConnectionID, route.Domain, msg); err != nil {
265 h.logger.Warn("desktop bridge notification send failed", "platform", route.Platform, "err", err)
266 }
267 }(route)
268 }
269 wg.Wait()
270 }
271
272 // tabLabel 把 tabID 解析成人类可读的会话名。
273 func (h *botBridgeHub) tabLabel(tabID string) string {
274 if s, ok := h.sessionByTabID(tabID); ok {
275 if label := strings.TrimSpace(s.Label); label != "" {
276 return label
277 }
278 if title := strings.TrimSpace(s.Topic); title != "" {
279 return title
280 }
281 }
282 return "(未命名会话)"
283 }
284
285 func (h *botBridgeHub) sessionByTabID(tabID string) (bot.DesktopSessionInfo, bool) {
286 if h.sessions == nil {
287 return bot.DesktopSessionInfo{}, false
288 }
289 for _, s := range h.sessions() {
290 if s.TabID == tabID {
291 return s, true
292 }
293 }
294 return bot.DesktopSessionInfo{}, false
295 }
296
297 func (h *botBridgeHub) approvalNotification(tabID string, approval event.Approval) desktopBridgeNotification {
298 label := h.tabLabel(tabID)
299 // The approval subject is the pending command line; only reveal it in a
300 // private chat. In a shared chat show the tool name and point the operator
301 // to the desktop / a DM instead of leaking the command to the whole group.
302 subjectFor := func(route bot.DesktopWatchRoute) string {
303 if isSharedChat(route.ChatType) {
304 return "(命令详情仅在桌面端或私聊显示)"
305 }
306 return truncateForBridge(approval.Subject, botBridgeSubjectLimit)
307 }
308 return desktopBridgeNotification{
309 text: func(route bot.DesktopWatchRoute) string {
310 return fmt.Sprintf("⚠️ 桌面会话「%s」需要批准操作\n工具: %s\n操作: %s\n\nID: `%s`\n用 /desktop approve %s 批准,/desktop deny %s 拒绝。桌面端先处理则以先到者为准。",
311 label, approval.Tool, subjectFor(route), approval.ID, approval.ID, approval.ID)
312 },
313 card: func(route bot.DesktopWatchRoute) *bot.InteractiveCard {
314 return &bot.InteractiveCard{
315 Header: "桌面会话需要批准",
316 Elements: []bot.InteractiveCardElement{
317 {Tag: "markdown", Content: fmt.Sprintf("**会话**: %s\n\n**工具**: %s\n\n**操作**: %s\n\nID: `%s`", label, approval.Tool, subjectFor(route), approval.ID)},
318 {Tag: "action", Extra: map[string]any{
319 "actions": []map[string]any{
320 desktopCardButton("允许一次", "primary", "/desktop approve "+approval.ID, route),
321 desktopCardButton("拒绝", "danger", "/desktop deny "+approval.ID, route),
322 },
323 }},
324 },
325 }
326 },
327 }
328 }
329
330 func (h *botBridgeHub) askNotification(tabID string, ask event.Ask) desktopBridgeNotification {
331 label := h.tabLabel(tabID)
332 var b strings.Builder
333 fmt.Fprintf(&b, "❓ 桌面会话「%s」在等待回答:\n", label)
334 for i, q := range ask.Questions {
335 fmt.Fprintf(&b, "\n**%d. %s**\n", i+1, truncateForBridge(q.Prompt, botBridgePromptPreview))
336 for j, opt := range q.Options {
337 fmt.Fprintf(&b, " %d. %s\n", j+1, opt.Label)
338 }
339 }
340 fmt.Fprintf(&b, "\nID: `%s`\n用 /desktop answer %s <选项编号或文本> 回答;桌面端先处理则以先到者为准。", ask.ID, ask.ID)
341 privateText := b.String()
342 sharedText := fmt.Sprintf("❓ 桌面会话「%s」正在等待回答(问题详情仅在桌面端或私聊显示)。\n\nID: `%s`", label, ask.ID)
343 textFor := func(route bot.DesktopWatchRoute) string {
344 if isSharedChat(route.ChatType) {
345 return sharedText
346 }
347 return privateText
348 }
349
350 var card func(route bot.DesktopWatchRoute) *bot.InteractiveCard
351 if len(ask.Questions) == 1 && len(ask.Questions[0].Options) > 0 {
352 options := ask.Questions[0].Options
353 card = func(route bot.DesktopWatchRoute) *bot.InteractiveCard {
354 if isSharedChat(route.ChatType) {
355 return nil
356 }
357 actions := make([]map[string]any, 0, len(options))
358 for i, opt := range options {
359 optLabel := strings.TrimSpace(opt.Label)
360 if optLabel == "" {
361 optLabel = fmt.Sprintf("选项 %d", i+1)
362 }
363 actions = append(actions, desktopCardButton(optLabel, "primary", fmt.Sprintf("/desktop answer %s %d", ask.ID, i+1), route))
364 }
365 return &bot.InteractiveCard{
366 Header: "桌面会话在等待回答",
367 Elements: []bot.InteractiveCardElement{
368 {Tag: "markdown", Content: privateText},
369 {Tag: "action", Extra: map[string]any{"actions": actions}},
370 },
371 }
372 }
373 }
374 return desktopBridgeNotification{text: textFor, card: card}
375 }
376
377 func (h *botBridgeHub) turnDoneNotification(tabID string, e event.Event) desktopBridgeNotification {
378 label := h.tabLabel(tabID)
379 if e.Outcome == event.TurnOutcomeIncompleteRead {
380 return desktopBridgeNotification{text: constText(fmt.Sprintf("⏸️ 桌面会话「%s」的读取任务尚未完成,已保留当前结果。请补充读取范围后继续。", label))}
381 }
382 if e.Outcome == event.TurnOutcomeRecoveryPaused {
383 return desktopBridgeNotification{text: constText(fmt.Sprintf(
384 "⏸️ 桌面会话「%s」已暂停自动重试。已完成的工作会保留;发送“继续”即可开始新一轮,也可以补充要求调整方向。",
385 label,
386 ))}
387 }
388 if e.Outcome == event.TurnOutcomeCompletionUncertain {
389 return desktopBridgeNotification{text: constText(fmt.Sprintf(
390 "⏸️ 桌面会话「%s」本轮完成状态未确认。当前结果和已完成工作均已保留;发送“继续”可接着完成,也可以补充说明需要调整的内容。",
391 label,
392 ))}
393 }
394 if e.Err != nil {
395 // Error text can contain paths/tokens; only detail it in a private chat.
396 return desktopBridgeNotification{text: func(route bot.DesktopWatchRoute) string {
397 if isSharedChat(route.ChatType) {
398 return fmt.Sprintf("❌ 桌面会话「%s」任务出错(详情见桌面端或私聊)。", label)
399 }
400 return fmt.Sprintf("❌ 桌面会话「%s」任务出错: %s", label, truncateForBridge(e.Err.Error(), botBridgeErrTextLimit))
401 }}
402 }
403 return desktopBridgeNotification{text: constText(fmt.Sprintf("✅ 桌面会话「%s」任务完成。", label))}
404 }
405
406 func desktopCardButton(label, style, command string, route bot.DesktopWatchRoute) map[string]any {
407 return map[string]any{
408 "tag": "button",
409 "text": map[string]string{"tag": "plain_text", "content": label},
410 "type": style,
411 "value": map[string]string{
412 "command": command,
413 "chat_type": string(route.ChatType),
414 },
415 }
416 }
417
418 func truncateForBridge(s string, limit int) string {
419 s = strings.TrimSpace(s)
420 runes := []rune(s)
421 if len(runes) <= limit {
422 return s
423 }
424 return string(runes[:limit]) + "…"
425 }
426
427 // bot.DesktopBridge 实现
428
429 func (h *botBridgeHub) Sessions() []bot.DesktopSessionInfo {
430 if h.sessions == nil {
431 return nil
432 }
433 sessions := h.sessions()
434 h.mu.Lock()
435 byTab := make(map[string][]bot.DesktopPendingInfo, len(h.pending))
436 for id, p := range h.pending {
437 byTab[p.tabID] = append(byTab[p.tabID], bot.DesktopPendingInfo{ID: id, Kind: p.kind, Tool: p.tool})
438 }
439 h.mu.Unlock()
440 for i := range sessions {
441 if pend := byTab[sessions[i].TabID]; len(pend) > 0 {
442 sort.Slice(pend, func(a, b int) bool { return pend[a].ID < pend[b].ID })
443 sessions[i].Pending = pend
444 }
445 }
446 return sessions
447 }
448
449 func (h *botBridgeHub) SetWatch(route bot.DesktopWatchRoute, enable bool) error {
450 h.mu.Lock()
451 if enable {
452 h.watchers[route.Key()] = route
453 } else {
454 delete(h.watchers, route.Key())
455 }
456 h.watchSeq++
457 h.watchPersistDirty = true
458 seq := h.watchSeq
459 routes := h.watcherRoutesLocked()
460 persist := h.persistWatchers
461 h.mu.Unlock()
462 if persist == nil {
463 return nil
464 }
465 // Serialize persists and drop stale ones: two concurrent SetWatch calls
466 // (different connections) compute snapshots under h.mu but write config
467 // outside it, so their writes could otherwise reorder and let an older
468 // snapshot clobber a newer one, silently losing a subscription.
469 h.persistMu.Lock()
470 defer h.persistMu.Unlock()
471 if seq <= h.lastPersistSeq {
472 return nil
473 }
474 if err := persist(routes); err != nil {
475 return err
476 }
477 h.lastPersistSeq = seq
478 h.mu.Lock()
479 if h.watchSeq == seq {
480 h.watchPersistDirty = false
481 }
482 h.mu.Unlock()
483 return nil
484 }
485
486 func (h *botBridgeHub) watcherVersion() uint64 {
487 h.mu.Lock()
488 defer h.mu.Unlock()
489 return h.watchSeq
490 }
491
492 // seedWatchers applies a config snapshot only if no watch command changed the
493 // runtime after the config read began. Fresh external config edits still apply;
494 // stale refreshes and failed local persists do not erase newer runtime state.
495 func (h *botBridgeHub) seedWatchers(routes []bot.DesktopWatchRoute, expectedSeq uint64) {
496 h.persistMu.Lock()
497 defer h.persistMu.Unlock()
498 h.mu.Lock()
499 defer h.mu.Unlock()
500 if h.watchSeq != expectedSeq || h.watchPersistDirty {
501 return
502 }
503 h.watchers = make(map[string]bot.DesktopWatchRoute, len(routes))
504 for _, r := range routes {
505 if strings.TrimSpace(r.ChatID) == "" {
506 continue
507 }
508 h.watchers[r.Key()] = r
509 }
510 }
511
512 func (h *botBridgeHub) watcherRoutesLocked() []bot.DesktopWatchRoute {
513 routes := make([]bot.DesktopWatchRoute, 0, len(h.watchers))
514 for _, r := range h.watchers {
515 routes = append(routes, r)
516 }
517 sort.Slice(routes, func(i, j int) bool { return routes[i].Key() < routes[j].Key() })
518 return routes
519 }
520
521 func (h *botBridgeHub) Watching(route bot.DesktopWatchRoute) bool {
522 h.mu.Lock()
523 defer h.mu.Unlock()
524 _, ok := h.watchers[route.Key()]
525 return ok
526 }
527
528 func (h *botBridgeHub) Approve(approvalID string, allow bool) (string, error) {
529 approvalID = strings.TrimSpace(approvalID)
530 h.mu.Lock()
531 p, ok := h.pending[approvalID]
532 if ok && p.kind == "approval" {
533 delete(h.pending, approvalID)
534 }
535 h.mu.Unlock()
536 if !ok || p.kind != "approval" {
537 return "", fmt.Errorf("未找到待处理的审批 %s(可能已在桌面端处理或已超时)。用 /desktop status 查看当前会话。", approvalID)
538 }
539 if h.approveTab == nil {
540 return "", fmt.Errorf("桌面端审批通道不可用。")
541 }
542 h.approveTab(p.tabID, approvalID, allow, false, false)
543 action := "批准"
544 if !allow {
545 action = "拒绝"
546 }
547 return fmt.Sprintf("已提交%s「%s」的操作(%s)。桌面端若已先处理,以先到者为准。", action, h.tabLabel(p.tabID), p.tool), nil
548 }
549
550 func (h *botBridgeHub) AskQuestions(askID string) ([]event.AskQuestion, bool) {
551 h.mu.Lock()
552 defer h.mu.Unlock()
553 p, ok := h.pending[strings.TrimSpace(askID)]
554 if !ok || p.kind != "ask" {
555 return nil, false
556 }
557 return p.questions, true
558 }
559
560 func (h *botBridgeHub) Answer(askID string, answers []event.AskAnswer) (string, error) {
561 askID = strings.TrimSpace(askID)
562 h.mu.Lock()
563 p, ok := h.pending[askID]
564 if ok && p.kind == "ask" {
565 delete(h.pending, askID)
566 }
567 h.mu.Unlock()
568 if !ok || p.kind != "ask" {
569 return "", fmt.Errorf("未找到待回答的提问 %s(可能已在桌面端回答或已超时)。", askID)
570 }
571 if h.answerTab == nil {
572 return "", fmt.Errorf("桌面端问答通道不可用。")
573 }
574 out := make([]QuestionAnswer, 0, len(answers))
575 for _, an := range answers {
576 out = append(out, QuestionAnswer{QuestionID: an.QuestionID, Selected: an.Selected})
577 }
578 h.answerTab(p.tabID, askID, out)
579 return fmt.Sprintf("已提交「%s」的回答。桌面端若已先处理,以先到者为准。", h.tabLabel(p.tabID)), nil
580 }
581
582 // 显式接管
583
584 func (h *botBridgeHub) Takeover(route bot.DesktopWatchRoute, tabID string) (string, error) {
585 tabID = strings.TrimSpace(tabID)
586 // DM only. In a group the binding is keyed on the group chat, so after an
587 // admin takes over, ANY allowlisted member's plain message would be diverted
588 // to drive the session — a privilege escalation past the admin gate that
589 // establishes the takeover. Restricting to DM keeps the driver identical to
590 // the operator who established it.
591 if route.ChatType != bot.ChatDM {
592 return "", fmt.Errorf("接管仅支持私聊:在群里接管会让其他成员也能驱动你的桌面会话。请在与 bot 的私聊中接管。")
593 }
594 session, ok := h.sessionByTabID(tabID)
595 if !ok {
596 return "", fmt.Errorf("未找到会话 %s。用 /desktop status 查看可接管的会话。", tabID)
597 }
598 if session.Detached {
599 return "", fmt.Errorf("会话「%s」在后台运行,暂不支持接管;请先在桌面端打开它。", h.tabLabel(tabID))
600 }
601 h.mu.Lock()
602 if holder, held := h.takeovers[tabID]; held && holder.Key() != route.Key() {
603 h.mu.Unlock()
604 return "", fmt.Errorf("会话「%s」已被另一个聊天接管。", h.tabLabel(tabID))
605 }
606 // 同一聊天换目标:先解除旧绑定,并记下旧 tab 以便公告解除。
607 released := ""
608 if prev, ok := h.takeoverTabs[route.Key()]; ok && prev != tabID {
609 delete(h.takeovers, prev)
610 released = prev
611 }
612 h.takeovers[tabID] = route
613 h.takeoverTabs[route.Key()] = tabID
614 announce := h.announce
615 changed := h.takeoverChanged
616 h.mu.Unlock()
617 if announce != nil {
618 if released != "" {
619 announce(released, "IM 远程接管已解除(接管方切换到了另一个会话)。")
620 }
621 announce(tabID, "此会话已被 IM 远程接管(bot 管理员)。在此本地发送任意消息即可收回控制。")
622 }
623 if changed != nil {
624 changed()
625 }
626 label := h.tabLabel(tabID)
627 return fmt.Sprintf("已接管「%s」。现在直接发消息即可驱动它,输出会流回本聊天;/desktop release 解除接管。桌面端本地发言会自动收回控制。", label), nil
628 }
629
630 func (h *botBridgeHub) Release(route bot.DesktopWatchRoute) (string, error) {
631 h.mu.Lock()
632 tabID, ok := h.takeoverTabs[route.Key()]
633 if ok {
634 delete(h.takeoverTabs, route.Key())
635 delete(h.takeovers, tabID)
636 }
637 announce := h.announce
638 changed := h.takeoverChanged
639 h.mu.Unlock()
640 if !ok {
641 return "", fmt.Errorf("本聊天当前没有接管任何桌面会话。")
642 }
643 if announce != nil {
644 announce(tabID, "IM 远程接管已解除。")
645 }
646 if changed != nil {
647 changed()
648 }
649 return fmt.Sprintf("已解除对「%s」的接管。", h.tabLabel(tabID)), nil
650 }
651
652 func (h *botBridgeHub) TakeoverTab(route bot.DesktopWatchRoute) string {
653 h.mu.Lock()
654 defer h.mu.Unlock()
655 return h.takeoverTabs[route.Key()]
656 }
657
658 func (h *botBridgeHub) DriveInput(route bot.DesktopWatchRoute, text string) (string, error) {
659 h.mu.Lock()
660 tabID := h.takeoverTabs[route.Key()]
661 h.mu.Unlock()
662 if tabID == "" {
663 return "", fmt.Errorf("本聊天没有接管任何桌面会话。")
664 }
665 session, ok := h.sessionByTabID(tabID)
666 if !ok || session.Detached {
667 // 会话被关闭或转入后台:自动解除绑定,避免消息黑洞。
668 h.mu.Lock()
669 delete(h.takeoverTabs, route.Key())
670 delete(h.takeovers, tabID)
671 h.mu.Unlock()
672 if changed := h.takeoverChanged; changed != nil {
673 changed()
674 }
675 return "", fmt.Errorf("被接管的会话已关闭或转入后台,接管已自动解除。")
676 }
677 if session.Running {
678 return "", h.busyError(tabID)
679 }
680 if h.drive == nil {
681 return "", fmt.Errorf("桌面端驱动通道不可用。")
682 }
683 if err := h.drive(tabID, text, route); err != nil {
684 if errors.Is(err, errDriveBusy) {
685 return "", h.busyError(tabID)
686 }
687 return "", fmt.Errorf("驱动失败: %w", err)
688 }
689 return "", nil
690 }
691
692 func (h *botBridgeHub) busyError(tabID string) error {
693 return fmt.Errorf("会话「%s」正在执行中,等它完成后再发;或用 /desktop watch on 订阅完成通知。", h.tabLabel(tabID))
694 }
695
696 // reclaimFromDesktop 在桌面用户本地提交输入时收回控制权:解除绑定并通知
697 // 远端聊天。由 App.SubmitToTab 调用(bridge 自己的驱动不走这条路)。
698 func (h *botBridgeHub) reclaimFromDesktop(tabID string) {
699 h.mu.Lock()
700 route, ok := h.takeovers[tabID]
701 if ok {
702 delete(h.takeovers, tabID)
703 delete(h.takeoverTabs, route.Key())
704 }
705 notify := h.notify
706 changed := h.takeoverChanged
707 h.mu.Unlock()
708 if !ok {
709 return
710 }
711 if changed != nil {
712 changed()
713 }
714 if notify == nil {
715 return
716 }
717 label := h.tabLabel(tabID)
718 // 直接入通知队列(不依赖 watch 订阅):接管者必须知道控制权没了。
719 h.enqueue(desktopBridgeNotification{
720 text: constText(fmt.Sprintf("🔓 桌面端已收回会话「%s」的控制权,接管已解除。", label)),
721 route: &route,
722 })
723 }
724
725 // remoteControlledTabs 返回当前被接管的 tabID 集合(TabMeta 标记用)。
726 func (h *botBridgeHub) remoteControlledTabs() map[string]bool {
727 h.mu.Lock()
728 defer h.mu.Unlock()
729 if len(h.takeovers) == 0 {
730 return nil
731 }
732 out := make(map[string]bool, len(h.takeovers))
733 for tabID := range h.takeovers {
734 out[tabID] = true
735 }
736 return out
737 }
738
738 lines GO