返回 DeepSeek-Reasonix
manual_creation_progress.go
根目录 / desktop / manual_creation_progress.go
1 package main
2
3 import (
4 "encoding/json"
5 "errors"
6 "log/slog"
7 "os"
8 "sort"
9 "time"
10
11 "reasonix/desktop/internal/sessionui"
12 )
13
14 // ManualCreationProgress is local observational state, never a lease or a
15 // persisted phase. Missing/unknown fields are safe for older clients.
16 type ManualCreationProgress struct {
17 Status string `json:"status"`
18 Stage string `json:"stage"`
19 StageStartedAt int64 `json:"stageStartedAt"`
20 ElapsedMs int64 `json:"elapsedMs"`
21 NextRetryAt int64 `json:"nextRetryAt,omitempty"`
22 ErrorCode string `json:"errorCode,omitempty"`
23 Slow bool `json:"slow"`
24 }
25
26 func manualCreationErrorCode(err error) string {
27 if errors.Is(err, sessionui.ErrFutureVersion) {
28 return "unsupported_ui_schema"
29 }
30 if errors.Is(err, errManualCreationIdentity) {
31 return "creation_state_conflict"
32 }
33 return "creation_storage_unavailable"
34 }
35
36 func (m *manualCreationManager) Snapshot(id string) *ManualCreationProgress {
37 m.mu.Lock()
38 defer m.mu.Unlock()
39 if t := m.tasks[id]; t != nil {
40 p := m.progressLocked(t)
41 return &p
42 }
43 return nil
44 }
45
46 func (m *manualCreationManager) progressLocked(t *manualCreationTask) ManualCreationProgress {
47 p := t.progress
48 p.ElapsedMs = max(0, m.now().UnixMilli()-p.StageStartedAt)
49 p.Slow = p.ElapsedMs >= 30000
50 return p
51 }
52
53 func (a *App) creationView(v ManualSessionCreationView) ManualSessionCreationView {
54 v.Progress = nil
55 v.SurfaceReady = false
56 a.mu.RLock()
57 for _, tab := range a.tabs {
58 if !tab.removed && tab.SessionID == v.Ref.SessionID && v.Ref.SessionID != "" {
59 v.SurfaceReady = true
60 break
61 }
62 }
63 a.mu.RUnlock()
64 a.manualCreationMu.Lock()
65 m := a.manualCreations
66 a.manualCreationMu.Unlock()
67 if m != nil {
68 v.Progress = m.Snapshot(v.OperationID)
69 }
70 return v
71 }
72
73 func (m *manualCreationManager) setStage(t *manualCreationTask, status, stage string) {
74 m.setStageForGeneration(t, status, stage, 0)
75 }
76
77 func (m *manualCreationManager) setStageForGeneration(t *manualCreationTask, status, stage string, generation uint64) {
78 if stage == "stopping" {
79 status = "stopping"
80 }
81 m.mu.Lock()
82 // A late stage callback must never alter a newer task for the same identity.
83 if m.tasks[t.id] != t || !t.running {
84 m.mu.Unlock()
85 return
86 }
87 if generation != 0 {
88 if generation < t.generation {
89 m.mu.Unlock()
90 return
91 }
92 t.generation = generation
93 }
94 now := m.now()
95 changed := t.progress.Stage != stage || t.progress.Status != status
96 if t.progress.Stage != stage {
97 t.progress.StageStartedAt = now.UnixMilli()
98 t.lastSlowLog = time.Time{}
99 }
100 t.progress.Status, t.progress.Stage, t.progress.NextRetryAt = status, stage, 0
101 attrs := m.logAttrsLocked(t)
102 if changed {
103 m.recordEventLocked(t)
104 }
105 m.mu.Unlock()
106 if changed {
107 slog.Info("desktop: manual creation stage", attrs...)
108 }
109 }
110
111 func (m *manualCreationManager) logAttrsLocked(t *manualCreationTask) []any {
112 return []any{"operation", t.id, "session", t.sessionID, "pid", os.Getpid(), "attempt", t.attempt,
113 "generation", t.generation, "revision", t.revision, "stage", t.progress.Stage,
114 "status", t.progress.Status, "elapsed_ms", max(0, m.now().UnixMilli()-t.progress.StageStartedAt)}
115 }
116
117 func (m *manualCreationManager) logSlowStages(now time.Time) {
118 m.mu.Lock()
119 var reports [][]any
120 for _, t := range m.tasks {
121 if t.counted && now.UnixMilli()-t.progress.StageStartedAt >= 30000 && (t.lastSlowLog.IsZero() || now.Sub(t.lastSlowLog) >= time.Minute) {
122 t.lastSlowLog = now
123 reports = append(reports, m.logAttrsLocked(t))
124 }
125 }
126 m.mu.Unlock()
127 for _, attrs := range reports {
128 slog.Warn("desktop: manual creation slow stage", attrs...)
129 }
130 }
131
132 func (m *manualCreationManager) scan() {
133 rows, err := m.store.List(m.ctx, "creation")
134 if err != nil {
135 slog.Warn("desktop: pending creation scan unavailable", "code", manualCreationErrorCode(err))
136 return
137 }
138 for _, r := range rows {
139 var v ManualSessionCreationView
140 if json.Unmarshal(r.Payload, &v) != nil {
141 m.Ensure(r.Key, "recovery", "")
142 continue
143 }
144 if (v.Phase != "ready" && v.Phase != "failed") || legacyManualCreationConflict(v) {
145 m.Ensure(r.Key, "recovery", "")
146 }
147 }
148 }
149
150 // Content-free local diagnostics can be exported even if the database is busy.
151 type manualCreationDiagnostic struct {
152 OperationID string `json:"operationId"`
153 SessionID string `json:"sessionId"`
154 PID int `json:"pid"`
155 Attempt uint64 `json:"attempt"`
156 Generation uint64 `json:"generation"`
157 Revision string `json:"revision"`
158 Progress ManualCreationProgress `json:"progress"`
159 }
160
161 func (m *manualCreationManager) recordEventLocked(t *manualCreationTask) {
162 m.events = append(m.events, manualCreationDiagnostic{t.id, t.sessionID, os.Getpid(), t.attempt, t.generation, t.revision, m.progressLocked(t)})
163 if len(m.events) > 256 {
164 m.events = append([]manualCreationDiagnostic(nil), m.events[len(m.events)-256:]...)
165 }
166 }
167
168 func (a *App) manualCreationDiagnostics() []manualCreationDiagnostic {
169 rows := []manualCreationDiagnostic{}
170 a.manualCreationMu.Lock()
171 m := a.manualCreations
172 a.manualCreationMu.Unlock()
173 if m == nil {
174 return rows
175 }
176 m.mu.Lock()
177 for _, t := range m.tasks {
178 rows = append(rows, manualCreationDiagnostic{t.id, t.sessionID, os.Getpid(), t.attempt, t.generation, t.revision, m.progressLocked(t)})
179 }
180 m.mu.Unlock()
181 sort.Slice(rows, func(i, j int) bool { return rows[i].OperationID < rows[j].OperationID })
182 return rows
183 }
184
184 lines GO