返回 DeepSeek-Reasonix
runtime_state.go
根目录 / desktop / runtime_state.go
1 package main
2
3 import (
4 "reflect"
5 "sort"
6 "strings"
7 "sync"
8
9 "reasonix/internal/control"
10 "reasonix/internal/event"
11 )
12
13 type RuntimeSessionState struct {
14 TabID string `json:"tabId"`
15 Scope string `json:"scope"`
16 WorkspaceRoot string `json:"workspaceRoot"`
17 TopicID string `json:"topicId"`
18 SessionID string `json:"sessionId,omitempty"`
19 SessionPath string `json:"sessionPath"`
20 SessionGeneration uint64 `json:"sessionGeneration"`
21 Open bool `json:"open"`
22 Remote bool `json:"remote"`
23 HostID string `json:"hostId,omitempty"`
24 Freshness string `json:"freshness"`
25 State event.RuntimeStateSnapshot `json:"state"`
26 }
27
28 type RuntimeStateProjection struct {
29 Epoch string `json:"epoch"`
30 Revision uint64 `json:"revision"`
31 Sessions []RuntimeSessionState `json:"sessions"`
32 Topics []ProjectRuntimeTopic `json:"topics"`
33 }
34
35 type desktopRuntimeProjection struct {
36 mu sync.Mutex
37 snapshot RuntimeStateProjection
38 bindings map[localRuntimeBindingKey]localRuntimeBinding
39 publishedEpoch string
40 publishedRevision uint64
41 events runtimeProjectionEvents
42 }
43
44 type localRuntimeUpdate struct {
45 tab *WorkspaceTab
46 ctrl control.SessionAPI
47 state event.RuntimeStateSnapshot
48 generation uint64
49 path string
50 sessionID string
51 }
52
53 type localRuntimeBindingKey struct {
54 key string
55 open bool
56 }
57
58 type localRuntimeBinding struct {
59 key localRuntimeBindingKey
60 tab *WorkspaceTab
61 view RuntimeSessionState
62 ctrl control.SessionAPI
63 catalog catalogRuntimeSnapshot
64 }
65
66 // localRuntimeBindingsLocked copies identity and display metadata as one
67 // binding. App.mu must be held; controller methods must stay outside that lock.
68 func (a *App) localRuntimeBindingsLocked() map[localRuntimeBindingKey]localRuntimeBinding {
69 bindings := make(map[localRuntimeBindingKey]localRuntimeBinding, len(a.tabs)+len(a.detachedSessions))
70 collect := func(key string, tab *WorkspaceTab, open bool) {
71 if tab == nil {
72 return
73 }
74 bindings[localRuntimeBindingKey{key, open}] = localRuntimeBinding{key: localRuntimeBindingKey{key, open}, tab: tab, ctrl: tab.Ctrl,
75 view: RuntimeSessionState{TabID: tab.ID, Scope: tab.Scope, WorkspaceRoot: tab.WorkspaceRoot,
76 TopicID: tab.TopicID, SessionID: tab.SessionID, SessionPath: tab.SessionPath, SessionGeneration: tab.SessionGeneration, Open: open, Freshness: "synced"},
77 catalog: catalogRuntimeSnapshot{tabID: tab.ID, scope: tab.Scope, workspaceRoot: tab.WorkspaceRoot, topicID: tab.TopicID, sessionPath: tab.SessionPath,
78 activity: tab.ActivityStatus, topicTitle: tab.TopicTitle, topicTitleSource: tab.topicTitleSource, open: open}}
79 if tab.SessionID != "" {
80 binding := bindings[localRuntimeBindingKey{key, open}]
81 binding.catalog.sessionPath = sessionRoute(tab.SessionID)
82 bindings[localRuntimeBindingKey{key, open}] = binding
83 }
84 }
85 for key, tab := range a.tabs {
86 collect(key, tab, true)
87 }
88 for key, tab := range a.detachedSessions {
89 collect(key, tab, false)
90 }
91 return bindings
92 }
93
94 func (a *App) sampleLocalRuntimeBindingsWithUpdate(update *localRuntimeUpdate) []localRuntimeBinding {
95 updates := map[*WorkspaceTab]localRuntimeUpdate{}
96 if update != nil {
97 updates[update.tab] = *update
98 }
99 return a.sampleLocalRuntimeBindingsWithUpdates(updates, update != nil)
100 }
101
102 func (a *App) sampleLocalRuntimeBindingsWithUpdates(updates map[*WorkspaceTab]localRuntimeUpdate, incremental bool) []localRuntimeBinding {
103 for {
104 a.mu.RLock()
105 bindings := a.localRuntimeBindingsLocked()
106 a.mu.RUnlock()
107 states := make(map[localRuntimeBindingKey]event.RuntimeStateSnapshot, len(bindings))
108 r := &a.runtimeStateProjection
109 r.mu.Lock()
110 cached := r.bindings
111 r.mu.Unlock()
112 for key, binding := range bindings {
113 if incremental {
114 update, exists := updates[binding.tab]
115 if exists && sameSessionAPI(binding.ctrl, update.ctrl) &&
116 binding.view.SessionGeneration == update.generation && binding.view.SessionPath == update.path && binding.view.SessionID == update.sessionID &&
117 (update.state.SessionID == "" || update.state.SessionID == binding.view.SessionID) {
118 states[key] = update.state
119 continue
120 }
121 if previous, ok := cached[key]; ok && sameLocalRuntimeBinding(previous, binding) {
122 states[key] = previous.view.State
123 continue
124 }
125 }
126 states[key] = controllerRuntimeState(binding.ctrl)
127 }
128 // A controller can rotate its session, be replaced, or move between
129 // open and detached while sampled. Retry the binding set so its state
130 // cannot be published under the previous session identity.
131 a.mu.RLock()
132 current := a.localRuntimeBindingsLocked()
133 valid := len(current) == len(bindings)
134 for key, binding := range bindings {
135 if !sameLocalRuntimeBinding(current[key], binding) {
136 valid = false
137 break
138 }
139 }
140 a.mu.RUnlock()
141 if !valid {
142 continue
143 }
144 result := make([]localRuntimeBinding, 0, len(bindings))
145 for key, binding := range bindings {
146 binding.view.State = states[key]
147 result = append(result, binding)
148 }
149 return result
150 }
151 }
152
153 // sameLocalRuntimeBinding compares the immutable identity and display fields
154 // copied while App.mu was held. reflect.DeepEqual is deliberately unsuitable
155 // here: following tab or controller pointers recursively reads their live
156 // mutex/atomic state and races the controller's runtime-state publisher.
157 func sameLocalRuntimeBinding(current, sampled localRuntimeBinding) bool {
158 return current.tab == sampled.tab &&
159 sameSessionAPI(current.ctrl, sampled.ctrl) &&
160 current.view.TabID == sampled.view.TabID &&
161 current.view.Scope == sampled.view.Scope &&
162 current.view.WorkspaceRoot == sampled.view.WorkspaceRoot &&
163 current.view.TopicID == sampled.view.TopicID &&
164 current.view.SessionID == sampled.view.SessionID &&
165 current.view.SessionPath == sampled.view.SessionPath &&
166 current.view.SessionGeneration == sampled.view.SessionGeneration &&
167 current.view.Open == sampled.view.Open &&
168 current.view.Remote == sampled.view.Remote &&
169 current.view.HostID == sampled.view.HostID &&
170 current.view.Freshness == sampled.view.Freshness &&
171 current.catalog.scope == sampled.catalog.scope &&
172 current.catalog.workspaceRoot == sampled.catalog.workspaceRoot &&
173 current.catalog.topicID == sampled.catalog.topicID &&
174 current.catalog.sessionPath == sampled.catalog.sessionPath &&
175 current.catalog.activity == sampled.catalog.activity &&
176 current.catalog.topicTitle == sampled.catalog.topicTitle &&
177 current.catalog.topicTitleSource == sampled.catalog.topicTitleSource &&
178 current.catalog.open == sampled.catalog.open
179 }
180
181 func sameSessionAPI(current, sampled control.SessionAPI) bool {
182 if current == nil || sampled == nil {
183 return current == nil && sampled == nil
184 }
185 currentValue, sampledValue := reflect.ValueOf(current), reflect.ValueOf(sampled)
186 if currentValue.Type() != sampledValue.Type() {
187 return false
188 }
189 if currentValue.Type().Comparable() {
190 return currentValue.Interface() == sampledValue.Interface()
191 }
192 return false
193 }
194
195 func controllerRuntimeState(ctrl control.SessionAPI) event.RuntimeStateSnapshot {
196 if reader, ok := ctrl.(control.PublishedRuntimeStateReader); ok {
197 return reader.PublishedRuntimeStateSnapshot()
198 }
199 // Compatibility for embedders without the committed observation boundary.
200 if reader, ok := ctrl.(control.RuntimeStateReader); ok {
201 return reader.RuntimeStateSnapshot()
202 }
203 if ctrl == nil {
204 return event.RuntimeStateSnapshot{Phase: "idle"}
205 }
206 legacy := ctrl.RuntimeStatus()
207 phase := "idle"
208 if legacy.Running {
209 phase = "executing"
210 }
211 return event.RuntimeStateSnapshot{Phase: phase, Running: legacy.Running, PendingPrompt: legacy.PendingPrompt,
212 BackgroundJobs: legacy.BackgroundJobs, Cancellable: legacy.Cancellable, CancelRequested: legacy.CancelRequested,
213 TurnID: legacy.TurnID, TurnStatus: legacy.Status, TurnEventSeq: legacy.TurnEventSeq}
214 }
215
216 func runtimeDisplayStatus(state event.RuntimeStateSnapshot, result string) string {
217 switch {
218 case state.Phase == "finishing":
219 return "finishing"
220 case state.CancelRequested:
221 return "cancelling"
222 case state.PendingPrompt:
223 return topicStatusWaitingConfirmation
224 case state.Phase == "executing":
225 if state.Activity == "streaming" {
226 return topicStatusStreaming
227 }
228 return topicStatusThinking
229 case state.BackgroundJobs > 0:
230 return topicStatusBackgroundJob
231 }
232 if result == topicStatusError || result == topicStatusPaused || result == topicStatusAwaitingDelivery {
233 return result
234 }
235 return ""
236 }
237
238 func catalogControllerStatus(ctrl control.SessionAPI, activity string) (string, bool) {
239 state := controllerRuntimeState(ctrl)
240 return catalogStateStatus(state, activity)
241 }
242
243 func catalogStateStatus(state event.RuntimeStateSnapshot, activity string) (string, bool) {
244 if state.SchemaVersion == 1 {
245 return runtimeDisplayStatus(state, activity), state.ActiveWork()
246 }
247 legacy := control.RuntimeStatus{Running: state.Running, PendingPrompt: state.PendingPrompt, BackgroundJobs: state.BackgroundJobs}
248 status := catalogRuntimeStatus(activity, legacy)
249 return status, status != "" || state.ActiveWork()
250 }
251
252 // GetRuntimeStateSnapshot reads committed controller snapshots after copying
253 // bindings off App.mu. Controller and remote-tab sampling stay outside the
254 // projection mutex: archive and runtime-state callbacks also enter here, and
255 // holding that mutex across a controller read deadlocks a running turn.
256 func (a *App) GetRuntimeStateSnapshot() RuntimeStateProjection {
257 return a.runtimeStateSnapshotWithUpdate(nil)
258 }
259
260 func (a *App) runtimeStateSnapshotWithUpdate(update *localRuntimeUpdate) RuntimeStateProjection {
261 bindings := a.sampleLocalRuntimeBindingsWithUpdate(update)
262 return a.projectRuntimeBindings(bindings)
263 }
264
265 func (a *App) projectRuntimeBindings(bindings []localRuntimeBinding) RuntimeStateProjection {
266 r := &a.runtimeStateProjection
267 var remote []RuntimeSessionState
268 for {
269 remote = a.sampleRemoteRuntimeSessions()
270 r.mu.Lock()
271 // Sampling can finish before archive/rebind and publish after its
272 // replacement projection. Validate under the publication lock so an
273 // obsolete binding set cannot receive a newer projection revision.
274 a.mu.RLock()
275 current := a.localRuntimeBindingsLocked()
276 valid := len(current) == len(bindings)
277 for _, binding := range bindings {
278 if !sameLocalRuntimeBinding(current[binding.key], binding) {
279 valid = false
280 break
281 }
282 }
283 a.mu.RUnlock()
284 if valid {
285 break
286 }
287 r.mu.Unlock()
288 // Controller reads must remain outside both locks, including retries.
289 bindings = a.sampleLocalRuntimeBindingsWithUpdate(nil)
290 }
291 defer r.mu.Unlock()
292 nextBindings := make(map[localRuntimeBindingKey]localRuntimeBinding, len(bindings))
293 next := RuntimeStateProjection{Epoch: r.snapshot.Epoch, Sessions: []RuntimeSessionState{}}
294 catalog := []catalogRuntimeSnapshot{}
295 if next.Epoch == "" {
296 next.Epoch = newSessionRuntimeID("projection")
297 }
298 for _, binding := range bindings {
299 key := binding.key
300 // Another publisher may have committed after this off-lock sample.
301 // Never regress a state within the same binding and runtime epoch.
302 if previous, ok := r.bindings[key]; ok && sameLocalRuntimeBinding(previous, binding) &&
303 previous.view.State.RuntimeEpoch == binding.view.State.RuntimeEpoch &&
304 previous.view.State.Revision > binding.view.State.Revision {
305 binding.view.State = previous.view.State
306 }
307 nextBindings[key] = binding
308 view := binding.view
309 next.Sessions = append(next.Sessions, view)
310 if binding.catalog.topicID != "" {
311 entry := binding.catalog
312 entry.state = &view.State
313 catalog = append(catalog, entry)
314 }
315 }
316 r.bindings = nextBindings
317 next.Topics = a.projectTreeRuntimeTopics(catalog)
318 next.Sessions = append(next.Sessions, remote...)
319 sort.Slice(next.Sessions, func(i, j int) bool {
320 if next.Sessions[i].TabID == next.Sessions[j].TabID {
321 return next.Sessions[i].SessionPath < next.Sessions[j].SessionPath
322 }
323 return next.Sessions[i].TabID < next.Sessions[j].TabID
324 })
325 next.Revision = r.snapshot.Revision
326 if !reflect.DeepEqual(next, r.snapshot) {
327 next.Revision++
328 r.snapshot = next
329 }
330 result := r.snapshot
331 result.Sessions = append([]RuntimeSessionState{}, result.Sessions...)
332 result.Topics = cloneRuntimeTopics(result.Topics)
333 return result
334 }
335
336 func (a *App) sampleRemoteRuntimeSessions() []RuntimeSessionState {
337 if a == nil {
338 return nil
339 }
340 a.remoteTabMu.Lock()
341 defer a.remoteTabMu.Unlock()
342 sessions := make([]RuntimeSessionState, 0)
343 for _, tab := range a.remoteTabs {
344 freshness := "synced"
345 if tab.state != "ready" || tab.session.takenOver || tab.runtime.syncFailed || tab.runtimeUnknown[tab.routing.currentPath] != 0 {
346 freshness = "unknown"
347 }
348 state := tab.runtimeStates[tab.routing.currentPath]
349 if state.SchemaVersion == 0 {
350 state = event.RuntimeStateSnapshot{Phase: "idle", Running: tab.runtime.running, PendingPrompt: tab.runtime.pendingPrompt,
351 BackgroundJobs: tab.runtime.backgroundJobs, Cancellable: tab.runtime.cancellable, CancelRequested: tab.runtime.cancelRequested}
352 if state.Running {
353 state.Phase = "executing"
354 }
355 }
356 sessions = append(sessions, RuntimeSessionState{TabID: tab.id, Scope: "remote", HostID: tab.ref.HostID, WorkspaceRoot: tab.ref.Workspace,
357 SessionID: remoteRuntimeSessionID(tab.routing.currentPath, tab.session.sessionID, state), SessionPath: tab.routing.currentPath,
358 Open: true, Remote: true, Freshness: freshness, State: state})
359 for path, background := range tab.runtimeStates {
360 if path == tab.routing.currentPath {
361 continue
362 }
363 freshness := "synced"
364 // A foreground takeover says nothing about another session, but a
365 // tab without a live stream or with a failed sync only holds the
366 // snapshot frozen at its last observation.
367 if tab.state != "ready" || tab.runtime.syncFailed || tab.runtimeUnknown[path] != 0 {
368 freshness = "unknown"
369 }
370 sessions = append(sessions, RuntimeSessionState{TabID: tab.id, Scope: "remote", HostID: tab.ref.HostID, WorkspaceRoot: tab.ref.Workspace,
371 SessionID: remoteRuntimeSessionID(path, "", background), SessionPath: path, Remote: true, Freshness: freshness, State: background})
372 }
373 }
374 return sessions
375 }
376
377 func remoteRuntimeSessionID(route, fallback string, state event.RuntimeStateSnapshot) string {
378 if id := strings.TrimSpace(state.SessionID); id != "" {
379 return id
380 }
381 if id, ok := parseSessionRoute(route); ok {
382 return id
383 }
384 return strings.TrimSpace(fallback)
385 }
386
387 func (a *App) emitRuntimeStateChanged() {
388 if a != nil {
389 a.emitRuntimeEvent("runtime-state:changed", a.GetRuntimeStateSnapshot())
390 }
391 }
392
393 func (s *tabEventSink) RuntimeStateChanged(snapshot event.RuntimeStateSnapshot) {
394 id, app := s.binding()
395 if app == nil {
396 return
397 }
398 app.mu.RLock()
399 tab := app.tabByEventSinkIDLocked(id)
400 var ctrl control.SessionAPI
401 var generation uint64
402 var path, sessionID string
403 if tab != nil {
404 ctrl = tab.Ctrl
405 generation, path, sessionID = tab.SessionGeneration, tab.SessionPath, tab.SessionID
406 }
407 app.mu.RUnlock()
408 if ctrl == nil {
409 return
410 }
411 current := controllerRuntimeState(ctrl)
412 if current.RuntimeEpoch != snapshot.RuntimeEpoch || current.Revision > snapshot.Revision {
413 return
414 }
415 app.queueRuntimeProjection(localRuntimeUpdate{tab: tab, ctrl: ctrl, state: current, generation: generation, path: path, sessionID: sessionID})
416 if !current.Running && app.deferredRebuildPending(id) {
417 app.kickDeferredRebuildRetry()
418 }
419 }
420
420 lines GO