返回 DeepSeek-Reasonix
manual_creation_manager.go
根目录 / desktop / manual_creation_manager.go
1 package main
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "log/slog"
8 "sync"
9 "time"
10
11 "reasonix/desktop/internal/sessionui"
12 )
13
14 type manualCreationStore interface {
15 Get(context.Context, string, string) (sessionui.Record, error)
16 List(context.Context, string) ([]sessionui.Record, error)
17 Save(context.Context, string, string, string, json.RawMessage, ...sessionui.Record) (sessionui.Record, error)
18 }
19
20 type manualCreationTask struct {
21 id, retryRevision string
22 running, counted bool
23 wakeRequested bool
24 next time.Time
25 retries int
26 progress ManualCreationProgress
27 sessionID, revision string
28 attempt, generation uint64
29 lastSlowLog time.Time
30 }
31
32 // One scheduler owns admission and delayed attempts. The OS lock remains the
33 // cross-process authority; neither progress nor elapsed time grants ownership.
34 type manualCreationManager struct {
35 mu sync.Mutex
36 tasks map[string]*manualCreationTask
37 events []manualCreationDiagnostic
38 ctx context.Context
39 cancel context.CancelFunc
40 wake, done chan struct{}
41 stopped, scanning bool
42 nextScan time.Time
43 a *App
44 store manualCreationStore
45 now func() time.Time
46 execute func(context.Context, ManualSessionCreationView, func(string)) error
47 }
48
49 func (a *App) creationManager() *manualCreationManager {
50 a.manualCreationMu.Lock()
51 defer a.manualCreationMu.Unlock()
52 if a.manualCreations == nil {
53 ctx, cancel := context.WithCancel(a.bootContext())
54 m := &manualCreationManager{a: a, ctx: ctx, cancel: cancel, tasks: make(map[string]*manualCreationTask),
55 wake: make(chan struct{}, 1), done: make(chan struct{}), store: a.sessionUIStore(), now: time.Now,
56 execute: a.createManualSessionRuntime}
57 m.stopped = a.shuttingDown.Load()
58 if m.stopped {
59 cancel()
60 }
61 a.manualCreations = m
62 a.goSafe("manualCreationScheduler", m.loop)
63 }
64 return a.manualCreations
65 }
66
67 // retryRevision is populated only by an explicit retry of the observed failed
68 // record. A delayed request cannot restart a newer failure.
69 func (m *manualCreationManager) Ensure(id, reason, retryRevision string) {
70 m.mu.Lock()
71 defer m.mu.Unlock()
72 if m.stopped || m.ctx.Err() != nil || m.a.shuttingDown.Load() {
73 return
74 }
75 t := m.tasks[id]
76 if t != nil {
77 if t.running {
78 if reason == "retry" {
79 t.retryRevision = retryRevision
80 t.wakeRequested = true
81 }
82 return
83 }
84 if t.progress.Status == "blocked" && reason != "retry" {
85 return
86 }
87 }
88 if t == nil {
89 t = &manualCreationTask{id: id}
90 m.tasks[id] = t
91 }
92 if reason == "retry" {
93 t.retryRevision = retryRevision
94 }
95 if !t.counted {
96 m.a.manualCreationTasks.Add(1)
97 t.counted = true
98 }
99 if t.next.IsZero() || reason == "retry" {
100 t.next = m.now()
101 t.progress = ManualCreationProgress{Status: "queued", Stage: "queued", StageStartedAt: t.next.UnixMilli()}
102 }
103 m.notify()
104 }
105
106 func (m *manualCreationManager) notify() {
107 select {
108 case m.wake <- struct{}{}:
109 default:
110 }
111 }
112
113 func (m *manualCreationManager) armRecovery() {
114 m.mu.Lock()
115 if !m.stopped {
116 m.scanning = true
117 m.nextScan = m.now()
118 }
119 m.mu.Unlock()
120 m.notify()
121 }
122
123 func (m *manualCreationManager) StopAdmission() {
124 m.mu.Lock()
125 m.stopped = true
126 m.scanning = false
127 m.mu.Unlock()
128 m.notify()
129 }
130
131 func (m *manualCreationManager) CancelAndWait(ctx context.Context) error {
132 m.StopAdmission()
133 m.cancel()
134 select {
135 case <-m.done:
136 return nil
137 case <-ctx.Done():
138 return ctx.Err()
139 }
140 }
141
142 func (m *manualCreationManager) loop() {
143 defer close(m.done)
144 for {
145 if m.ctx.Err() != nil {
146 m.mu.Lock()
147 for id, t := range m.tasks {
148 if !t.running {
149 m.finishLocked(t)
150 delete(m.tasks, id)
151 }
152 }
153 m.mu.Unlock()
154 m.a.manualCreationTasks.Wait()
155 return
156 }
157 m.mu.Lock()
158 now := m.now()
159 scan := !m.stopped && m.scanning && !now.Before(m.nextScan)
160 if scan {
161 m.nextScan = now.Add(30 * time.Second)
162 }
163 delay := time.Second // also drives rate-limited slow-stage diagnostics
164 for _, t := range m.tasks {
165 if m.stopped || t.running || !t.counted {
166 continue
167 }
168 if now.Before(t.next) {
169 if d := t.next.Sub(now); d < delay {
170 delay = d
171 }
172 continue
173 }
174 t.running = true
175 t.attempt++
176 t.generation = 0
177 m.a.goSafe("manualCreationAttempt", func() { m.run(t) })
178 }
179 m.mu.Unlock()
180 m.logSlowStages(now)
181 if scan {
182 m.scan()
183 }
184 timer := time.NewTimer(delay)
185 select {
186 case <-m.ctx.Done():
187 case <-m.wake:
188 case <-timer.C:
189 }
190 timer.Stop()
191 }
192 }
193
194 func (m *manualCreationManager) finishLocked(t *manualCreationTask) {
195 if t.counted {
196 t.counted = false
197 m.a.manualCreationTasks.Done()
198 }
199 }
200
201 func (m *manualCreationManager) run(t *manualCreationTask) {
202 retry, err := m.attempt(t)
203 m.mu.Lock()
204 t.running = false
205 if m.ctx.Err() == nil && !m.stopped && t.wakeRequested && (retry || err == nil) {
206 t.wakeRequested = false
207 t.next = m.now()
208 t.retries = 0
209 } else if m.ctx.Err() == nil && !m.stopped && retry {
210 t.next = m.now().Add(manualCreationBackoff(t.retries))
211 t.retries++
212 t.progress.NextRetryAt = t.next.UnixMilli()
213 } else {
214 if err != nil && m.ctx.Err() == nil && !errors.Is(err, context.Canceled) {
215 t.progress.Status, t.progress.ErrorCode = "blocked", manualCreationErrorCode(err)
216 } else {
217 t.progress.Status = "completed"
218 if m.ctx.Err() != nil || errors.Is(err, context.Canceled) {
219 t.progress.Status = "stopped"
220 }
221 delete(m.tasks, t.id)
222 }
223 m.recordEventLocked(t)
224 m.finishLocked(t)
225 }
226 attrs := m.logAttrsLocked(t)
227 attrs = append(attrs, "error_code", t.progress.ErrorCode)
228 m.mu.Unlock()
229 slog.Info("desktop: manual creation attempt finished", attrs...)
230 m.notify()
231 }
232
233 func (m *manualCreationManager) mayRetryFailure(t *manualCreationTask, revision string) bool {
234 m.mu.Lock()
235 defer m.mu.Unlock()
236 return t.retryRevision == revision
237 }
238
238 lines GO