返回 DeepSeek-Reasonix
manual_creation_attempt.go
根目录 / desktop / manual_creation_attempt.go
1 package main
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "math/rand/v2"
10 "os"
11 "path/filepath"
12 "time"
13
14 "reasonix/desktop/internal/sessionui"
15 "reasonix/internal/identitylock"
16 )
17
18 var errManualCreationIdentity = errors.New("invalid manual creation identity or state")
19
20 func decodeManualCreation(r sessionui.Record) (ManualSessionCreationView, error) {
21 var v ManualSessionCreationView
22 if err := json.Unmarshal(r.Payload, &v); err != nil {
23 return v, err
24 }
25 sum := sha256.Sum256([]byte(r.Key))
26 if v.OperationID != r.Key || v.WorkspaceID == "" || v.Ref.HostID != localDesktopHostID ||
27 v.Ref.SessionID != fmt.Sprintf("desktop-manual-%x", sum[:16]) || v.TopicID != fmt.Sprintf("manual-%x", sum[:16]) {
28 return v, errManualCreationIdentity
29 }
30 switch v.Phase {
31 case "reserved", "starting", "ready", "failed":
32 default:
33 return v, errManualCreationIdentity
34 }
35 return v, nil
36 }
37
38 func manualCreationBackoff(n int) time.Duration {
39 base := 500 * time.Millisecond * time.Duration(1<<min(n, 4))
40 // Small downward jitter preserves the five-second upper bound.
41 return min(base, 5*time.Second) * time.Duration(900+rand.IntN(101)) / 1000
42 }
43
44 func manualCreationRetryable(err error) bool {
45 if errors.Is(err, sessionui.ErrConflict) {
46 return true
47 }
48 var sqlErr interface{ Code() int }
49 if errors.As(err, &sqlErr) {
50 switch sqlErr.Code() & 255 {
51 case 5, 6, 10:
52 return true
53 } // BUSY, LOCKED, IOERR
54 }
55 return false
56 }
57
58 // Target availability can be retried with the same identity. Result-save retries
59 // retain the lock and the computed result, never re-entering the runtime builder.
60 func (m *manualCreationManager) attempt(t *manualCreationTask) (bool, error) {
61 r, err := m.store.Get(m.ctx, "creation", t.id)
62 if err != nil {
63 m.setStage(t, "retrying_storage", "queued")
64 return manualCreationRetryable(err), err
65 }
66 v, err := decodeManualCreation(r)
67 if err != nil {
68 return false, err
69 }
70 m.mu.Lock()
71 t.sessionID, t.revision = v.Ref.SessionID, r.Revision
72 m.mu.Unlock()
73 if v.Phase == "ready" || (v.Phase == "failed" && !legacyManualCreationConflict(v) && !m.mayRetryFailure(t, r.Revision)) {
74 return false, nil
75 }
76 lockRoot := filepath.Join(filepath.Dir(m.a.sessionUIStore().Path()), "manual-creation-locks")
77 if err := os.MkdirAll(lockRoot, 0700); err != nil {
78 return false, err
79 }
80 release, err := identitylock.TryAcquire(filepath.Join(lockRoot, v.Ref.SessionID+".lock"))
81 if errors.Is(err, identitylock.ErrHeld) {
82 m.setStage(t, "waiting_lock", "waiting_lock")
83 return true, nil
84 }
85 if err != nil {
86 return false, err
87 }
88 defer release()
89 current, err := m.store.Get(m.ctx, "creation", t.id)
90 if err != nil {
91 m.setStage(t, "retrying_storage", "preparing_storage")
92 return manualCreationRetryable(err), err
93 }
94 fresh, err := decodeManualCreation(current)
95 if err != nil || fresh.Ref != v.Ref || fresh.WorkspaceID != v.WorkspaceID {
96 return false, errManualCreationIdentity
97 }
98 if fresh.Phase == "ready" || (fresh.Phase == "failed" && !legacyManualCreationConflict(fresh) && !m.mayRetryFailure(t, current.Revision)) {
99 return false, nil
100 }
101 if err := m.ctx.Err(); err != nil {
102 return false, err
103 }
104 r, v = current, fresh
105 r, err = m.savePhase(r, "starting", "")
106 if err != nil {
107 m.setStage(t, "retrying_storage", "queued")
108 return manualCreationRetryable(err), err
109 }
110 m.mu.Lock()
111 t.revision = r.Revision
112 m.mu.Unlock()
113 err = m.execute(m.ctx, v, func(stage string) { m.setStage(t, "running", stage) })
114 // Shutdown interruption is recoverable on the next start, not a user failure.
115 if m.ctx.Err() != nil || m.a.shuttingDown.Load() {
116 return false, context.Canceled
117 }
118 var targetErr *SessionOperationError
119 if errors.As(err, &targetErr) && targetErr.Code == "workspace_unavailable" {
120 m.setStage(t, "waiting_workspace", "preparing_storage")
121 return true, nil
122 }
123 phase, message := "ready", ""
124 if err != nil {
125 phase, message = "failed", sessionOperationErrorForTarget(err, v.Ref.SessionID, v.OperationID).Error()
126 }
127 m.setStage(t, "running", "persisting_result")
128 err = m.persistResult(t, r, v, phase, message)
129 if err == nil {
130 m.a.emitProjectTreeChanged()
131 }
132 return false, err
133 }
134
135 func (m *manualCreationManager) savePhase(r sessionui.Record, phase, message string) (sessionui.Record, error) {
136 var fields map[string]json.RawMessage
137 if err := json.Unmarshal(r.Payload, &fields); err != nil {
138 return r, err
139 }
140 fields["phase"], _ = json.Marshal(phase)
141 delete(fields, "error")
142 if message != "" {
143 fields["error"], _ = json.Marshal(message)
144 }
145 // Only phase/error are replaced; nested and unknown fields stay opaque.
146 payload, err := json.Marshal(fields)
147 if err != nil {
148 return r, err
149 }
150 return m.store.Save(m.ctx, "creation", r.Key, r.Revision, payload)
151 }
152
153 func (m *manualCreationManager) persistResult(t *manualCreationTask, r sessionui.Record, identity ManualSessionCreationView, phase, message string) error {
154 for n := 0; ; n++ {
155 if err := m.ctx.Err(); err != nil {
156 return err
157 }
158 saved, err := m.savePhase(r, phase, message)
159 if err == nil {
160 m.mu.Lock()
161 t.revision = saved.Revision
162 m.mu.Unlock()
163 return nil
164 }
165 if !manualCreationRetryable(err) {
166 return err
167 }
168 m.setStage(t, "retrying_storage", "persisting_result")
169 if err := m.waitStorageRetry(t, n); err != nil {
170 return err
171 }
172 // Retry reads as well as writes without ever discarding the runtime result.
173 for {
174 r, err = m.store.Get(m.ctx, "creation", t.id)
175 if err == nil {
176 break
177 }
178 if !manualCreationRetryable(err) {
179 return err
180 }
181 if err = m.waitStorageRetry(t, n); err != nil {
182 return err
183 }
184 }
185 v, err := decodeManualCreation(r)
186 if err != nil || v.Ref != identity.Ref || v.WorkspaceID != identity.WorkspaceID {
187 return errManualCreationIdentity
188 }
189 if v.Phase == phase {
190 m.mu.Lock()
191 t.revision = r.Revision
192 m.mu.Unlock()
193 return nil
194 }
195 if v.Phase != "starting" {
196 return errManualCreationIdentity
197 }
198 m.mu.Lock()
199 t.revision = r.Revision
200 m.mu.Unlock()
201 }
202 }
203
204 func (m *manualCreationManager) waitStorageRetry(t *manualCreationTask, n int) error {
205 delay := manualCreationBackoff(n)
206 m.mu.Lock()
207 t.progress.NextRetryAt = m.now().Add(delay).UnixMilli()
208 m.mu.Unlock()
209 timer := time.NewTimer(delay)
210 defer timer.Stop()
211 select {
212 case <-m.ctx.Done():
213 return m.ctx.Err()
214 case <-timer.C:
215 return nil
216 }
217 }
218
218 lines GO