| 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 |