返回 DeepSeek-Reasonix
runtime_projection_events.go
根目录 / desktop / runtime_projection_events.go
1 package main
2
3 import (
4 "sync"
5 "time"
6 )
7
8 type runtimeProjectionEvents struct {
9 mu sync.Mutex
10 pending map[*WorkspaceTab]localRuntimeUpdate
11 done chan struct{}
12 flush chan struct{}
13 }
14
15 // The per-controller stream remains immediate. Only the rebuildable global
16 // sidebar projection is batched, with one publisher for concurrent producers.
17 func (a *App) queueRuntimeProjection(update localRuntimeUpdate) {
18 c := &a.runtimeStateProjection.events
19 c.mu.Lock()
20 if a.shuttingDown.Load() {
21 c.mu.Unlock()
22 return
23 }
24 if c.pending == nil {
25 c.pending = map[*WorkspaceTab]localRuntimeUpdate{}
26 }
27 previous, exists := c.pending[update.tab]
28 if !exists || previous.state.RuntimeEpoch != update.state.RuntimeEpoch || previous.state.Revision <= update.state.Revision {
29 c.pending[update.tab] = update
30 }
31 if c.done != nil {
32 c.mu.Unlock()
33 return
34 }
35 c.done, c.flush = make(chan struct{}), make(chan struct{}, 1)
36 flush := c.flush
37 c.mu.Unlock()
38 go func() {
39 for {
40 timer := time.NewTimer(16 * time.Millisecond)
41 select {
42 case <-timer.C:
43 case <-flush:
44 }
45 timer.Stop()
46 c.mu.Lock()
47 updates := c.pending
48 c.pending = nil
49 c.mu.Unlock()
50 bindings := a.sampleLocalRuntimeBindingsWithUpdates(updates, true)
51 a.emitRuntimeProjection(a.projectRuntimeBindings(bindings), true)
52 c.mu.Lock()
53 if len(c.pending) == 0 {
54 close(c.done)
55 c.done, c.flush = nil, nil
56 c.mu.Unlock()
57 return
58 }
59 c.mu.Unlock()
60 }
61 }()
62 }
63
64 func (a *App) flushRuntimeProjections() {
65 c := &a.runtimeStateProjection.events
66 c.mu.Lock()
67 done := c.done
68 if done != nil {
69 select {
70 case c.flush <- struct{}{}:
71 default:
72 }
73 }
74 c.mu.Unlock()
75 if done != nil {
76 <-done
77 }
78 }
79
79 lines GO