| 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 |