返回 DeepSeek-Reasonix
background_scope.go
根目录 / internal / jobs / background_scope.go
1 package jobs
2
3 import (
4 "errors"
5 "log/slog"
6 "sync"
7 "time"
8
9 "reasonix/internal/event"
10 "reasonix/internal/workspacelease"
11 )
12
13 // SessionBackgroundScope owns resources which outlive a controller generation.
14 // Build candidates acquire a reference before borrowing them; failed candidates
15 // release only that reference. Jobs do not retain the scope themselves.
16 type SessionBackgroundScope struct {
17 Manager *Manager
18 WorkspaceLease *workspacelease.Owner
19 mu sync.Mutex
20 refs int
21 closed bool
22 }
23
24 func NewSessionBackgroundScope(manager *Manager, lease *workspacelease.Owner) *SessionBackgroundScope {
25 return &SessionBackgroundScope{Manager: manager, WorkspaceLease: lease, refs: 1}
26 }
27
28 func (s *SessionBackgroundScope) Acquire() error {
29 s.mu.Lock()
30 defer s.mu.Unlock()
31 if s.closed {
32 return errors.New("background scope is closed")
33 }
34 s.refs++
35 return nil
36 }
37
38 func (s *SessionBackgroundScope) Release(async bool) {
39 s.mu.Lock()
40 if s.refs == 0 {
41 s.mu.Unlock()
42 return
43 }
44 s.refs--
45 closeNow := s.refs == 0
46 if closeNow {
47 s.closed = true
48 }
49 s.mu.Unlock()
50 if closeNow {
51 if async {
52 s.Manager.CloseAsync()
53 } else {
54 s.Manager.Close()
55 }
56 }
57 }
58
59 // Bind is called only on publication, never while staging a replacement.
60 func (s *SessionBackgroundScope) Bind(sink event.Sink, recorder TaskRecorder) {
61 s.Manager.bindingMu.Lock()
62 s.Manager.sink = sink
63 if s.Manager.taskRecorder == nil {
64 s.Manager.taskRecorder = recorder
65 }
66 s.Manager.bindingMu.Unlock()
67 }
68
69 type Lifetime string
70
71 const (
72 RuntimeBound Lifetime = "runtime_bound"
73 SessionProcess Lifetime = "session_process"
74 )
75
76 var ErrRebuildInProgress = errors.New("background task admission is paused for model configuration replacement")
77
78 // BeginReplacement checks and seals task registration under the same lock.
79 // Completion and cancellation remain available throughout the reservation.
80 func (m *Manager) BeginReplacement(session string) (func(), error) {
81 started := time.Now()
82 m.mu.Lock()
83 defer m.mu.Unlock()
84 if m.replacing {
85 return nil, ErrRebuildInProgress
86 }
87 if m.root.Err() != nil {
88 return nil, errors.New("background scope is closed")
89 }
90 if len(m.blockingJobsLocked(session)) > 0 {
91 return nil, errors.New("runtime-dependent background jobs are still running")
92 }
93 m.replacing = true
94 m.eventMu.Lock()
95 m.eventPaused = true
96 m.eventMu.Unlock()
97 return sync.OnceFunc(func() {
98 m.mu.Lock()
99 m.replacing = false
100 m.mu.Unlock()
101 m.eventMu.Lock()
102 m.eventPaused = false
103 m.eventMu.Unlock()
104 go func() {
105 m.drainEvents()
106 // Observers suppress candidate/old-generation callbacks while sealed.
107 // Resample after publication or rollback even if the final process
108 // exited during the build and no later job transition will occur.
109 m.notifyRuntime("", "")
110 slog.Debug("model replacement background reservation released", "phase", "release", "task_class", SessionProcess, "duration", time.Since(started))
111 }()
112 }), nil
113 }
114
115 func (m *Manager) BlockingJobs(session string) []View {
116 m.mu.Lock()
117 defer m.mu.Unlock()
118 return m.blockingJobsLocked(session)
119 }
120
121 func (m *Manager) blockingJobsLocked(session string) []View {
122 out := []View{}
123 for _, key := range m.order {
124 j := m.jobs[key]
125 if !sessionMatches(session, j.SessionID) || j.lifetime == SessionProcess {
126 continue
127 }
128 select {
129 case <-j.done:
130 continue
131 default:
132 }
133 j.mu.Lock()
134 out = append(out, View{ID: j.ID, Kind: j.Kind, Label: j.Label, Status: string(Running), StartedAt: j.clock.startedAt})
135 j.mu.Unlock()
136 }
137 return out
138 }
139
140 func (m *Manager) boundSink() event.Sink {
141 return m
142 }
143
144 // Emit queues lifecycle notices across replacement. Only one drainer invokes
145 // the current sink, always outside registry/binding locks.
146 func (m *Manager) Emit(e event.Event) {
147 m.eventMu.Lock()
148 m.eventQueue = append(m.eventQueue, e)
149 m.eventMu.Unlock()
150 m.drainEvents()
151 }
152
153 func (m *Manager) drainEvents() {
154 m.eventMu.Lock()
155 if m.eventPaused || m.eventDraining {
156 m.eventMu.Unlock()
157 return
158 }
159 m.eventDraining = true
160 for len(m.eventQueue) > 0 && !m.eventPaused {
161 e := m.eventQueue[0]
162 m.eventQueue[0] = event.Event{}
163 m.eventQueue = m.eventQueue[1:]
164 m.eventMu.Unlock()
165 m.bindingMu.RLock()
166 sink := m.sink
167 m.bindingMu.RUnlock()
168 if sink != nil {
169 sink.Emit(e)
170 }
171 m.eventMu.Lock()
172 }
173 m.eventDraining = false
174 m.eventMu.Unlock()
175 }
176
177 func (m *Manager) boundRecorder() TaskRecorder {
178 m.bindingMu.RLock()
179 defer m.bindingMu.RUnlock()
180 return m.taskRecorder
181 }
182
183 func (m *Manager) ActiveSessionID() string {
184 m.mu.Lock()
185 defer m.mu.Unlock()
186 return m.active
187 }
188
189 func (m *Manager) ReplacementInProgress() bool {
190 m.mu.Lock()
191 defer m.mu.Unlock()
192 return m.replacing
193 }
194
194 lines GO