返回 DeepSeek-Reasonix
start.go
根目录 / internal / jobs / start.go
1 package jobs
2
3 import (
4 "context"
5 "fmt"
6 "io"
7 "reasonix/internal/event"
8 "reasonix/internal/nilutil"
9 "strings"
10 )
11
12 // StartForSession launches a job owned by parentSession. Session-scoped readers
13 // only see jobs whose owner matches the active session.
14 func (m *Manager) StartForSession(parentSession, kind, label string, run func(ctx context.Context, out io.Writer) (string, error)) *Job {
15 j, err := m.TryStartForSession(parentSession, kind, label, run)
16 if err != nil {
17 return m.startInvalid(parentSession, kind, label, err)
18 }
19 return j
20 }
21
22 // TryStartForSession leaves captured resources with the caller on rejection.
23 func (m *Manager) TryStartForSession(parentSession, kind, label string, run func(context.Context, io.Writer) (string, error)) (*Job, error) {
24 return m.startForSession(parentSession, kind, label, RuntimeBound, run)
25 }
26
27 // StartSessionProcess is reserved for shell processes whose resources belong
28 // to the logical session rather than the controller which launched them.
29 func (m *Manager) StartSessionProcess(parentSession, kind, label string, run func(context.Context, io.Writer) (string, error)) *Job {
30 j, err := m.TryStartSessionProcess(parentSession, kind, label, run)
31 if err != nil {
32 return m.startInvalid(parentSession, kind, label, err)
33 }
34 return j
35 }
36
37 func (m *Manager) TryStartSessionProcess(parentSession, kind, label string, run func(context.Context, io.Writer) (string, error)) (*Job, error) {
38 return m.startForSession(parentSession, kind, label, SessionProcess, run)
39 }
40
41 func (m *Manager) startForSession(parentSession, kind, label string, lifetime Lifetime, run func(context.Context, io.Writer) (string, error)) (*Job, error) {
42 parentSession = strings.TrimSpace(parentSession)
43 kind = strings.TrimSpace(kind)
44 if err := validatePathSegment(parentSession, "parentSession"); err != nil {
45 return nil, err
46 }
47 if err := validatePathSegment(kind, "kind"); err != nil {
48 return nil, err
49 }
50 m.mu.Lock()
51 if m.replacing || m.root.Err() != nil {
52 m.mu.Unlock()
53 return nil, ErrRebuildInProgress
54 }
55 m.seq++
56 id := fmt.Sprintf("%s-%d", kind, m.seq)
57 ctx, cancel := context.WithCancel(m.root)
58 startedAt := nowMs()
59 logPath, metaPath, file, artifactErr := m.openArtifactLocked(parentSession, id)
60 j := &Job{
61 lifetime: lifetime,
62 ID: id,
63 Kind: kind,
64 Label: label,
65 SessionID: parentSession,
66 status: Running,
67 clock: jobClock{startedAt: startedAt, activityAt: startedAt},
68 cancel: cancel,
69 done: make(chan struct{}),
70 artifactPath: logPath,
71 artifactMetaPath: metaPath,
72 artifactFile: file,
73 artifactComplete: artifactErr == "",
74 artifactErr: artifactErr,
75 }
76 ctx = WithSession(ctx, parentSession)
77 ctx = context.WithValue(ctx, jobCtxKey{}, j)
78 key := jobKey(parentSession, id)
79 m.jobs[key] = j
80 m.order = append(m.order, key)
81 m.wg.Add(1)
82 if m.stalledWarning > 0 {
83 m.wg.Add(1)
84 }
85 m.mu.Unlock()
86 j.mu.Lock()
87 if err := m.writeJobMetaLocked(j, Running); err != nil {
88 j.artifactComplete = false
89 j.artifactErr = err.Error()
90 }
91 j.mu.Unlock()
92 if m.onJobStart != nil {
93 m.onJobStart(j.done)
94 }
95
96 m.emitIfActive(parentSession, event.Event{Kind: event.Notice, Level: event.LevelInfo, Text: startedText(kind, id, label)})
97 m.notifyRuntime(parentSession, id)
98
99 if recorder := m.boundRecorder(); !nilutil.IsNil(recorder) {
100 recorder.RecordStart(id, kind, label)
101 }
102
103 if m.stalledWarning > 0 {
104 go m.monitorStalled(parentSession, j)
105 }
106 go m.runJob(ctx, j, run)
107 return j, nil
108 }
109
110 func (m *Manager) runJob(ctx context.Context, j *Job, run func(context.Context, io.Writer) (string, error)) {
111 defer m.wg.Done()
112 result, err := runRecovered(ctx, jobWriter{j}, run)
113 j.mu.Lock()
114 j.outcome.returned = true
115 j.mu.Unlock()
116
117 var st Status
118 switch {
119 case ctx.Err() != nil:
120 st = Killed
121 case err != nil:
122 st = Failed
123 if result == "" {
124 result = err.Error()
125 }
126 default:
127 st = Done
128 }
129 finishedAt := nowMs()
130 if result != "" {
131 j.mu.Lock()
132 if j.artifactFile != nil {
133 if _, writeErr := j.artifactFile.WriteString(result); writeErr != nil {
134 j.artifactErr = writeErr.Error()
135 }
136 } else {
137 j.outcome.text = result
138 }
139 j.tail = appendTail(j.tail, []byte(result), defaultTailBytes)
140 j.mu.Unlock()
141 }
142 targetDir := m.artifactTargetDirForJob(j)
143 j.mu.Lock()
144 if j.artifactFile != nil {
145 if closeErr := j.artifactFile.Close(); closeErr != nil && j.artifactErr == "" {
146 j.artifactErr = closeErr.Error()
147 }
148 j.artifactFile = nil
149 }
150 if j.artifactErr != "" {
151 j.artifactComplete = false
152 }
153 j.clock.finishedAt = finishedAt
154 if targetDir != "" {
155 if moveErr := j.moveArtifactToDirLocked(targetDir); moveErr != nil {
156 j.noteArtifactErr("migration: " + moveErr.Error())
157 }
158 }
159 metaErr := m.writeJobMetaLocked(j, st)
160 if metaErr != nil {
161 j.noteArtifactErr("metadata: " + metaErr.Error())
162 }
163 j.mu.Unlock()
164 // Queue the drain note and closing Notice before terminal status so Wait
165 // cannot observe completion before DrainCompletedNote sees its bookkeeping.
166 // The structured runtime notification follows the actual done boundary.
167 parentSession := m.recordCompletion(j, st, err)
168
169 j.mu.Lock()
170 if j.status != Killed { // a concurrent Kill already published Killed — keep it
171 j.status = st
172 }
173 if j.artifactPath != "" && j.artifactComplete {
174 j.outcome.text = ""
175 j.tail = nil
176 }
177 j.mu.Unlock()
178 close(j.done)
179 m.notifyRuntime(parentSession, j.ID)
180 }
181
181 lines GO