返回 DeepSeek-Reasonix
gateway_session_rebuild.go
根目录 / internal / bot / gateway_session_rebuild.go
1 package bot
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "os"
8 "strings"
9 "time"
10
11 "reasonix/internal/agent"
12 "reasonix/internal/boot"
13 "reasonix/internal/control"
14 "reasonix/internal/event"
15 "reasonix/internal/secrets"
16 "reasonix/internal/session"
17 )
18
19 type builtBotSession struct {
20 state *sessionState
21 reusedLease bool
22 reusedRuntime bool
23 }
24
25 func botRuntimeSwitchBusyText() string {
26 return "当前会话仍有正在运行、等待确认或后台执行的任务。请先完成或停止这些任务,再切换项目或 attach 会话。"
27 }
28
29 func botRuntimeSwitchFailedText(action string) string {
30 return action + "失败,当前会话保持不变。请检查配置后重试。"
31 }
32
33 func (gw *BotGateway) buildBotController(ctx context.Context, opts boot.Options) (*control.Controller, error) {
34 if opts.SessionService == nil {
35 var err error
36 opts.SessionService, err = gw.botSessionService(opts.SessionDir)
37 if err != nil {
38 return nil, err
39 }
40 opts.SessionHostID = "local"
41 }
42 if gw.buildController != nil {
43 return gw.buildController(ctx, opts)
44 }
45 return boot.Build(ctx, opts)
46 }
47
48 func (gw *BotGateway) botSessionService(sessionDir string) (*session.Service, error) {
49 root := session.RootForLegacyDir(sessionDir)
50 if root == "" {
51 return nil, nil
52 }
53 if gw.cfg.SessionServiceForRoot != nil {
54 return gw.cfg.SessionServiceForRoot(root)
55 }
56 gw.sessionServicesMu.Lock()
57 defer gw.sessionServicesMu.Unlock()
58 if gw.sessionServices == nil {
59 gw.sessionServices = make(map[string]*session.Service)
60 }
61 if service := gw.sessionServices[root]; service != nil {
62 return service, nil
63 }
64 service, err := session.NewService("local", session.NewFilesystemPersistence(root))
65 if err != nil {
66 return nil, err
67 }
68 gw.sessionServices[root] = service
69 return service, nil
70 }
71
72 // buildSessionState prepares a complete replacement without publishing it.
73 // When the transcript path is unchanged, the candidate reuses the old keeper
74 // so the session lease never has an unowned window during a model/profile swap.
75 func (gw *BotGateway) buildSessionState(ctx context.Context, key string, msg InboundMessage, profile sessionRuntimeProfile, previous *sessionState) (*builtBotSession, error) {
76 leases := control.NewSessionLeaseKeeper()
77 reusedLease := false
78 if previous != nil && previous.leases != nil {
79 heldPath := agent.CanonicalSessionPath(previous.leases.HeldPath())
80 if heldPath != "" && heldPath == agent.CanonicalSessionPath(profile.sessionPath) {
81 leases = previous.leases
82 reusedLease = true
83 }
84 }
85
86 sessionSink := &sessionEventSink{}
87 state := &sessionState{
88 sink: sessionSink,
89 leases: leases,
90 platform: msg.Platform,
91 connectionID: strings.TrimSpace(msg.ConnectionID),
92 model: profile.model,
93 workspaceRoot: profile.workspaceRoot,
94 toolApprovalMode: profile.toolApprovalMode,
95 sessionPath: profile.sessionPath,
96 sessionRef: profile.sessionRef,
97 pendingAsks: make(map[string][]event.AskQuestion),
98 createdAt: time.Now(),
99 lastActive: time.Now(),
100 }
101 state.onSessionTransition = gw.botSessionTransitionHandler(key, msg, state)
102 buildOptions := boot.Options{
103 Model: profile.model,
104 MaxSteps: gw.cfg.MaxSteps,
105 MaxStepsKey: "bot.max_steps",
106 RequireKey: true,
107 Sink: sessionSink,
108 StatsSource: "bot",
109 WorkspaceRoot: profile.workspaceRoot,
110 SessionDir: botSessionDir(profile.workspaceRoot),
111 ApprovalTimeout: gw.approvalTimeout(),
112 OnSessionRecovered: gw.botSessionRecoveredHandler(key, msg, state),
113 OnSessionTransition: state.onSessionTransition,
114 }
115 reusedRuntime := false
116 if previous != nil {
117 if binding, ok := previous.ctrl.(interface {
118 SessionBinding() (*session.Service, *session.Runtime, bool)
119 }); ok {
120 if service, runtime, bound := binding.SessionBinding(); bound {
121 reusedRuntime = profile.sessionPath == "" && (profile.sessionRef.SessionID == "" || profile.sessionRef.SessionID == runtime.Ref().SessionID)
122 if reusedRuntime {
123 buildOptions.SessionService = service
124 buildOptions.SessionRuntime = runtime
125 buildOptions.SessionHostID = runtime.Ref().HostID
126 }
127 }
128 }
129 }
130 ctrl, err := gw.buildBotController(ctx, buildOptions)
131 if err != nil {
132 if !reusedLease {
133 leases.Release()
134 }
135 return nil, err
136 }
137 state.ctrl = ctrl
138 fail := func(buildErr error) (*builtBotSession, error) {
139 if reusedRuntime {
140 ctrl.ReleaseResources()
141 } else {
142 ctrl.Close()
143 }
144 if reusedLease {
145 if restoreErr := bindBotSessionWriteAuthority(previous); restoreErr != nil {
146 gw.logger.Error("restore bot session write authority failed", "err", secrets.RedactError(restoreErr))
147 }
148 } else {
149 leases.Release()
150 }
151 return nil, buildErr
152 }
153 if identity, ok := any(ctrl).(control.IdentityLifecycle); ok && identity.UsesExclusiveSession() {
154 ref, bindErr := bindBotSessionIdentity(ctx, identity, profile, msg)
155 if bindErr != nil {
156 if (profile.sessionRefOptional || profile.sessionPathOptional) && !reusedRuntime {
157 gw.logger.Warn("mapped bot session unavailable; starting fresh", "err", bindErr)
158 ref, bindErr = identity.BindFreshSession(ctx, "")
159 state.mappingDegraded = bindErr == nil
160 }
161 if bindErr != nil {
162 return fail(bindErr)
163 }
164 }
165 state.sessionRef = ref
166 state.sessionPath = ""
167 ctrl.EnableInteractiveApproval()
168 ctrl.SetToolApprovalMode(profile.toolApprovalMode)
169 return &builtBotSession{state: state, reusedRuntime: reusedRuntime}, nil
170 }
171
172 if profile.sessionPath != "" {
173 degrade := func(reason string, loadErr error) bool {
174 if !profile.sessionPathOptional {
175 return false
176 }
177 gw.logger.Warn("mapped bot session unavailable; starting fresh", "reason", reason, "session_path", profile.sessionPath, "err", loadErr)
178 profile.sessionPath = ""
179 state.sessionPath = ""
180 state.mappingDegraded = true
181 return true
182 }
183 if err := leases.Rebind(profile.sessionPath); err != nil {
184 if !degrade("lease held elsewhere", err) {
185 return fail(fmt.Errorf("attached bot session is in use: %w", err))
186 }
187 } else if loaded, err := agent.LoadSession(profile.sessionPath); err != nil {
188 if os.IsNotExist(err) && profile.sessionPathOptional {
189 ctrl.SetSessionPath(profile.sessionPath)
190 } else if !degrade("load failed", err) {
191 return fail(fmt.Errorf("load attached bot session: %w", err))
192 }
193 } else {
194 ctrl.Resume(loaded, profile.sessionPath)
195 }
196 }
197 ctrl.EnableInteractiveApproval()
198 ctrl.SetToolApprovalMode(profile.toolApprovalMode)
199 ctrl.EnsureSessionPath()
200 if reusedLease && agent.CanonicalSessionPath(ctrl.SessionPath()) != agent.CanonicalSessionPath(leases.HeldPath()) {
201 return fail(errors.New("replacement session path changed while reusing the current lease"))
202 }
203 if err := rebindBotSessionWriteAuthority(state, ctrl.SessionPath()); err != nil {
204 return fail(fmt.Errorf("bind bot session write authority: %w", err))
205 }
206 return &builtBotSession{state: state, reusedLease: reusedLease}, nil
207 }
208
209 func bindBotSessionIdentity(ctx context.Context, identity control.IdentityLifecycle, profile sessionRuntimeProfile, msg InboundMessage) (session.SessionRef, error) {
210 service := identity.SessionService()
211 if service == nil {
212 return session.SessionRef{}, errors.New("bot v3 session service is unavailable")
213 }
214 if current, ok := identity.SessionRef(); ok {
215 if profile.sessionRef.SessionID == "" || profile.sessionRef.SessionID == current.SessionID {
216 return current, nil
217 }
218 }
219 if profile.sessionRef.SessionID != "" {
220 ref := profile.sessionRef
221 if ref.HostID == "" {
222 ref.HostID = service.HostID()
223 }
224 opened, err := identity.OpenSession(ctx, ref)
225 if err == nil {
226 return opened, nil
227 }
228 if !errors.Is(err, session.ErrSessionNotFound) || !profile.sessionRefOptional {
229 return session.SessionRef{}, err
230 }
231 return identity.BindFreshSession(ctx, ref.SessionID)
232 }
233 if profile.sessionPath != "" {
234 return identity.ContinueLegacySession(ctx, profile.sessionPath, "")
235 }
236 stableID := ""
237 if strings.TrimSpace(msg.ChatID) != "" {
238 stableID = "bot-" + BuildSessionKey(msg.Session())
239 }
240 return identity.BindFreshSession(ctx, stableID)
241 }
242
243 func (gw *BotGateway) discardBuiltSession(built *builtBotSession, previous *sessionState) {
244 if built == nil || built.state == nil {
245 return
246 }
247 if built.state.ctrl != nil {
248 if built.reusedRuntime {
249 if releaser, ok := built.state.ctrl.(interface{ ReleaseResources() }); ok {
250 releaser.ReleaseResources()
251 } else {
252 built.state.ctrl.Close()
253 }
254 } else {
255 built.state.ctrl.Close()
256 }
257 }
258 if built.reusedLease {
259 if previous != nil {
260 previous.lifecycleMu.Lock()
261 retired := previous.retired
262 previous.lifecycleMu.Unlock()
263 if !retired {
264 if err := bindBotSessionWriteAuthority(previous); err != nil {
265 gw.logger.Error("restore bot session write authority failed", "err", secrets.RedactError(err))
266 }
267 }
268 }
269 return
270 }
271 if built.state.leases != nil {
272 built.state.leases.Release()
273 }
274 }
275
276 func (gw *BotGateway) setSessionRuntimeOverride(ctx context.Context, key string, msg InboundMessage, override sessionRuntimeOverride, enabled bool) (bool, error) {
277 if _, ok := parseBotSessionRefTarget(override.sessionPath); !ok {
278 override.sessionPath = canonicalBotPath(override.sessionPath)
279 }
280 override.channel.WorkspaceRoot = canonicalBotPath(override.channel.WorkspaceRoot)
281 profile := gw.sessionProfileForResolvedOverride(msg, override, enabled)
282 var switchErr error
283 switched := gw.sessions.runIfIdle(key, func() bool {
284 gw.mu.Lock()
285 previous := gw.controllers[key]
286 if previous == nil {
287 if enabled {
288 gw.sessionOverrides[key] = override
289 } else {
290 delete(gw.sessionOverrides, key)
291 }
292 gw.mu.Unlock()
293 return true
294 }
295 if previous != nil && botSessionHasActiveWork(previous) {
296 gw.mu.Unlock()
297 return false
298 }
299 if previous != nil && sessionStateMatchesRuntime(previous, profile) {
300 if enabled {
301 gw.sessionOverrides[key] = override
302 } else {
303 delete(gw.sessionOverrides, key)
304 }
305 updateSessionStateRuntime(previous, msg, profile)
306 gw.mu.Unlock()
307 safeBotSetToolApprovalMode(previous.ctrl, profile.toolApprovalMode)
308 return true
309 }
310 gw.mu.Unlock()
311
312 built, err := gw.buildSessionState(ctx, key, msg, profile, previous)
313 if err != nil {
314 switchErr = err
315 gw.logger.Error("bot session runtime switch failed", "err", secrets.RedactError(err))
316 return false
317 }
318
319 gw.mu.Lock()
320 if gw.controllers[key] != previous {
321 gw.mu.Unlock()
322 gw.discardBuiltSession(built, previous)
323 switchErr = errors.New("bot session changed while replacement was building")
324 return false
325 }
326 if enabled {
327 gw.sessionOverrides[key] = override
328 } else {
329 delete(gw.sessionOverrides, key)
330 }
331 gw.controllers[key] = built.state
332 if built.reusedLease && previous != nil {
333 previous.leases = nil
334 }
335 if built.reusedRuntime && previous != nil {
336 previous.releaseRuntimeOnly = true
337 }
338 gw.mu.Unlock()
339 gw.closeSessionState(previous)
340 return true
341 })
342 return switched, switchErr
343 }
344
344 lines GO