返回 DeepSeek-Reasonix
jobs.go
根目录 / internal / jobs / jobs.go
1 // Package jobs is the session-scoped background-job registry behind the agent's
2 // background tools (bash run_in_background, task run_in_background) and the
3 // job_output / job_kill tools (plus replay-only legacy aliases). A Manager owns a context whose lifetime
4 // is the session, NOT a single turn — so a job started in one turn keeps running
5 // across turns and is cancelled only when the controller closes (or job_kill is
6 // called). Tools reach the Manager through the call context (WithManager /
7 // FromContext), the same injection pattern the `ask` tool uses for the asker.
8 //
9 // The Manager emits a user-visible Notice when a job starts and finishes, and
10 // accumulates a one-line completion summary that the controller drains into the
11 // next turn (DrainCompletedNote) so the model itself learns of completions.
12 package jobs
13
14 import (
15 "context"
16 "crypto/rand"
17 "encoding/hex"
18 "fmt"
19 "io"
20 "os"
21 "path/filepath"
22 "runtime/debug"
23 "slices"
24 "sort"
25 "strings"
26 "sync"
27 "sync/atomic"
28 "time"
29
30 "reasonix/internal/event"
31 "reasonix/internal/evidence"
32 "reasonix/internal/nilutil"
33 "reasonix/internal/tool"
34 )
35
36 var renamePath = os.Rename
37 var repairArtifactMeta = writeMeta
38
39 var (
40 managerOwnerSeq atomic.Uint64
41 liveManagerOwners = struct {
42 sync.RWMutex
43 ids map[string]struct{}
44 }{ids: map[string]struct{}{}}
45 )
46
47 // Status is a job's lifecycle state.
48 type Status string
49
50 const (
51 Running Status = "running"
52 Done Status = "done"
53 Failed Status = "failed"
54 Killed Status = "killed"
55 Interrupted Status = "interrupted"
56 )
57
58 // DefaultTeardownGrace bounds Close and destroy waits for non-cooperative jobs.
59 const DefaultTeardownGrace = 15 * time.Second
60
61 // View is a read-only snapshot of a job for the status bar.
62 type View struct {
63 ID string `json:"id"`
64 Kind string `json:"kind"`
65 Label string `json:"label"`
66 Status string `json:"status"`
67 StartedAt int64 `json:"startedAt"` // unix milliseconds
68 }
69
70 // Result is one job's terminal (or current) state returned by Wait.
71 type Result struct {
72 ID string
73 Kind string
74 Label string
75 Status Status
76 Output string // the terminal result text, or the streamed buffer when no result was set
77 }
78
79 // TeardownJob identifies a job that is still unwinding after teardown waited.
80 type TeardownJob struct {
81 ID string
82 Kind string
83 Label string
84 Waited time.Duration
85 }
86
87 // TeardownResult reports jobs that did not unwind within the teardown grace.
88 type TeardownResult struct {
89 TimedOut []TeardownJob
90 }
91
92 type teardownTarget struct {
93 info TeardownJob
94 done <-chan struct{}
95 }
96
97 // SessionTeardown is the destroy handle for a session's owned background jobs.
98 type SessionTeardown struct {
99 SessionID string
100 targets []teardownTarget
101 }
102
103 // Async reports whether the handle has jobs to wait on.
104 func (h SessionTeardown) Async() bool { return len(h.targets) > 0 }
105
106 // DoneChannels returns each target's completion channel for legacy callers.
107 func (h SessionTeardown) DoneChannels() []<-chan struct{} {
108 out := make([]<-chan struct{}, 0, len(h.targets))
109 for _, target := range h.targets {
110 out = append(out, target.done)
111 }
112 return out
113 }
114
115 // Job is one background job. The mutex guards the streaming buffer and the
116 // terminal fields; the run goroutine writes them, readers (Output/Wait/snapshots)
117 // take the same lock.
118 type Job struct {
119 lifetime Lifetime
120 ID string
121 Kind string // "bash" | "pwsh" | "task"
122 Label string
123 SessionID string
124
125 mu sync.Mutex
126 tail []byte
127 readOffset int64
128 status Status
129 clock jobClock
130 outcome jobOutcome
131 cancel context.CancelFunc
132 done chan struct{}
133
134 artifactPath string
135 artifactMetaPath string
136 artifactStatus Status // last metadata phase; may precede published terminal status
137 artifactFile *os.File
138 artifactComplete bool
139 artifactErr string
140 tombstone bool
141
142 evidence evidence.ChildEvidenceSummary
143 evidenceCommitted bool
144 execution *tool.ShellExecution
145 }
146
147 // Manager is the session's background-job table. It is safe for concurrent use.
148 type Manager struct {
149 eventMu sync.Mutex
150 eventQueue []event.Event
151 eventPaused bool
152 eventDraining bool
153 bindingMu sync.RWMutex
154 replacing bool // guarded by mu; seals registration during replacement
155 runtimeObservers runtimeObservers
156 sink event.Sink
157 root context.Context
158 cancel context.CancelFunc
159 wg sync.WaitGroup
160 onJobStart func(done <-chan struct{})
161 ownerID string
162 ownerDone sync.Once
163 // sessionOwnershipProbe authorizes destructive repair of persisted running
164 // artifacts. A nil probe is conservative: an observer that cannot prove it
165 // owns the transcript must never publish an interrupted tombstone.
166 sessionOwnershipProbe func(path string) bool
167
168 mu sync.Mutex
169 seq int
170 jobs map[string]*Job
171 order []string
172 completed []completion // finished-job summaries awaiting drain into the next turn
173 active string
174 destroying map[string]bool
175 artifactDirs map[string]string
176 loaded map[string]bool
177 tempRoot string
178 reservations map[string]int
179
180 stalledWarning time.Duration
181 teardownGrace time.Duration
182
183 taskRecorder TaskRecorder // optional task-monitoring lifecycle hook
184 }
185
186 type completion struct {
187 sessionID string
188 text string
189 }
190
191 // Option configures a Manager.
192 type Option func(*Manager)
193
194 // TaskRecorder observes background-job lifecycle for task monitoring. The
195 // store-backed write side lives outside jobs (typically internal/taskmonitor);
196 // jobs only calls the hooks. RecordStart runs on the caller's goroutine,
197 // RecordDone on the job's own goroutine — implementations must be safe for
198 // concurrent use and must not block or fail the job pipeline (best-effort).
199 type TaskRecorder interface {
200 RecordStart(id, kind, label string)
201 RecordDone(id string, st Status, err error)
202 }
203
204 // WithStalledWarningAfter enables one stalled warning per job after d without
205 // job-owned visible output. A non-positive duration disables stalled warnings.
206 func WithStalledWarningAfter(d time.Duration) Option {
207 return func(m *Manager) {
208 if d > 0 {
209 m.stalledWarning = d
210 }
211 }
212 }
213
214 // WithTeardownGrace overrides the Close/destroy grace window. Tests can set a
215 // short value; production uses DefaultTeardownGrace.
216 func WithTeardownGrace(d time.Duration) Option {
217 return func(m *Manager) {
218 if d >= 0 {
219 m.teardownGrace = d
220 }
221 }
222 }
223
224 // WithJobStartObserver observes every registered background job before its
225 // goroutine starts. Delivery uses this to retain a workspace writer lease over
226 // the job's opening writes. The callback must return quickly.
227 func WithJobStartObserver(observer func(done <-chan struct{})) Option {
228 return func(m *Manager) { m.onJobStart = observer }
229 }
230
231 // WithSessionOwnershipProbe supplies the runtime ownership check used when
232 // loading persisted Running artifacts. The probe must return true only when the
233 // current runtime owns the session transcript for writing.
234 func WithSessionOwnershipProbe(probe func(path string) bool) Option {
235 return func(m *Manager) { m.sessionOwnershipProbe = probe }
236 }
237
238 // WithTaskRecorder installs an optional background-job lifecycle recorder for
239 // task monitoring. A nil recorder disables recording.
240 func WithTaskRecorder(r TaskRecorder) Option {
241 return func(m *Manager) { m.taskRecorder = r }
242 }
243
244 // SetTaskRecorder installs (or clears, with nil) the lifecycle recorder after
245 // construction. Controllers that assemble their job manager before the
246 // recorder's dependencies (workspace root, session id) are known use this.
247 func (m *Manager) SetTaskRecorder(r TaskRecorder) {
248 m.bindingMu.Lock()
249 m.taskRecorder = r
250 m.bindingMu.Unlock()
251 }
252
253 // TeardownGrace reports the manager's configured close/destroy wait window.
254 func (m *Manager) TeardownGrace() time.Duration { return m.teardownGrace }
255
256 // NewManager returns a Manager whose jobs run under a fresh session-scoped
257 // context (cancelled by Close). sink receives job-lifecycle notices; pass the
258 // session's synchronized sink (event.Sync) since jobs emit from goroutines.
259 func NewManager(sink event.Sink, opts ...Option) *Manager {
260 if nilutil.IsNil(sink) {
261 sink = event.Discard
262 }
263 root, cancel := context.WithCancel(context.Background())
264 tempRoot, _ := os.MkdirTemp("", "reasonix-jobs-*")
265 m := &Manager{
266 sink: sink,
267 root: root,
268 cancel: cancel,
269 jobs: map[string]*Job{},
270 destroying: map[string]bool{},
271 artifactDirs: map[string]string{},
272 reservations: map[string]int{},
273 loaded: map[string]bool{},
274 tempRoot: tempRoot,
275 teardownGrace: DefaultTeardownGrace,
276 ownerID: newManagerOwnerID(),
277 }
278 registerManagerOwner(m.ownerID)
279 for _, opt := range opts {
280 if opt != nil {
281 opt(m)
282 }
283 }
284 return m
285 }
286
287 func newManagerOwnerID() string {
288 var token [16]byte
289 if _, err := rand.Read(token[:]); err == nil {
290 return hex.EncodeToString(token[:])
291 }
292 return fmt.Sprintf("%d-%d-%d", os.Getpid(), time.Now().UnixNano(), managerOwnerSeq.Add(1))
293 }
294
295 func registerManagerOwner(ownerID string) {
296 ownerID = strings.TrimSpace(ownerID)
297 if ownerID == "" {
298 return
299 }
300 liveManagerOwners.Lock()
301 liveManagerOwners.ids[ownerID] = struct{}{}
302 liveManagerOwners.Unlock()
303 }
304
305 func managerOwnerIsLive(ownerID string) bool {
306 ownerID = strings.TrimSpace(ownerID)
307 if ownerID == "" {
308 return false
309 }
310 liveManagerOwners.RLock()
311 _, ok := liveManagerOwners.ids[ownerID]
312 liveManagerOwners.RUnlock()
313 return ok
314 }
315
316 func (m *Manager) releaseOwner() {
317 if m == nil {
318 return
319 }
320 m.ownerDone.Do(func() {
321 liveManagerOwners.Lock()
322 delete(liveManagerOwners.ids, m.ownerID)
323 liveManagerOwners.Unlock()
324 })
325 }
326
327 // jobWriter appends a job's streamed output under its lock so a concurrent
328 // Output read never races the producing goroutine.
329 type jobWriter struct{ j *Job }
330
331 func (w jobWriter) Write(p []byte) (int, error) {
332 w.j.mu.Lock()
333 defer w.j.mu.Unlock()
334 w.j.clock.activityAt = nowMs()
335 w.j.tail = appendTail(w.j.tail, p, defaultTailBytes)
336 if w.j.artifactFile != nil {
337 if _, err := w.j.artifactFile.Write(p); err != nil {
338 w.j.artifactErr = err.Error()
339 }
340 }
341 return len(p), nil
342 }
343
344 // Start launches run on a goroutine under the manager's session context and
345 // returns the job immediately. run streams output to the writer and returns the
346 // terminal result text (a task's final answer; a bash job streams everything to
347 // the buffer and returns ""). The job is marked killed when its context was
348 // cancelled, failed on any other error, else done.
349 func (m *Manager) Start(kind, label string, run func(ctx context.Context, out io.Writer) (string, error)) *Job {
350 return m.StartForSession("", kind, label, run)
351 }
352
353 // validatePathSegment rejects values that would let parentSession or kind
354 // escape the temp-root fallback built by artifactDirLocked. Persistent artifact
355 // directories bound by SetActiveSessionPath are trusted store paths and are
356 // intentionally outside that temp root. The check is intentionally conservative:
357 // it forbids any path-separator character (forward slash, backslash), NUL, and
358 // any control character. Empty parentSession is allowed (the unscoped default);
359 // kind must be non-empty.
360 //
361 // See #6932. Before this check existed, a malicious or malformed parentSession
362 // such as "../../etc" combined with filepath.Join(tempRoot, parentSession, id)
363 // resolved to a directory outside the manager's temp root, allowing the
364 // subsequent os.MkdirAll + os.OpenFile to create files at locations controlled
365 // by the caller (subject to the running process's filesystem permissions).
366 func validatePathSegment(name, field string) error {
367 if field == "kind" && name == "" {
368 return fmt.Errorf("jobs: %s must not be empty", field)
369 }
370 for i, r := range name {
371 switch {
372 case r < 0x20 || r == 0x7f:
373 return fmt.Errorf("jobs: %s contains control character 0x%02x at index %d", field, r, i)
374 case r == '/' || r == '\\':
375 return fmt.Errorf("jobs: %s contains path separator %q at index %d", field, r, i)
376 }
377 }
378 if name == "." || name == ".." {
379 return fmt.Errorf("jobs: %s is reserved (%q)", field, name)
380 }
381 return nil
382 }
383
384 // startInvalid registers a job that failed validation BEFORE any goroutine or
385 // artifact was created. The job is observable to Wait / list calls as Failed
386 // with the validation error recorded in artifactErr, and no run goroutine is
387 // started so the manager's wg is unaffected.
388 func (m *Manager) startInvalid(parentSession, kind, label string, validationErr error) *Job {
389 finishedAt := nowMs()
390 m.mu.Lock()
391 m.seq++
392 id := fmt.Sprintf("invalid-%d", m.seq)
393 j := &Job{
394 ID: id,
395 Kind: kind,
396 Label: label,
397 SessionID: parentSession,
398 status: Failed,
399 clock: jobClock{startedAt: finishedAt, activityAt: finishedAt, finishedAt: finishedAt},
400 outcome: jobOutcome{returned: true},
401 cancel: func() {},
402 done: make(chan struct{}),
403 artifactComplete: false,
404 artifactErr: validationErr.Error(),
405 }
406 key := jobKey(parentSession, id)
407 m.jobs[key] = j
408 m.order = append(m.order, key)
409 m.mu.Unlock()
410 close(j.done)
411 m.recordCompletion(j, Failed, validationErr)
412 m.notifyRuntime(parentSession, id)
413 return j
414 }
415
416 func runRecovered(ctx context.Context, out io.Writer, run func(context.Context, io.Writer) (string, error)) (result string, err error) {
417 defer func() {
418 if r := recover(); r != nil {
419 err = fmt.Errorf("internal error: panic: %v\n%s", r, debug.Stack())
420 }
421 }()
422 return run(ctx, out)
423 }
424
425 func (m *Manager) openArtifactLocked(parentSession, id string) (logPath, metaPath string, file *os.File, artifactErr string) {
426 dir := m.artifactDirLocked(parentSession)
427 if dir == "" {
428 return "", "", nil, "artifact directory unavailable"
429 }
430 if err := ensurePrivateArtifactDir(dir); err != nil {
431 return filepath.Join(dir, id+jobLogExt), filepath.Join(dir, id+jobMetaExt), nil, err.Error()
432 }
433 logPath = filepath.Join(dir, id+jobLogExt)
434 metaPath = filepath.Join(dir, id+jobMetaExt)
435 f, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o600)
436 if err != nil {
437 return logPath, metaPath, nil, err.Error()
438 }
439 // O_TRUNC does not apply the requested mode to an existing artifact. Tighten
440 // it before any raw tool output is written so upgrades cannot append secrets
441 // to a legacy 0644 log.
442 if err := f.Chmod(0o600); err != nil {
443 _ = f.Close()
444 return logPath, metaPath, nil, err.Error()
445 }
446 return logPath, metaPath, f, ""
447 }
448
449 func ensurePrivateArtifactDir(dir string) error {
450 if err := os.MkdirAll(dir, 0o700); err != nil {
451 return err
452 }
453 // MkdirAll leaves an existing 0755 directory unchanged.
454 return os.Chmod(dir, 0o700)
455 }
456
457 func (m *Manager) artifactDirLocked(parentSession string) string {
458 parentSession = strings.TrimSpace(parentSession)
459 if parentSession != "" {
460 if dir := strings.TrimSpace(m.artifactDirs[parentSession]); dir != "" {
461 return dir
462 }
463 }
464 if strings.TrimSpace(m.tempRoot) == "" {
465 return ""
466 }
467 if parentSession == "" {
468 return filepath.Join(m.tempRoot, "default")
469 }
470 return filepath.Join(m.tempRoot, parentSession)
471 }
472
473 func (m *Manager) writeJobMetaLocked(j *Job, st Status) error {
474 j.artifactStatus = st
475 if j.artifactMetaPath == "" {
476 return nil
477 }
478 meta := artifactMeta{
479 ID: j.ID,
480 Kind: j.Kind,
481 Label: j.Label,
482 SessionID: j.SessionID,
483 OwnerID: m.ownerID,
484 Status: st,
485 StartedAt: j.clock.startedAt,
486 FinishedAt: j.clock.finishedAt,
487 ArtifactComplete: st != Running && j.artifactComplete && j.artifactErr == "",
488 ArtifactError: j.artifactErr,
489 LogPath: filepath.Base(j.artifactPath),
490 }
491 if j.Kind == "task" {
492 meta.MutationEvidenceVersion = mutationEvidenceVersion
493 meta.MutationEvidence = mutationEvidenceForArtifact(j.evidence)
494 }
495 return writeMeta(j.artifactMetaPath, meta)
496 }
497
498 func mutationEvidenceForArtifact(summary evidence.ChildEvidenceSummary) *artifactMutationEvidence {
499 firstMutation := -1
500 for i, receipt := range summary.Receipts {
501 if receipt.Success && receipt.Mutation {
502 firstMutation = i
503 break
504 }
505 }
506 if firstMutation < 0 {
507 return nil
508 }
509 return &artifactMutationEvidence{
510 Paths: summary.MutationPaths(),
511 }
512 }
513
514 func mutationEvidenceFromArtifact(meta artifactMeta) evidence.ChildEvidenceSummary {
515 if meta.Kind != "task" {
516 return evidence.ChildEvidenceSummary{}
517 }
518 if meta.MutationEvidenceVersion != mutationEvidenceVersion {
519 // Any version this build cannot parse — a pre-feature artifact
520 // (version 0) or one written by a newer build — is treated as an
521 // opaque mutation. A missing summary only proves the mutation state
522 // was not recorded, not that the task made no changes: a legacy
523 // background writer task collected after upgrade could carry real,
524 // edits. Preserve the existing unknown-mutation compatibility record;
525 // it never implies verification or creates an acceptance requirement.
526 return opaqueRecoveredTaskMutation()
527 }
528 if meta.MutationEvidence == nil {
529 // Same-version artifact with no summary: this build DID record the
530 // mutation state and found none, so there is genuinely nothing to
531 // recover.
532 return evidence.ChildEvidenceSummary{}
533 }
534
535 paths := append([]string(nil), meta.MutationEvidence.Paths...)
536 // Historical risk labels do not erase observed paths or create obligations.
537 return evidence.ChildEvidenceSummary{Receipts: []evidence.Receipt{{
538 ToolName: recoveredBackgroundTaskToolName,
539 Success: true,
540 Write: true,
541 Mutation: true,
542 Paths: paths,
543 }}}
544 }
545
546 func opaqueRecoveredTaskMutation() evidence.ChildEvidenceSummary {
547 return evidence.ChildEvidenceSummary{Receipts: []evidence.Receipt{{
548 ToolName: recoveredBackgroundTaskToolName,
549 Success: true,
550 Write: true,
551 Mutation: true,
552 }}}
553 }
554
555 func (m *Manager) artifactTargetDirForJob(j *Job) string {
556 if j == nil {
557 return ""
558 }
559 m.mu.Lock()
560 defer m.mu.Unlock()
561 session := strings.TrimSpace(j.SessionID)
562 if session == "" {
563 return ""
564 }
565 return strings.TrimSpace(m.artifactDirs[session])
566 }
567
568 func (j *Job) noteArtifactErr(msg string) {
569 msg = strings.TrimSpace(msg)
570 if msg == "" {
571 return
572 }
573 if j.artifactErr == "" {
574 j.artifactErr = msg
575 } else {
576 j.artifactErr += "; " + msg
577 }
578 j.artifactComplete = false
579 }
580
581 func (j *Job) moveArtifactToDirLocked(dir string) error {
582 dir = strings.TrimSpace(dir)
583 if dir == "" || j.artifactPath == "" {
584 return nil
585 }
586 if filepath.Clean(filepath.Dir(j.artifactPath)) == filepath.Clean(dir) {
587 return nil
588 }
589 if err := ensurePrivateArtifactDir(dir); err != nil {
590 return err
591 }
592 newLogPath := filepath.Join(dir, filepath.Base(j.artifactPath))
593 if err := moveArtifactFile(j.artifactPath, newLogPath); err != nil {
594 return err
595 }
596 j.artifactPath = newLogPath
597 if j.artifactMetaPath != "" {
598 j.artifactMetaPath = filepath.Join(dir, filepath.Base(j.artifactMetaPath))
599 }
600 return nil
601 }
602
603 func (m *Manager) monitorStalled(parentSession string, j *Job) {
604 defer m.wg.Done()
605 timer := time.NewTimer(m.stalledWarning)
606 defer timer.Stop()
607 for {
608 select {
609 case <-j.done:
610 return
611 case <-timer.C:
612 j.mu.Lock()
613 if j.outcome.returned || j.status != Running {
614 j.mu.Unlock()
615 return
616 }
617 idle := time.Since(time.UnixMilli(j.clock.activityAt))
618 if idle >= m.stalledWarning && !j.clock.stalled {
619 j.clock.stalled = true
620 j.mu.Unlock()
621 m.recordStalled(parentSession, j.ID, j.Kind, j.Label)
622 return
623 }
624 wait := m.stalledWarning - idle
625 if wait <= 0 {
626 wait = m.stalledWarning
627 }
628 j.mu.Unlock()
629 timer.Reset(wait)
630 }
631 }
632 }
633
634 // recordCompletion queues the finished-job summary for DrainCompletedNote and
635 // emits a closing Notice (warn for a failure, info otherwise).
636 func (m *Manager) recordCompletion(j *Job, st Status, err error) string {
637 id, kind, label := j.ID, j.Kind, j.Label
638 tag := id
639 if label != "" {
640 tag = fmt.Sprintf("%s (%s)", id, label)
641 }
642 shouldEmit := false
643 m.mu.Lock()
644 parentSession := strings.TrimSpace(j.SessionID)
645 if parentSession != "" && m.destroying[parentSession] {
646 m.mu.Unlock()
647 return parentSession
648 }
649 m.completed = append(m.completed, completion{
650 sessionID: parentSession,
651 text: fmt.Sprintf("%s — %s", tag, st),
652 })
653 active := m.active
654 shouldEmit = active == "" || parentSession == "" || active == parentSession
655 m.mu.Unlock()
656
657 if recorder := m.boundRecorder(); !nilutil.IsNil(recorder) {
658 recorder.RecordDone(id, st, err)
659 }
660
661 level, text := event.LevelInfo, fmt.Sprintf("background %s finished: %s", kind, id)
662 detail := ""
663 switch st {
664 case Failed:
665 level, text = event.LevelWarn, fmt.Sprintf("background %s failed: needs attention", kind)
666 detail = fmt.Sprintf("background %s failed: %s — %v", kind, id, err)
667 case Killed:
668 text = fmt.Sprintf("background %s killed: %s", kind, id)
669 }
670 if shouldEmit {
671 m.boundSink().Emit(event.Event{Kind: event.Notice, Code: event.NoticeCodeBackgroundJobFinished, Level: level, Text: text, Detail: detail})
672 }
673 return parentSession
674 }
675
676 func (m *Manager) recordStalled(parentSession, id, kind, label string) {
677 tag := id
678 if label != "" {
679 tag = fmt.Sprintf("%s (%s)", id, label)
680 }
681 parentSession = strings.TrimSpace(parentSession)
682 m.mu.Lock()
683 if parentSession != "" && m.destroying[parentSession] {
684 m.mu.Unlock()
685 return
686 }
687 quietFor := m.stalledWarning.Round(time.Second)
688 text := fmt.Sprintf("%s is still running after %s with no visible output — a quiet long-running job can look like this and is not necessarily stuck. If it should have finished, inspect it with job_output, or stop it with job_kill. Tune or disable this check with tools.background_jobs.stalled_warning_seconds in your config (0 disables).", tag, quietFor)
689 m.completed = append(m.completed, completion{sessionID: parentSession, text: text})
690 active := m.active
691 shouldEmit := active == "" || parentSession == "" || active == parentSession
692 notice := event.Event{Kind: event.Notice, Level: event.LevelWarn,
693 Text: fmt.Sprintf("background %s still running after %s with no visible output: %s", kind, quietFor, id),
694 Detail: "A quiet long-running job can look like this, so this is a heads-up, not an error. If it should have finished, inspect with job_output, or stop it with job_kill. Set tools.background_jobs.stalled_warning_seconds to 0 in your config to disable this notice."}
695 m.mu.Unlock()
696 if shouldEmit {
697 m.boundSink().Emit(notice)
698 }
699 }
700
701 func (m *Manager) get(parentSession, id string) *Job {
702 m.mu.Lock()
703 defer m.mu.Unlock()
704 return m.findJobLocked(parentSession, id)
705 }
706
707 func (m *Manager) findJobLocked(parentSession, id string) *Job {
708 parentSession = strings.TrimSpace(parentSession)
709 id = strings.TrimSpace(id)
710 if parentSession != "" {
711 return m.jobs[jobKey(parentSession, id)]
712 }
713 for _, key := range m.order {
714 j := m.jobs[key]
715 if j != nil && j.ID == id {
716 return j
717 }
718 }
719 return nil
720 }
721
722 // Output returns the job's output produced since the last Output call plus its
723 // current status. ok is false when the id is unknown.
724 func (m *Manager) Output(id string) (text string, status Status, ok bool) {
725 return m.OutputForSession("", id)
726 }
727
728 // OutputForSession returns output only when id belongs to parentSession. Empty
729 // parentSession preserves the legacy unscoped behavior.
730 func (m *Manager) OutputForSession(parentSession, id string) (text string, status Status, ok bool) {
731 j := m.get(parentSession, id)
732 if j == nil {
733 return "", "", false
734 }
735 j.mu.Lock()
736 defer j.mu.Unlock()
737 if j.artifactPath != "" {
738 text = j.readArtifactSinceOffsetLocked()
739 } else {
740 full := string(j.tail)
741 if j.readOffset < int64(len(full)) {
742 text = full[j.readOffset:]
743 j.readOffset = int64(len(full))
744 }
745 }
746 // A task job streams nothing to the tail buffer — its answer lands in
747 // outcome.text. Surface it once when terminal with no buffered output, so a
748 // task's answer is visible here too (bash_output promises task support).
749 if text == "" && j.status != Running && j.outcome.text != "" && !j.outcome.read {
750 text = j.outcome.text
751 j.outcome.read = true
752 }
753 if j.artifactErr != "" {
754 if text != "" {
755 text += "\n"
756 }
757 text += "job artifact incomplete: " + j.artifactErr
758 }
759 return text, j.status, true
760 }
761
762 func (j *Job) readArtifactSinceOffsetLocked() string {
763 f, err := os.Open(j.artifactPath)
764 if err != nil {
765 if j.artifactErr == "" {
766 j.artifactErr = err.Error()
767 }
768 return ""
769 }
770 defer f.Close()
771 info, err := f.Stat()
772 if err != nil {
773 if j.artifactErr == "" {
774 j.artifactErr = err.Error()
775 }
776 return ""
777 }
778 size := info.Size()
779 if j.readOffset > size {
780 j.readOffset = size
781 return ""
782 }
783 if _, err := f.Seek(j.readOffset, io.SeekStart); err != nil {
784 if j.artifactErr == "" {
785 j.artifactErr = err.Error()
786 }
787 return ""
788 }
789 b, err := io.ReadAll(f)
790 if err != nil {
791 if j.artifactErr == "" {
792 j.artifactErr = err.Error()
793 }
794 return ""
795 }
796 text := string(b)
797 j.readOffset = size
798 return text
799 }
800
801 // readArtifactAllLocked deliberately reads raw bytes: the artifact is captured
802 // subprocess output (possibly binary), not a user-edited config file, and the
803 // incremental reader (readArtifactSinceOffsetLocked) is raw byte-offset based —
804 // decoding only the whole-file path would render the same artifact in two
805 // different encodings and could garble binary output via UTF-16 misdetection.
806 func (j *Job) readArtifactAllLocked() string {
807 if j.artifactPath == "" {
808 return ""
809 }
810 b, err := os.ReadFile(j.artifactPath)
811 if err != nil {
812 if j.artifactErr == "" {
813 j.artifactErr = err.Error()
814 }
815 return ""
816 }
817 return string(b)
818 }
819
820 // Kill cancels a running job. Returns false when the id is unknown or the job has
821 // already finished.
822 func (m *Manager) Kill(id string) bool {
823 return m.KillForSession("", id)
824 }
825
826 // KillForSession cancels a running job only when it belongs to parentSession.
827 // Empty parentSession preserves the legacy unscoped behavior.
828 func (m *Manager) KillForSession(parentSession, id string) bool {
829 j := m.get(parentSession, id)
830 if j == nil {
831 return false
832 }
833 j.mu.Lock()
834 running := j.status == Running
835 if running {
836 // Flip to Killed synchronously so Output/Wait reflect the kill the instant
837 // it's requested, not whenever the run goroutine's cmd.Run returns (which
838 // trails by WaitDelay while a cancelled process tree tears down). The
839 // goroutine still sets Killed + records completion on return; this only
840 // fires when the job is actually Running, so a job that just finished
841 // keeps its real terminal status.
842 j.status = Killed
843 }
844 j.mu.Unlock()
845 if !running {
846 return false
847 }
848 j.cancel()
849 return true
850 }
851
852 // Wait blocks until the named jobs (or every currently-running job when ids is
853 // empty) reach a terminal state, or ctx is cancelled, or timeoutSec elapses
854 // (0 = no timeout). It returns each target's snapshot regardless of why it
855 // returned, so a timeout still reports partial progress.
856 func (m *Manager) Wait(ctx context.Context, ids []string, timeoutSec int) []Result {
857 return m.WaitForSession(ctx, "", ids, timeoutSec)
858 }
859
860 // WaitForSession waits only on jobs owned by parentSession. Empty parentSession
861 // preserves the legacy unscoped behavior.
862 func (m *Manager) WaitForSession(ctx context.Context, parentSession string, ids []string, timeoutSec int) []Result {
863 targets := m.resolve(parentSession, ids)
864 if len(targets) == 0 {
865 return nil
866 }
867 var timeout <-chan time.Time
868 if timeoutSec > 0 {
869 t := time.NewTimer(time.Duration(timeoutSec) * time.Second)
870 defer t.Stop()
871 timeout = t.C
872 }
873 for _, j := range targets {
874 select {
875 case <-j.done:
876 case <-ctx.Done():
877 return m.results(targets)
878 case <-timeout:
879 return m.results(targets)
880 }
881 }
882 return m.results(targets)
883 }
884
885 // resolve maps requested ids to jobs; an empty list selects all running jobs.
886 func (m *Manager) resolve(parentSession string, ids []string) []*Job {
887 m.mu.Lock()
888 defer m.mu.Unlock()
889 var out []*Job
890 if len(ids) == 0 {
891 for _, key := range m.order {
892 j := m.jobs[key]
893 if !sessionMatches(parentSession, j.SessionID) {
894 continue
895 }
896 j.mu.Lock()
897 running := j.status == Running
898 j.mu.Unlock()
899 if running {
900 out = append(out, j)
901 }
902 }
903 return out
904 }
905 for _, id := range ids {
906 if j := m.findJobLocked(parentSession, id); j != nil {
907 out = append(out, j)
908 }
909 }
910 return out
911 }
912
913 func (m *Manager) results(targets []*Job) []Result {
914 out := make([]Result, 0, len(targets))
915 for _, j := range targets {
916 j.mu.Lock()
917 text := j.outcome.text
918 if text == "" && j.artifactPath != "" {
919 text = j.readArtifactAllLocked()
920 }
921 if text == "" {
922 text = string(j.tail)
923 }
924 if j.artifactErr != "" {
925 if text != "" {
926 text += "\n"
927 }
928 text += "job artifact incomplete: " + j.artifactErr
929 }
930 out = append(out, Result{ID: j.ID, Kind: j.Kind, Label: j.Label, Status: j.status, Output: text})
931 j.mu.Unlock()
932 }
933 return out
934 }
935
936 // Running returns a snapshot of the still-running jobs (for the status bar).
937 func (m *Manager) Running() []View {
938 return m.RunningForSession("")
939 }
940
941 // RunningForSession returns still-running jobs owned by parentSession. Empty
942 // parentSession preserves the legacy unscoped behavior.
943 func (m *Manager) RunningForSession(parentSession string) []View {
944 m.mu.Lock()
945 defer m.mu.Unlock()
946 var out []View
947 for _, key := range m.order {
948 j := m.jobs[key]
949 if !sessionMatches(parentSession, j.SessionID) {
950 continue
951 }
952 select {
953 case <-j.done:
954 continue
955 default:
956 }
957 j.mu.Lock()
958 // A cancellation request flips the persisted/result status to Killed
959 // synchronously, but the process tree may still be unwinding. Keep the job
960 // on the operational running surface until its done channel closes so
961 // Desktop rebuild guards and Delivery workspace leases cannot declare the
962 // runtime idle early. The public view remains "running" while a stop is
963 // in flight; clients may render a local "stopping" state after they
964 // request cancellation.
965 out = append(out, View{ID: j.ID, Kind: j.Kind, Label: j.Label, Status: string(Running), StartedAt: j.clock.startedAt})
966 j.mu.Unlock()
967 }
968 return out
969 }
970
971 // ReserveStartForSession atomically reserves capacity for a job start. The
972 // caller must release the reservation after StartForSession has registered the
973 // job (or when setup fails). Running jobs and in-flight start reservations both
974 // count toward limit, so concurrent callers cannot overshoot it.
975 func (m *Manager) ReserveStartForSession(parentSession, kind string, limit int) (release func(), running int, ok bool) {
976 if limit <= 0 {
977 return func() {}, 0, true
978 }
979 parentSession = strings.TrimSpace(parentSession)
980 key := jobKey(parentSession, kind)
981 m.mu.Lock()
982 for _, jobKey := range m.order {
983 j := m.jobs[jobKey]
984 if j == nil || !sessionMatches(parentSession, j.SessionID) || j.Kind != kind {
985 continue
986 }
987 select {
988 case <-j.done:
989 default:
990 running++
991 }
992 }
993 running += m.reservations[key]
994 if running >= limit {
995 m.mu.Unlock()
996 return func() {}, running, false
997 }
998 m.reservations[key]++
999 m.mu.Unlock()
1000
1001 var once sync.Once
1002 release = func() {
1003 once.Do(func() {
1004 m.mu.Lock()
1005 m.reservations[key]--
1006 if m.reservations[key] == 0 {
1007 delete(m.reservations, key)
1008 }
1009 m.mu.Unlock()
1010 })
1011 }
1012 return release, running, true
1013 }
1014
1015 // HasUnfinishedForSession reports whether parentSession owns any job whose
1016 // goroutine has not fully exited yet. Empty parentSession preserves the legacy
1017 // unscoped behavior.
1018 func (m *Manager) HasUnfinishedForSession(parentSession string) bool {
1019 m.mu.Lock()
1020 defer m.mu.Unlock()
1021 for _, key := range m.order {
1022 j := m.jobs[key]
1023 if !sessionMatches(parentSession, j.SessionID) {
1024 continue
1025 }
1026 select {
1027 case <-j.done:
1028 default:
1029 return true
1030 }
1031 }
1032 return false
1033 }
1034
1035 // DrainCompletedNote returns (and clears) a one-line summary of jobs that
1036 // finished since the last drain, for the controller to fold into the next turn
1037 // so the model learns of completions. "" when nothing finished.
1038 func (m *Manager) DrainCompletedNote() string {
1039 return m.DrainCompletedNoteForSession("")
1040 }
1041
1042 // DrainCompletedNoteForSession drains completion notes for parentSession only.
1043 // Notes for other sessions stay queued until that session becomes active again.
1044 // Empty parentSession preserves the legacy unscoped behavior.
1045 func (m *Manager) DrainCompletedNoteForSession(parentSession string) string {
1046 m.mu.Lock()
1047 var c []string
1048 if strings.TrimSpace(parentSession) == "" {
1049 for _, item := range m.completed {
1050 c = append(c, item.text)
1051 }
1052 m.completed = nil
1053 } else {
1054 remaining := m.completed[:0]
1055 for _, item := range m.completed {
1056 if item.sessionID == parentSession {
1057 c = append(c, item.text)
1058 } else {
1059 remaining = append(remaining, item)
1060 }
1061 }
1062 m.completed = remaining
1063 }
1064 m.mu.Unlock()
1065 if len(c) == 0 {
1066 return ""
1067 }
1068 return "Background job updates since your last message: " + strings.Join(c, "; ") +
1069 ". Read their output with job_output if you still need it."
1070 }
1071
1072 // SetActiveSession controls which session receives lifecycle notices for jobs
1073 // that finish asynchronously. Empty active session preserves legacy behavior.
1074 func (m *Manager) SetActiveSession(parentSession string) {
1075 m.mu.Lock()
1076 m.active = strings.TrimSpace(parentSession)
1077 m.mu.Unlock()
1078 }
1079
1080 // validateTrustedSessionPath performs defense-in-depth syntax validation on a
1081 // transcript path already trusted by the store/controller layer. It rejects
1082 // control characters, but deliberately preserves separators and `..`: those are
1083 // valid host-path syntax, and rejecting them without a trusted root would break
1084 // legitimate relative paths without establishing filesystem containment.
1085 func validateTrustedSessionPath(sessionPath string) error {
1086 if sessionPath == "" {
1087 return fmt.Errorf("jobs: sessionPath must not be empty")
1088 }
1089 for i, r := range sessionPath {
1090 if r < 0x20 || r == 0x7f {
1091 return fmt.Errorf("jobs: sessionPath contains control character 0x%02x at index %d", r, i)
1092 }
1093 }
1094 return nil
1095 }
1096
1097 // SetActiveSessionPath binds a parent session id to its persistent transcript
1098 // path, migrates any temporary artifacts, and loads completed job tombstones from
1099 // the session sidecar. sessionPath must come from the trusted store/controller
1100 // path; this method does not establish filesystem containment on its own.
1101 func (m *Manager) SetActiveSessionPath(parentSession, sessionPath string) {
1102 defer func() { m.notifyRuntime("", "") }()
1103 parentSession = strings.TrimSpace(parentSession)
1104 sessionPath = strings.TrimSpace(sessionPath)
1105 // Preserve the legacy active-only behavior for calls without a complete
1106 // binding. In particular, an empty path is not an error or filesystem input.
1107 if parentSession == "" || sessionPath == "" {
1108 m.mu.Lock()
1109 m.active = parentSession
1110 m.mu.Unlock()
1111 return
1112 }
1113 // Reject malformed trusted paths before any filesystem side effect. This is
1114 // syntax hardening, not a boundary for arbitrary caller-controlled paths.
1115 if err := validateTrustedSessionPath(sessionPath); err != nil {
1116 m.mu.Lock()
1117 m.active = parentSession
1118 // A rejected rebinding must not leave future jobs writing to a stale
1119 // transcript that happened to use the same parent session id.
1120 delete(m.artifactDirs, parentSession)
1121 delete(m.loaded, parentSession)
1122 m.mu.Unlock()
1123 m.boundSink().Emit(event.Event{
1124 Kind: event.Notice,
1125 Level: event.LevelWarn,
1126 Text: "Ignoring SetActiveSessionPath with invalid session path",
1127 Detail: fmt.Sprintf("session %q: %v", parentSession, err),
1128 })
1129 return
1130 }
1131 m.mu.Lock()
1132 m.active = parentSession
1133 oldDir := m.artifactDirLocked(parentSession)
1134 adoptDefault := false
1135 if _, hasDir := m.artifactDirs[parentSession]; !hasDir && m.hasUnscopedJobsLocked() {
1136 oldDir = m.artifactDirLocked("")
1137 adoptDefault = true
1138 }
1139 newDir := ArtifactDir(sessionPath)
1140 m.artifactDirs[parentSession] = newDir
1141 loaded := m.loaded[parentSession]
1142 m.mu.Unlock()
1143
1144 var migrationErr error
1145 if oldDir != "" && newDir != "" && oldDir != newDir {
1146 oldSession := parentSession
1147 if adoptDefault {
1148 oldSession = ""
1149 }
1150 migrationErr = m.migrateArtifactDirForSession(oldSession, oldDir, newDir)
1151 }
1152 if adoptDefault {
1153 m.mu.Lock()
1154 adopted := m.adoptUnscopedJobsLocked(parentSession)
1155 m.mu.Unlock()
1156 for _, j := range adopted {
1157 j.mu.Lock()
1158 st := j.artifactStatus
1159 if st == "" {
1160 st = j.status
1161 }
1162 if err := m.writeJobMetaLocked(j, st); err != nil {
1163 j.noteArtifactErr("ownership metadata: " + err.Error())
1164 }
1165 j.mu.Unlock()
1166 }
1167 }
1168 if migrationErr != nil {
1169 m.recordArtifactMigrationError(parentSession, migrationErr)
1170 }
1171 if !loaded {
1172 m.loadSessionArtifacts(parentSession, sessionPath, newDir)
1173 }
1174 }
1175
1176 func (m *Manager) hasUnscopedJobsLocked() bool {
1177 for _, j := range m.jobs {
1178 if j != nil && strings.TrimSpace(j.SessionID) == "" {
1179 return true
1180 }
1181 }
1182 return false
1183 }
1184
1185 func (m *Manager) adoptUnscopedJobsLocked(parentSession string) []*Job {
1186 var adopted []*Job
1187 parentSession = strings.TrimSpace(parentSession)
1188 if parentSession == "" {
1189 return adopted
1190 }
1191 for i := range m.completed {
1192 if strings.TrimSpace(m.completed[i].sessionID) == "" {
1193 m.completed[i].sessionID = parentSession
1194 }
1195 }
1196 for oldKey, j := range m.jobs {
1197 if j == nil || strings.TrimSpace(j.SessionID) != "" {
1198 continue
1199 }
1200 newKey := jobKey(parentSession, j.ID)
1201 if existing := m.jobs[newKey]; existing != nil && existing != j {
1202 j.mu.Lock()
1203 j.artifactErr = "migration: job id collision while adopting temporary session"
1204 j.artifactComplete = false
1205 j.mu.Unlock()
1206 continue
1207 }
1208 delete(m.jobs, oldKey)
1209 // Manager readers and artifact writers use distinct locks. Ownership
1210 // changes hold both, always in manager-before-job order.
1211 j.mu.Lock()
1212 j.SessionID = parentSession
1213 j.mu.Unlock()
1214 adopted = append(adopted, j)
1215 m.jobs[newKey] = j
1216 for i, key := range m.order {
1217 if key == oldKey {
1218 m.order[i] = newKey
1219 }
1220 }
1221 }
1222 return adopted
1223 }
1224
1225 func (m *Manager) recordArtifactMigrationError(parentSession string, err error) {
1226 text := "job artifact migration failed: " + err.Error()
1227 m.mu.Lock()
1228 for _, j := range m.jobs {
1229 if j == nil || !sessionMatches(parentSession, j.SessionID) {
1230 continue
1231 }
1232 j.mu.Lock()
1233 if j.artifactErr == "" {
1234 j.artifactErr = "migration: " + err.Error()
1235 j.artifactComplete = false
1236 }
1237 j.mu.Unlock()
1238 }
1239 active := m.active
1240 m.mu.Unlock()
1241 if active == "" || active == parentSession {
1242 m.boundSink().Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: "Job artifact migration failed.", Detail: text})
1243 }
1244 }
1245
1246 type artifactMigrationJob struct {
1247 job *Job
1248 wasOpen bool
1249 }
1250
1251 func (m *Manager) migrateArtifactDirForSession(parentSession, oldDir, newDir string) error {
1252 locked := m.lockArtifactJobsForMigration(parentSession, oldDir)
1253 defer unlockArtifactMigrationJobs(locked)
1254 skip := openArtifactMigrationFiles(locked)
1255 migrateErr := migrateArtifactDirSkipping(oldDir, newDir, skip)
1256 if migrateErr == nil {
1257 rebaseArtifactMigrationJobs(locked, newDir)
1258 }
1259 return migrateErr
1260 }
1261
1262 func (m *Manager) lockArtifactJobsForMigration(parentSession, dir string) []artifactMigrationJob {
1263 parentSession = strings.TrimSpace(parentSession)
1264 dir = filepath.Clean(strings.TrimSpace(dir))
1265 m.mu.Lock()
1266 jobs := make([]*Job, 0, len(m.jobs))
1267 for _, j := range m.jobs {
1268 if j == nil || strings.TrimSpace(j.SessionID) != parentSession {
1269 continue
1270 }
1271 jobs = append(jobs, j)
1272 }
1273 m.mu.Unlock()
1274 sort.Slice(jobs, func(i, k int) bool {
1275 return jobs[i].ID < jobs[k].ID
1276 })
1277 locked := make([]artifactMigrationJob, 0, len(jobs))
1278 for _, j := range jobs {
1279 j.mu.Lock()
1280 if !artifactPathInDir(j.artifactPath, dir) {
1281 j.mu.Unlock()
1282 continue
1283 }
1284 locked = append(locked, artifactMigrationJob{job: j, wasOpen: j.artifactFile != nil})
1285 }
1286 return locked
1287 }
1288
1289 func artifactPathInDir(path, dir string) bool {
1290 path = filepath.Clean(strings.TrimSpace(path))
1291 dir = filepath.Clean(strings.TrimSpace(dir))
1292 if path == "." || dir == "." {
1293 return false
1294 }
1295 return filepath.Dir(path) == dir
1296 }
1297
1298 func openArtifactMigrationFiles(jobs []artifactMigrationJob) map[string]bool {
1299 skip := map[string]bool{}
1300 for _, item := range jobs {
1301 j := item.job
1302 if j == nil || j.artifactFile == nil {
1303 continue
1304 }
1305 if j.artifactPath != "" {
1306 skip[filepath.Base(j.artifactPath)] = true
1307 }
1308 if j.artifactMetaPath != "" {
1309 skip[filepath.Base(j.artifactMetaPath)] = true
1310 }
1311 }
1312 return skip
1313 }
1314
1315 func rebaseArtifactMigrationJobs(jobs []artifactMigrationJob, dir string) {
1316 for _, item := range jobs {
1317 j := item.job
1318 if j == nil || item.wasOpen {
1319 continue
1320 }
1321 if j.artifactPath != "" {
1322 j.artifactPath = filepath.Join(dir, filepath.Base(j.artifactPath))
1323 }
1324 if j.artifactMetaPath != "" {
1325 j.artifactMetaPath = filepath.Join(dir, filepath.Base(j.artifactMetaPath))
1326 }
1327 }
1328 }
1329
1330 func unlockArtifactMigrationJobs(jobs []artifactMigrationJob) {
1331 for _, v := range slices.Backward(jobs) {
1332 if v.job != nil {
1333 v.job.mu.Unlock()
1334 }
1335 }
1336 }
1337
1338 func migrateArtifactDir(src, dst string) error {
1339 return migrateArtifactDirSkipping(src, dst, nil)
1340 }
1341
1342 func migrateArtifactDirSkipping(src, dst string, skip map[string]bool) error {
1343 entries, err := os.ReadDir(src)
1344 if err != nil {
1345 if os.IsNotExist(err) {
1346 return nil
1347 }
1348 return err
1349 }
1350 if err := ensurePrivateArtifactDir(dst); err != nil {
1351 return err
1352 }
1353 for _, entry := range entries {
1354 if entry.IsDir() {
1355 continue
1356 }
1357 if skip[entry.Name()] {
1358 continue
1359 }
1360 if err := moveArtifactFile(filepath.Join(src, entry.Name()), filepath.Join(dst, entry.Name())); err != nil {
1361 return err
1362 }
1363 }
1364 _ = os.Remove(src)
1365 return nil
1366 }
1367
1368 func moveArtifactFile(src, dst string) error {
1369 // A rename preserves the source mode, so tighten legacy artifacts before
1370 // either the fast rename or the cross-device copy fallback.
1371 if err := os.Chmod(src, 0o600); err != nil {
1372 return err
1373 }
1374 if err := renamePath(src, dst); err == nil {
1375 return nil
1376 }
1377 if err := copyArtifactFile(src, dst); err != nil {
1378 return err
1379 }
1380 return os.Remove(src)
1381 }
1382
1383 func copyArtifactFile(src, dst string) error {
1384 in, err := os.Open(src)
1385 if err != nil {
1386 return err
1387 }
1388 defer in.Close()
1389 out, err := os.OpenFile(dst, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o600)
1390 if err != nil {
1391 return err
1392 }
1393 if err := out.Chmod(0o600); err != nil {
1394 _ = out.Close()
1395 _ = os.Remove(dst)
1396 return err
1397 }
1398 _, copyErr := io.Copy(out, in)
1399 closeErr := out.Close()
1400 if copyErr != nil {
1401 _ = os.Remove(dst)
1402 return copyErr
1403 }
1404 if closeErr != nil {
1405 _ = os.Remove(dst)
1406 return closeErr
1407 }
1408 return nil
1409 }
1410
1411 func (m *Manager) loadSessionArtifacts(parentSession, sessionPath, dir string) {
1412 entries, err := os.ReadDir(dir)
1413 if err != nil {
1414 m.mu.Lock()
1415 m.loaded[parentSession] = true
1416 m.mu.Unlock()
1417 return
1418 }
1419 var loaded []*Job
1420 deferredLiveOwner := false
1421 var repairErrors []string
1422 maxSeq := 0
1423 for _, entry := range entries {
1424 if entry.IsDir() || filepath.Ext(entry.Name()) != jobMetaExt {
1425 continue
1426 }
1427 metaPath := filepath.Join(dir, entry.Name())
1428 meta, err := readMeta(metaPath)
1429 if err != nil || strings.TrimSpace(meta.ID) == "" {
1430 continue
1431 }
1432 id := strings.TrimSpace(meta.ID)
1433 if seq := maxJobSeq(id); seq > maxSeq {
1434 maxSeq = seq
1435 }
1436 // A persisted Running record may belong to another manager in this
1437 // process or to another Reasonix process entirely. Only the runtime that
1438 // owns the session lease may repair an abandoned record as Interrupted.
1439 // Observers without proof of ownership defer the artifact and leave the
1440 // session reloadable for a later owned bind.
1441 if meta.Status == Running {
1442 if managerOwnerIsLive(meta.OwnerID) {
1443 deferredLiveOwner = true
1444 continue
1445 }
1446 if m.sessionOwnershipProbe == nil || !m.sessionOwnershipProbe(sessionPath) {
1447 deferredLiveOwner = true
1448 continue
1449 }
1450 meta.Status = Interrupted
1451 if meta.FinishedAt == 0 {
1452 meta.FinishedAt = nowMs()
1453 }
1454 meta.ArtifactComplete = false
1455 if err := repairArtifactMeta(metaPath, meta); err != nil {
1456 // Do not publish an in-memory Interrupted tombstone when the durable
1457 // state still says Running. Keep the session reloadable so a later bind
1458 // can retry the repair, and surface the failure instead of letting live
1459 // and machine-facing status silently disagree.
1460 deferredLiveOwner = true
1461 repairErrors = append(repairErrors, fmt.Sprintf("repair job %s metadata: %v", id, err))
1462 continue
1463 }
1464 }
1465 done := make(chan struct{})
1466 close(done)
1467 logPath := filepath.Join(dir, id+jobLogExt)
1468 if strings.TrimSpace(meta.LogPath) != "" {
1469 logPath = filepath.Join(dir, filepath.Base(meta.LogPath))
1470 }
1471 loaded = append(loaded, &Job{
1472 ID: id,
1473 Kind: meta.Kind,
1474 Label: meta.Label,
1475 SessionID: parentSession,
1476 status: meta.Status,
1477 clock: jobClock{startedAt: meta.StartedAt, finishedAt: meta.FinishedAt, activityAt: meta.FinishedAt},
1478 done: done,
1479 artifactPath: logPath,
1480 artifactMetaPath: filepath.Join(dir, id+jobMetaExt),
1481 artifactComplete: meta.ArtifactComplete,
1482 artifactErr: meta.ArtifactError,
1483 tombstone: true,
1484 evidence: mutationEvidenceFromArtifact(meta),
1485 })
1486 }
1487 if len(repairErrors) > 0 {
1488 m.boundSink().Emit(event.Event{
1489 Kind: event.Notice,
1490 Level: event.LevelWarn,
1491 Text: "Background job recovery did not complete.",
1492 Detail: strings.Join(repairErrors, "; "),
1493 })
1494 }
1495 m.mu.Lock()
1496 defer m.mu.Unlock()
1497 for _, j := range loaded {
1498 key := jobKey(parentSession, j.ID)
1499 if _, exists := m.jobs[key]; exists {
1500 continue
1501 }
1502 m.jobs[key] = j
1503 m.order = append(m.order, key)
1504 }
1505 if maxSeq > m.seq {
1506 m.seq = maxSeq
1507 }
1508 m.loaded[parentSession] = !deferredLiveOwner
1509 }
1510
1511 // BeginDestroySession marks a parent session as being removed from active use
1512 // and cancels its running jobs. WaitTeardown waits for the returned handle.
1513 func (m *Manager) BeginDestroySession(parentSession string) SessionTeardown {
1514 parentSession = strings.TrimSpace(parentSession)
1515 if parentSession == "" {
1516 return SessionTeardown{}
1517 }
1518 var cancels []context.CancelFunc
1519 var targets []teardownTarget
1520 m.mu.Lock()
1521 m.destroying[parentSession] = true
1522 remaining := m.completed[:0]
1523 for _, item := range m.completed {
1524 if item.sessionID != parentSession {
1525 remaining = append(remaining, item)
1526 }
1527 }
1528 m.completed = remaining
1529 for _, key := range m.order {
1530 j := m.jobs[key]
1531 if !sessionMatches(parentSession, j.SessionID) {
1532 continue
1533 }
1534 j.mu.Lock()
1535 switch j.status {
1536 case Running:
1537 j.status = Killed
1538 cancels = append(cancels, j.cancel)
1539 targets = append(targets, teardownTarget{info: TeardownJob{ID: j.ID, Kind: j.Kind, Label: j.Label}, done: j.done})
1540 case Killed:
1541 targets = append(targets, teardownTarget{info: TeardownJob{ID: j.ID, Kind: j.Kind, Label: j.Label}, done: j.done})
1542 }
1543 j.mu.Unlock()
1544 }
1545 m.mu.Unlock()
1546 for _, cancel := range cancels {
1547 cancel()
1548 }
1549 return SessionTeardown{SessionID: parentSession, targets: targets}
1550 }
1551
1552 // DestroySession preserves the legacy channel-based destroy API.
1553 func (m *Manager) DestroySession(parentSession string) []<-chan struct{} {
1554 return m.BeginDestroySession(parentSession).DoneChannels()
1555 }
1556
1557 // WaitTeardown waits for a destroy handle to unwind up to grace. A timed-out
1558 // result means the caller should defer physical cleanup until the jobs exit.
1559 func (m *Manager) WaitTeardown(ctx context.Context, h SessionTeardown, grace time.Duration) TeardownResult {
1560 result, timedOut := waitTeardownTargets(ctx, h.targets, grace)
1561 if timedOut {
1562 m.emitTeardownTimeout("destroy session "+h.SessionID, result)
1563 }
1564 return result
1565 }
1566
1567 // IsDestroying reports whether parentSession is in the destroy window. Empty
1568 // parent sessions are never considered destroyed.
1569 func (m *Manager) IsDestroying(parentSession string) bool {
1570 parentSession = strings.TrimSpace(parentSession)
1571 if parentSession == "" {
1572 return false
1573 }
1574 m.mu.Lock()
1575 defer m.mu.Unlock()
1576 return m.destroying[parentSession]
1577 }
1578
1579 // FinishDestroySession ends the destroy window after all owned jobs have unwound
1580 // and persistent cleanup/move work has completed.
1581 func (m *Manager) FinishDestroySession(parentSession string) {
1582 parentSession = strings.TrimSpace(parentSession)
1583 if parentSession == "" {
1584 return
1585 }
1586 m.mu.Lock()
1587 delete(m.destroying, parentSession)
1588 delete(m.artifactDirs, parentSession)
1589 delete(m.loaded, parentSession)
1590 m.purgeSessionLocked(parentSession)
1591 m.mu.Unlock()
1592 }
1593
1594 func (m *Manager) purgeSessionLocked(parentSession string) {
1595 kept := m.order[:0]
1596 for _, key := range m.order {
1597 j := m.jobs[key]
1598 if j == nil || sessionMatches(parentSession, j.SessionID) {
1599 delete(m.jobs, key)
1600 continue
1601 }
1602 kept = append(kept, key)
1603 }
1604 m.order = kept
1605 }
1606
1607 func waitTeardownTargets(ctx context.Context, targets []teardownTarget, grace time.Duration, allDone ...<-chan struct{}) (TeardownResult, bool) {
1608 if ctx == nil {
1609 ctx = context.Background()
1610 }
1611 start := time.Now()
1612 var timeout <-chan time.Time
1613 if grace >= 0 {
1614 timer := time.NewTimer(grace)
1615 defer timer.Stop()
1616 timeout = timer.C
1617 }
1618 if len(allDone) > 0 && allDone[0] != nil {
1619 select {
1620 case <-allDone[0]:
1621 return TeardownResult{}, false
1622 case <-ctx.Done():
1623 return teardownTimedOut(targets, time.Since(start)), false
1624 case <-timeout:
1625 return teardownTimedOut(targets, time.Since(start)), true
1626 }
1627 }
1628 for _, target := range targets {
1629 select {
1630 case <-target.done:
1631 case <-ctx.Done():
1632 return teardownTimedOut(targets, time.Since(start)), false
1633 case <-timeout:
1634 return teardownTimedOut(targets, time.Since(start)), true
1635 }
1636 }
1637 return TeardownResult{}, false
1638 }
1639
1640 func teardownTimedOut(targets []teardownTarget, waited time.Duration) TeardownResult {
1641 var out []TeardownJob
1642 for _, target := range targets {
1643 select {
1644 case <-target.done:
1645 continue
1646 default:
1647 }
1648 info := target.info
1649 info.Waited = waited
1650 out = append(out, info)
1651 }
1652 return TeardownResult{TimedOut: out}
1653 }
1654
1655 func (m *Manager) closeTargets() []teardownTarget {
1656 m.mu.Lock()
1657 defer m.mu.Unlock()
1658 var targets []teardownTarget
1659 for _, key := range m.order {
1660 j := m.jobs[key]
1661 if j == nil {
1662 continue
1663 }
1664 select {
1665 case <-j.done:
1666 continue
1667 default:
1668 }
1669 j.mu.Lock()
1670 switch j.status {
1671 case Running:
1672 j.status = Killed
1673 targets = append(targets, teardownTarget{info: TeardownJob{ID: j.ID, Kind: j.Kind, Label: j.Label}, done: j.done})
1674 case Killed:
1675 targets = append(targets, teardownTarget{info: TeardownJob{ID: j.ID, Kind: j.Kind, Label: j.Label}, done: j.done})
1676 }
1677 j.mu.Unlock()
1678 }
1679 return targets
1680 }
1681
1682 func (m *Manager) emitTeardownTimeout(action string, result TeardownResult) {
1683 if len(result.TimedOut) == 0 {
1684 return
1685 }
1686 var b strings.Builder
1687 fmt.Fprintf(&b, "background job teardown timed out during %s", strings.TrimSpace(action))
1688 for i, job := range result.TimedOut {
1689 if i == 0 {
1690 b.WriteString(": ")
1691 } else {
1692 b.WriteString("; ")
1693 }
1694 fmt.Fprintf(&b, "%s kind=%s", job.ID, job.Kind)
1695 if strings.TrimSpace(job.Label) != "" {
1696 fmt.Fprintf(&b, " label=%q", job.Label)
1697 }
1698 if job.Waited > 0 {
1699 fmt.Fprintf(&b, " waited=%s", job.Waited.Round(time.Millisecond))
1700 }
1701 }
1702 m.boundSink().Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: "Background job teardown timed out.", Detail: b.String()})
1703 }
1704
1705 func (m *Manager) removeTempRoot() {
1706 if m.tempRoot != "" {
1707 _ = os.RemoveAll(m.tempRoot)
1708 }
1709 }
1710
1711 func nowMs() int64 { return time.Now().UnixMilli() }
1712
1713 func startedText(kind, id, label string) string {
1714 if label != "" {
1715 return fmt.Sprintf("background %s started: %s (%s)", kind, id, label)
1716 }
1717 return fmt.Sprintf("background %s started: %s", kind, id)
1718 }
1719
1720 func (m *Manager) emitIfActive(parentSession string, ev event.Event) {
1721 m.mu.Lock()
1722 active := m.active
1723 m.mu.Unlock()
1724 if active == "" || strings.TrimSpace(parentSession) == "" || active == strings.TrimSpace(parentSession) {
1725 m.boundSink().Emit(ev)
1726 }
1727 }
1728
1729 func sessionMatches(filter, jobSession string) bool {
1730 filter = strings.TrimSpace(filter)
1731 return filter == "" || strings.TrimSpace(jobSession) == filter
1732 }
1733
1734 func jobKey(parentSession, id string) string {
1735 return strings.TrimSpace(parentSession) + "\x00" + strings.TrimSpace(id)
1736 }
1737
1738 // call-context injection (mirrors agent.CallContext)
1739
1740 type ctxKey struct{}
1741 type sessionCtxKey struct{}
1742 type jobCtxKey struct{}
1743 type noManager struct{}
1744
1745 // WithManager stamps ctx with the job manager so tools can reach it via
1746 // FromContext. The agent sets this on every tool call's context.
1747 func WithManager(ctx context.Context, m *Manager) context.Context {
1748 return context.WithValue(ctx, ctxKey{}, m)
1749 }
1750
1751 // WithoutManager shadows an ancestor manager. Child agents without an owned
1752 // Jobs manager must not operate the parent's background jobs by inheritance.
1753 func WithoutManager(ctx context.Context) context.Context {
1754 return context.WithValue(ctx, ctxKey{}, noManager{})
1755 }
1756
1757 // FromContext returns the job manager set by the agent, if any. ok is false for a
1758 // plain context (headless tests, calls outside the run loop).
1759 func FromContext(ctx context.Context) (*Manager, bool) {
1760 m, ok := ctx.Value(ctxKey{}).(*Manager)
1761 return m, ok && m != nil
1762 }
1763
1764 // WithSession stamps ctx with the active parent session ID for session-scoped job
1765 // operations.
1766 func WithSession(ctx context.Context, parentSession string) context.Context {
1767 return context.WithValue(ctx, sessionCtxKey{}, strings.TrimSpace(parentSession))
1768 }
1769
1770 // SessionFromContext returns the active parent session ID for job ownership and
1771 // filtering. Empty means no session scope is available.
1772 func SessionFromContext(ctx context.Context) string {
1773 session, _ := ctx.Value(sessionCtxKey{}).(string)
1774 return strings.TrimSpace(session)
1775 }
1776
1777 // PublishEvidence attaches a background agent's host-observed receipts to its
1778 // job. The receipts stay independent of the parent turn ledger until the
1779 // parent collects the terminal result with wait or bash_output.
1780 func PublishEvidence(ctx context.Context, summary evidence.ChildEvidenceSummary) {
1781 j, _ := ctx.Value(jobCtxKey{}).(*Job)
1782 if j == nil || len(summary.Receipts) == 0 {
1783 return
1784 }
1785 j.mu.Lock()
1786 mergePublishedEvidence(&j.evidence, summary)
1787 j.mu.Unlock()
1788 }
1789
1790 // LeaseEvidenceForSession returns a copy of a terminal job's evidence without
1791 // consuming it. Collection is only provisional: the receipts merge into the
1792 // collecting turn's ledger, but that ledger is discarded if the turn is
1793 // cancelled, errors, or the process exits before the turn commits. Consuming
1794 // here would then lose the mutation for good — the parent's next turn resets its
1795 // ledger and this job would report nothing, so a background change would ship
1796 // unreviewed. The evidence is drained only by CommitEvidenceForSession, which
1797 // the agent calls after the collecting turn passes its delivery gates. A
1798 // committed job returns empty so a re-poll after successful delivery does not
1799 // re-demand review.
1800 func (m *Manager) LeaseEvidenceForSession(parentSession, id string) evidence.ChildEvidenceSummary {
1801 summary, _ := m.tryLeaseEvidenceForSession(parentSession, id)
1802 return summary
1803 }
1804
1805 // TryLeaseEvidenceForSession is LeaseEvidenceForSession plus a ready flag that
1806 // separates "terminal evidence available" (possibly empty — a committed job or
1807 // one with no mutations) from "not ready to lease yet": unknown job, still
1808 // running, or killed but its run goroutine has not yet flushed PublishEvidence
1809 // and closed done. KillForSession flips status to Killed synchronously, well
1810 // before the goroutine actually returns, so a bash_output poll that lands in
1811 // that window must not treat the empty read as final. Callers that record a
1812 // lease (collectBackgroundEvidence) must gate on ready so they never note a
1813 // lease before the evidence exists — noting it early would let a later commit
1814 // drain evidence nobody ever merged or reviewed.
1815 func (m *Manager) TryLeaseEvidenceForSession(parentSession, id string) (evidence.ChildEvidenceSummary, bool) {
1816 return m.tryLeaseEvidenceForSession(parentSession, id)
1817 }
1818
1819 func (m *Manager) tryLeaseEvidenceForSession(parentSession, id string) (evidence.ChildEvidenceSummary, bool) {
1820 j := m.get(parentSession, id)
1821 if j == nil {
1822 return evidence.ChildEvidenceSummary{}, false
1823 }
1824 j.mu.Lock()
1825 defer j.mu.Unlock()
1826 select {
1827 case <-j.done:
1828 default:
1829 return evidence.ChildEvidenceSummary{}, false
1830 }
1831 if j.evidenceCommitted {
1832 return evidence.ChildEvidenceSummary{}, true
1833 }
1834 out := make([]evidence.Receipt, len(j.evidence.Receipts))
1835 copy(out, j.evidence.Receipts)
1836 return evidence.ChildEvidenceSummary{Receipts: out, WorkspaceRoot: j.evidence.WorkspaceRoot}, true
1837 }
1838
1839 // PendingEvidenceJobIDsForSession returns the IDs of parentSession's terminal
1840 // jobs that carry uncommitted mutation evidence — a prior turn leased it but
1841 // never delivered (the turn failed or was cancelled, and the next turn's Reset
1842 // wiped it from the per-turn ledger), or the process restarted before any turn
1843 // collected it at all. The agent re-leases these at the start of every turn so
1844 // a turn that never calls wait/bash_output still surfaces the pending mutation
1845 // to its final-readiness checks instead of silently shipping it unreviewed.
1846 func (m *Manager) PendingEvidenceJobIDsForSession(parentSession string) []string {
1847 m.mu.Lock()
1848 defer m.mu.Unlock()
1849 var ids []string
1850 for _, key := range m.order {
1851 j := m.jobs[key]
1852 if j == nil || !sessionMatches(parentSession, j.SessionID) {
1853 continue
1854 }
1855 j.mu.Lock()
1856 terminal := false
1857 select {
1858 case <-j.done:
1859 terminal = true
1860 default:
1861 }
1862 pending := terminal && !j.evidenceCommitted && len(j.evidence.Receipts) > 0
1863 j.mu.Unlock()
1864 if pending {
1865 ids = append(ids, j.ID)
1866 }
1867 }
1868 return ids
1869 }
1870
1871 // CommitEvidenceForSession permanently consumes a terminal job's evidence after
1872 // the collecting turn has accounted for it (passed final-readiness). It clears
1873 // the in-memory copy and drains the persisted mutation summary so neither a
1874 // same-process re-poll nor a restart resurrects receipts the delivered turn
1875 // already reviewed. Best-effort on the disk rewrite — a failed rewrite merely
1876 // restores the conservative resurrection behavior.
1877 func (m *Manager) CommitEvidenceForSession(parentSession, id string) {
1878 j := m.get(parentSession, id)
1879 if j == nil {
1880 return
1881 }
1882 j.mu.Lock()
1883 defer j.mu.Unlock()
1884 select {
1885 case <-j.done:
1886 default:
1887 return
1888 }
1889 if j.evidenceCommitted {
1890 return
1891 }
1892 hadEvidence := len(j.evidence.Receipts) > 0
1893 j.evidenceCommitted = true
1894 j.evidence = evidence.ChildEvidenceSummary{}
1895 if hadEvidence {
1896 if err := m.writeJobMetaLocked(j, j.status); err != nil {
1897 j.noteArtifactErr("evidence drain: " + err.Error())
1898 }
1899 }
1900 }
1901
1901 lines GO