| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "log/slog" |
| 6 | "maps" |
| 7 | "sort" |
| 8 | "strings" |
| 9 | "sync" |
| 10 | "time" |
| 11 | |
| 12 | "reasonix/internal/control" |
| 13 | "reasonix/internal/session" |
| 14 | ) |
| 15 | |
| 16 | type projectTreeRuntimeState struct { |
| 17 | // activityAt records the last activity-status event per tab ID. The TTL |
| 18 | // watchdog reaps a live status that sees no terminating event; the state |
| 19 | // lives here (not on WorkspaceTab) to keep tabs.go within budget. |
| 20 | activityMu sync.Mutex |
| 21 | activityAt map[string]time.Time |
| 22 | } |
| 23 | |
| 24 | // noteActivityStatus refreshes a tab's activity timestamp on every status |
| 25 | // event, even an unchanged one: the TTL watchdog measures silence, not how |
| 26 | // long a status has been displayed, so a long but active turn is never reaped. |
| 27 | func (s *projectTreeRuntimeState) noteActivityStatus(a *App, tabID string) { |
| 28 | s.setActivityAt(tabID, time.Now()) |
| 29 | // Runtime snapshots, not silence, determine whether work remains active. |
| 30 | } |
| 31 | |
| 32 | func (s *projectTreeRuntimeState) setActivityAt(tabID string, at time.Time) { |
| 33 | s.activityMu.Lock() |
| 34 | defer s.activityMu.Unlock() |
| 35 | if s.activityAt == nil { |
| 36 | s.activityAt = map[string]time.Time{} |
| 37 | } |
| 38 | s.activityAt[tabID] = at |
| 39 | } |
| 40 | |
| 41 | func (s *projectTreeRuntimeState) activityAtSnapshot() map[string]time.Time { |
| 42 | s.activityMu.Lock() |
| 43 | defer s.activityMu.Unlock() |
| 44 | out := make(map[string]time.Time, len(s.activityAt)) |
| 45 | maps.Copy(out, s.activityAt) |
| 46 | return out |
| 47 | } |
| 48 | |
| 49 | func syncRuntimeWorkspaceRootSpelling(tab *WorkspaceTab, projects []desktopProject) bool { |
| 50 | if tab == nil || tab.Scope != "project" { |
| 51 | return false |
| 52 | } |
| 53 | i := projectIndexByRoot(projects, tab.WorkspaceRoot) |
| 54 | if i < 0 || tab.WorkspaceRoot == projects[i].Root { |
| 55 | return false |
| 56 | } |
| 57 | tab.WorkspaceRoot = projects[i].Root |
| 58 | return true |
| 59 | } |
| 60 | |
| 61 | // catalogRuntimeSnapshots copies runtime identity under App.mu, then lets all |
| 62 | // controller calls happen after the app lock is released. Controllers own |
| 63 | // their own locks and must never become part of the App.mu lock order. |
| 64 | func (a *App) catalogRuntimeSnapshots() []catalogRuntimeSnapshot { |
| 65 | if a == nil { |
| 66 | return []catalogRuntimeSnapshot{} |
| 67 | } |
| 68 | a.mu.RLock() |
| 69 | snapshots := make([]catalogRuntimeSnapshot, 0, len(a.tabs)+len(a.detachedSessions)) |
| 70 | collect := func(tab *WorkspaceTab, open bool) { |
| 71 | if tab == nil || strings.TrimSpace(tab.TopicID) == "" { |
| 72 | return |
| 73 | } |
| 74 | snapshots = append(snapshots, catalogRuntimeSnapshot{ |
| 75 | tabID: tab.ID, |
| 76 | scope: tab.Scope, workspaceRoot: tab.WorkspaceRoot, topicID: tab.TopicID, |
| 77 | sessionPath: tab.SessionPath, sessionHeadID: tab.SessionHeadID, activity: tab.ActivityStatus, topicTitle: tab.TopicTitle, |
| 78 | topicTitleSource: tab.topicTitleSource, ctrl: tab.Ctrl, open: open, |
| 79 | }) |
| 80 | } |
| 81 | for _, tab := range a.tabs { |
| 82 | collect(tab, true) |
| 83 | } |
| 84 | for _, tab := range a.detachedSessions { |
| 85 | collect(tab, false) |
| 86 | } |
| 87 | a.mu.RUnlock() |
| 88 | return snapshots |
| 89 | } |
| 90 | |
| 91 | // GetProjectTreeRuntimeSnapshot returns the complete in-memory runtime |
| 92 | // projection. The frontend subscribes first and then calls this method; the |
| 93 | // independent revision makes either arrival order deterministic. |
| 94 | func (a *App) GetProjectTreeRuntimeSnapshot() ProjectTreeRuntimeSnapshot { |
| 95 | if a == nil { |
| 96 | return ProjectTreeRuntimeSnapshot{Topics: []ProjectRuntimeTopic{}} |
| 97 | } |
| 98 | snapshot := a.GetRuntimeStateSnapshot() |
| 99 | return ProjectTreeRuntimeSnapshot{Revision: snapshot.Revision, Topics: snapshot.Topics} |
| 100 | } |
| 101 | |
| 102 | func cloneRuntimeTopics(topics []ProjectRuntimeTopic) []ProjectRuntimeTopic { |
| 103 | var cloneNode func(ProjectNode) ProjectNode |
| 104 | cloneNode = func(node ProjectNode) ProjectNode { |
| 105 | next := node |
| 106 | next.Children = make([]ProjectNode, len(node.Children)) |
| 107 | for i, child := range node.Children { |
| 108 | next.Children[i] = cloneNode(child) |
| 109 | } |
| 110 | if node.Remote != nil { |
| 111 | remote := *node.Remote |
| 112 | next.Remote = &remote |
| 113 | } |
| 114 | return next |
| 115 | } |
| 116 | result := make([]ProjectRuntimeTopic, len(topics)) |
| 117 | for i, topic := range topics { |
| 118 | result[i] = topic |
| 119 | result[i].Node = cloneNode(topic.Node) |
| 120 | } |
| 121 | return result |
| 122 | } |
| 123 | |
| 124 | func (a *App) projectTreeRuntimeTopics(snapshots []catalogRuntimeSnapshot) []ProjectRuntimeTopic { |
| 125 | bySession := map[string]ProjectRuntimeTopic{} |
| 126 | state, _ := a.workspaceRegistry().LoadProjection(a.bootContext()) |
| 127 | workspaceAliases := map[string]map[string][]string{} |
| 128 | for _, snapshot := range snapshots { |
| 129 | scope, root := normalizeDesktopTopicScope(snapshot.scope, snapshot.workspaceRoot) |
| 130 | snapshot.scope, snapshot.workspaceRoot = scope, root |
| 131 | if snapshot.sessionPath == "" && snapshot.ctrl != nil { |
| 132 | snapshot.sessionPath = snapshot.ctrl.SessionPath() |
| 133 | } |
| 134 | nodes, _ := a.runtimeProjectTopicNodes(scope, root, []catalogRuntimeSnapshot{snapshot}, false) |
| 135 | if len(nodes) == 0 { |
| 136 | continue |
| 137 | } |
| 138 | node := nodes[0] |
| 139 | node.TabID = snapshot.tabID |
| 140 | path := strings.TrimSpace(snapshot.sessionPath) |
| 141 | if id, ok := parseSessionRoute(path); ok { |
| 142 | ref := session.SessionRef{HostID: localDesktopHostID, SessionID: id} |
| 143 | node.Session, node.SessionPath, node.Key = &ref, sessionRoute(id), "canonical_"+id |
| 144 | } else if identity, ok := snapshot.ctrl.(control.IdentityLifecycle); ok { |
| 145 | if ref, bound := identity.SessionRef(); bound && ref.SessionID != "" { |
| 146 | node.Session, node.SessionPath, node.Key = &ref, sessionRoute(ref.SessionID), "canonical_"+ref.SessionID |
| 147 | } |
| 148 | } |
| 149 | if node.Session == nil && path != "" { |
| 150 | node.SessionPath, node.Key = path, projectSessionNodeKey(scope, path) |
| 151 | // A legacy transcript is listed by its source identity; publishing the |
| 152 | // same one lets the renderer overlay this tab onto that row. |
| 153 | node.Source = &SessionSourceRef{HostID: localDesktopHostID, Path: path, HeadID: snapshot.sessionHeadID, SourceKey: desktopSourceKey(path, snapshot.sessionHeadID)} |
| 154 | } |
| 155 | if node.Session != nil { |
| 156 | workspaceID := desktopWorkspaceOwnerID(state, scope, root) |
| 157 | aliases, exists := workspaceAliases[workspaceID] |
| 158 | if !exists { |
| 159 | aliases = workspaceSourceAliases(state, workspaceID) |
| 160 | workspaceAliases[workspaceID] = aliases |
| 161 | } |
| 162 | node.IdentityAliases = append([]string{}, aliases[node.Session.SessionID]...) |
| 163 | node.LifecycleGeneration = state.SessionStates[node.Session.SessionID].Generation |
| 164 | if snapshot.tabID != "" { |
| 165 | node.IdentityAliases = append(node.IdentityAliases, "tab\x00local\x00"+snapshot.tabID) |
| 166 | } |
| 167 | } else if path == "" && snapshot.tabID != "" { |
| 168 | node.Key = "tab_" + snapshot.tabID |
| 169 | // This tab was opened from an unbacked topic. Publish that exact |
| 170 | // placeholder binding instead of asking the renderer to guess by title. |
| 171 | node.IdentityAliases = []string{"topic\x00" + snapshot.topicID} |
| 172 | } |
| 173 | key := scope + "\x00" + root + "\x00" + node.Key |
| 174 | bySession[key] = ProjectRuntimeTopic{Scope: scope, WorkspaceRoot: root, Node: node} |
| 175 | } |
| 176 | keys := make([]string, 0, len(bySession)) |
| 177 | for key := range bySession { |
| 178 | keys = append(keys, key) |
| 179 | } |
| 180 | sort.Strings(keys) |
| 181 | topics := make([]ProjectRuntimeTopic, 0, len(keys)) |
| 182 | for _, key := range keys { |
| 183 | topics = append(topics, bySession[key]) |
| 184 | } |
| 185 | return topics |
| 186 | } |
| 187 | |
| 188 | func (a *App) attachExistingSessionRuntime(tab *WorkspaceTab, path string, appCtx context.Context) bool { |
| 189 | attached := a.attachExistingSessionRuntimeCore(tab, path, appCtx) |
| 190 | if attached { |
| 191 | a.emitProjectTreeRuntimeChangedWithLegacy() |
| 192 | } |
| 193 | return attached |
| 194 | } |
| 195 | |
| 196 | func (a *App) closeTab(tabID string, allowDetach bool) error { |
| 197 | err := a.closeTabRuntime(tabID, allowDetach) |
| 198 | if err == nil { |
| 199 | a.emitProjectTreeRuntimeChangedWithLegacy() |
| 200 | } |
| 201 | return err |
| 202 | } |
| 203 | |
| 204 | func (a *App) emitProjectTreeChangedEvent() { |
| 205 | if a.projectTreeChangedHook != nil { |
| 206 | a.projectTreeChangedHook() |
| 207 | return |
| 208 | } |
| 209 | a.emitProjectTreeRuntimeChanged() |
| 210 | a.emitRuntimeEvent("project-tree:changed") |
| 211 | } |
| 212 | |
| 213 | func (a *App) emitProjectTreeRuntimeChanged() { |
| 214 | if a == nil { |
| 215 | return |
| 216 | } |
| 217 | snapshot := a.GetRuntimeStateSnapshot() |
| 218 | a.emitRuntimeProjection(snapshot, false) |
| 219 | } |
| 220 | |
| 221 | func (a *App) emitRuntimeProjection(snapshot RuntimeStateProjection, legacy bool) { |
| 222 | r := &a.runtimeStateProjection |
| 223 | r.mu.Lock() |
| 224 | if r.publishedEpoch == snapshot.Epoch && r.publishedRevision >= snapshot.Revision { |
| 225 | r.mu.Unlock() |
| 226 | return |
| 227 | } |
| 228 | r.publishedEpoch, r.publishedRevision = snapshot.Epoch, snapshot.Revision |
| 229 | r.mu.Unlock() |
| 230 | a.emitRuntimeEvent("project-tree:runtime-changed", ProjectTreeRuntimeSnapshot{Revision: snapshot.Revision, Topics: snapshot.Topics}) |
| 231 | a.emitRuntimeEvent("runtime-state:changed", snapshot) |
| 232 | if legacy { |
| 233 | a.emitRuntimeEvent("project-tree:changed", map[string]string{"reason": "runtime"}) |
| 234 | } |
| 235 | } |
| 236 | |
| 237 | // The tagged legacy event keeps the previous frontend usable for one release. |
| 238 | func (a *App) emitProjectTreeRuntimeChangedWithLegacy() { |
| 239 | a.emitProjectTreeRuntimeChanged() |
| 240 | a.emitRuntimeEvent("project-tree:changed", map[string]string{"reason": "runtime"}) |
| 241 | } |
| 242 | |
| 243 | func (a *App) emitRuntimeEvent(name string, payload ...any) { |
| 244 | if a != nil && a.ctx != nil { |
| 245 | a.runtimeEvents.Emit(a.ctx, name, payload...) |
| 246 | } |
| 247 | } |
| 248 | |
| 249 | const ( |
| 250 | // topicActivityStatusTTL bounds how long a live spinner status may go |
| 251 | // without any turn event before it is treated as orphaned. A flat TTL is |
| 252 | // used instead of a per-session P99: turn durations are not tracked |
| 253 | // per session in this package, and simplicity wins (#8528/#8555/#8859). |
| 254 | topicActivityStatusTTL = 10 * time.Minute |
| 255 | ) |
| 256 | |
| 257 | // liveTopicActivityStatus reports statuses that must be terminated by a |
| 258 | // TurnDone event; if that event is lost the spinner would run forever. |
| 259 | // waiting_confirmation is excluded: it waits on the user, not on turn events. |
| 260 | func liveTopicActivityStatus(status string) bool { |
| 261 | switch status { |
| 262 | case topicStatusThinking, topicStatusStreaming: |
| 263 | return true |
| 264 | } |
| 265 | return false |
| 266 | } |
| 267 | |
| 268 | func (a *App) reapStaleTopicActivityStatus(now time.Time) { |
| 269 | activityAt := a.projectTreeRuntime.activityAtSnapshot() |
| 270 | a.mu.Lock() |
| 271 | var reaped []string |
| 272 | for _, tab := range a.runtimeTabsLocked() { |
| 273 | if tab == nil || !liveTopicActivityStatus(tab.ActivityStatus) { |
| 274 | continue |
| 275 | } |
| 276 | if at := activityAt[tab.ID]; at.IsZero() || now.Sub(at) < topicActivityStatusTTL { |
| 277 | continue |
| 278 | } |
| 279 | reaped = append(reaped, tab.ID+":"+tab.ActivityStatus) |
| 280 | tab.ActivityStatus = "" |
| 281 | } |
| 282 | a.mu.Unlock() |
| 283 | if len(reaped) == 0 { |
| 284 | return |
| 285 | } |
| 286 | slog.Warn("desktop: cleared stale topic activity status (no turn event within TTL)", |
| 287 | "tabs", reaped, "ttl", topicActivityStatusTTL.String()) |
| 288 | a.emitProjectTreeRuntimeChangedWithLegacy() |
| 289 | } |
| 290 | |
| 291 | // reconcileTabActivityStatus clears a live spinner status the session's |
| 292 | // controller does not corroborate — a TurnDone missed while the session was |
| 293 | // detached would otherwise spin forever after reopen. The controller query |
| 294 | // takes the controller's own lock, so it runs outside App.mu (the same lock |
| 295 | // order rule as catalogRuntimeSnapshots). |
| 296 | func (a *App) reconcileTabActivityStatus(tab *WorkspaceTab) bool { |
| 297 | if a == nil || tab == nil { |
| 298 | return false |
| 299 | } |
| 300 | a.mu.RLock() |
| 301 | status := tab.ActivityStatus |
| 302 | ctrl := tab.Ctrl |
| 303 | a.mu.RUnlock() |
| 304 | if ctrl == nil || !liveTopicActivityStatus(status) || ctrl.Running() { |
| 305 | return false |
| 306 | } |
| 307 | a.mu.Lock() |
| 308 | defer a.mu.Unlock() |
| 309 | if tab.ActivityStatus != status { |
| 310 | // A new turn (or its TurnDone) landed while the locks were dropped. |
| 311 | return false |
| 312 | } |
| 313 | tab.ActivityStatus = "" |
| 314 | slog.Warn("desktop: cleared stale topic activity status on session open", "tab", tab.ID, "status", status) |
| 315 | return true |
| 316 | } |
| 317 | |
| 318 | // setTabActivityStatus records the project-tree status for a tab's in-flight |
| 319 | // turn and notes the event time for the TTL watchdog. |
| 320 | func (a *App) setTabActivityStatus(tabID, status string) bool { |
| 321 | a.mu.Lock() |
| 322 | defer a.mu.Unlock() |
| 323 | tab := a.tabByEventSinkIDLocked(tabID) |
| 324 | if tab == nil { |
| 325 | return false |
| 326 | } |
| 327 | a.projectTreeRuntime.noteActivityStatus(a, tab.ID) |
| 328 | status = normalizeTopicStatus(status) |
| 329 | if tab.ActivityStatus == status { |
| 330 | return false |
| 331 | } |
| 332 | tab.ActivityStatus = status |
| 333 | return true |
| 334 | } |
| 335 |