返回 DeepSeek-Reasonix
bot_bridge_app.go
根目录 / desktop / bot_bridge_app.go
1 package main
2
3 import (
4 "errors"
5 "fmt"
6 "log/slog"
7 "strings"
8
9 "reasonix/internal/bot"
10 "reasonix/internal/config"
11 "reasonix/internal/control"
12 "reasonix/internal/event"
13 )
14
15 // 本文件是 botBridgeHub 对 App 的全部胶水:会话枚举(含后台 detached)、
16 // 按 tab 寻址的审批/问答/驱动、transcript 公告、订阅持久化。
17
18 func (a *App) newBotBridge() *botBridgeHub {
19 return newBotBridgeHub(botBridgeDeps{
20 sessions: a.bridgeSessions,
21 approveTab: a.bridgeApprove,
22 answerTab: a.bridgeAnswer,
23 notify: a.botRuntime.SendToAdapter,
24 drive: a.bridgeDrive,
25 announce: a.bridgeAnnounce,
26 persistWatchers: a.bridgePersistWatchers,
27 takeoverChanged: a.emitProjectTreeChanged,
28 logger: slog.Default(),
29 })
30 }
31
32 func (a *App) stopBotBridge() {
33 if a.botBridge != nil {
34 a.botBridge.Close()
35 }
36 }
37
38 // bridgeSessions 枚举所有 live 会话:可见 tab 用完整 TabMeta,后台 detached
39 // 会话补一份轻量快照(controller 仍存活,审批/问答仍可路由)。
40 func (a *App) bridgeSessions() []bot.DesktopSessionInfo {
41 tabs := a.ListTabs()
42 out := make([]bot.DesktopSessionInfo, 0, len(tabs)+4)
43 seen := make(map[string]bool, len(tabs))
44 for _, t := range tabs {
45 seen[t.ID] = true
46 out = append(out, bot.DesktopSessionInfo{
47 TabID: t.ID,
48 Label: t.Label,
49 Workspace: t.WorkspaceName,
50 Topic: t.TopicTitle,
51 Ready: t.Ready,
52 Running: t.Running,
53 PendingPrompt: t.PendingPrompt,
54 })
55 }
56 a.mu.RLock()
57 for _, tab := range a.detachedSessions {
58 if tab == nil || seen[tab.ID] {
59 continue
60 }
61 seen[tab.ID] = true
62 out = append(out, bot.DesktopSessionInfo{
63 TabID: tab.ID,
64 Label: tab.TopicTitle,
65 Topic: tab.TopicTitle,
66 Ready: tab.Ctrl != nil,
67 Running: strings.TrimSpace(tab.ActivityStatus) != "",
68 Detached: true,
69 })
70 }
71 a.mu.RUnlock()
72 return out
73 }
74
75 // bridgeCtrlByTabID 解析可见与后台 detached 两张表(区别于 ctrlByTabID:
76 // 那是前端语义,空 tabID 落到活跃 tab,且不看 detached)。
77 func (a *App) bridgeCtrlByTabID(tabID string) control.SessionAPI {
78 a.mu.RLock()
79 defer a.mu.RUnlock()
80 if tab := a.tabByEventSinkIDLocked(tabID); tab != nil {
81 return tab.Ctrl
82 }
83 return nil
84 }
85
86 func (a *App) bridgeApprove(tabID, id string, allow, session, persist bool) {
87 if ctrl := a.bridgeCtrlByTabID(tabID); ctrl != nil {
88 ctrl.Approve(id, allow, session, persist)
89 }
90 }
91
92 func (a *App) bridgeAnswer(tabID, id string, answers []QuestionAnswer) {
93 ctrl := a.bridgeCtrlByTabID(tabID)
94 if ctrl == nil {
95 return
96 }
97 out := make([]event.AskAnswer, len(answers))
98 for i, an := range answers {
99 out[i] = event.AskAnswer{QuestionID: an.QuestionID, Selected: an.Selected}
100 }
101 ctrl.AnswerQuestion(id, out)
102 }
103
104 // bridgeAnnounce 往会话 transcript 发一条 Notice,桌面用户在聊天流里可见。
105 func (a *App) bridgeAnnounce(tabID, text string) {
106 a.mu.RLock()
107 tab := a.tabByEventSinkIDLocked(tabID)
108 var sink *tabEventSink
109 if tab != nil {
110 sink = tab.sink
111 }
112 a.mu.RUnlock()
113 if sink == nil {
114 return
115 }
116 sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: text})
117 }
118
119 // bridgeDrive 把远程文本提交为可见 tab 的新 turn,并为这一轮挂上事件转发器,
120 // 让输出流回接管聊天(转发器在 TurnDone 自动卸载)。
121 func (a *App) bridgeDrive(tabID, text string, route bot.DesktopWatchRoute) error {
122 admission, ctrl, err := a.beginTabTurn(tabID, false)
123 if err != nil {
124 if errors.Is(err, control.ErrTurnRunning) {
125 return errDriveBusy
126 }
127 return err
128 }
129 defer admission.abort()
130 tab := admission.tab
131 if tab.sink == nil {
132 return fmt.Errorf("会话事件通道不可用,无法驱动")
133 }
134 // A local submission may have reclaimed the tab while this drive was waiting
135 // for the per-tab admission gate. Revalidate ownership only after the gate is
136 // held, immediately before attaching the route-specific forwarder.
137 if a.botBridge == nil || a.botBridge.TakeoverTab(route) != tabID {
138 return fmt.Errorf("接管已解除,请重新接管会话")
139 }
140 target := botForwardTarget{
141 ConnID: route.ConnectionID,
142 Domain: route.Domain,
143 ChatID: route.ChatID,
144 ChatType: route.ChatType,
145 }
146 generation := tab.sink.SetBotSink(newBotEventForwarder(a.botRuntime, []botForwardTarget{target}))
147 if err := a.ensureTabTopicIndexedForUserTurn(tab); err != nil {
148 tab.sink.clearBotSink(generation)
149 return err
150 }
151 ctrl.SubmitDisplay(text, text)
152 // Confirm the submit actually started a turn. If nothing is running now, the
153 // controller was rotating and the submit no-oped — detach this exact
154 // generation so a later turn's output does not leak.
155 if !admission.finish(ctrl) {
156 tab.sink.clearBotSink(generation)
157 return errDriveBusy
158 }
159 return nil
160 }
161
162 // bridgePersistWatchers 把订阅全集回写用户配置(bot.desktop_watchers),
163 // 桌面重启后由 refreshBotRuntime 重新种子。
164 func (a *App) bridgePersistWatchers(routes []bot.DesktopWatchRoute) error {
165 return a.applyConfigOnly(func(c *config.Config) error {
166 watchers := make([]config.BotDesktopWatcherConfig, 0, len(routes))
167 for _, r := range routes {
168 watchers = append(watchers, config.BotDesktopWatcherConfig{
169 Platform: string(r.Platform),
170 ConnectionID: r.ConnectionID,
171 Domain: r.Domain,
172 ChatType: string(r.ChatType),
173 ChatID: r.ChatID,
174 })
175 }
176 c.Bot.DesktopWatchers = watchers
177 return nil
178 })
179 }
180
181 func bridgeRoutesFromConfig(watchers []config.BotDesktopWatcherConfig) []bot.DesktopWatchRoute {
182 routes := make([]bot.DesktopWatchRoute, 0, len(watchers))
183 for _, w := range watchers {
184 routes = append(routes, bot.DesktopWatchRoute{
185 Platform: bot.Platform(strings.TrimSpace(w.Platform)),
186 ConnectionID: strings.TrimSpace(w.ConnectionID),
187 Domain: strings.TrimSpace(w.Domain),
188 ChatType: bot.ChatType(strings.TrimSpace(w.ChatType)),
189 ChatID: strings.TrimSpace(w.ChatID),
190 })
191 }
192 return routes
193 }
194
194 lines GO