| 1 | package jobs |
| 2 | |
| 3 | import ( |
| 4 | "errors" |
| 5 | "log/slog" |
| 6 | "sync" |
| 7 | "time" |
| 8 | |
| 9 | "reasonix/internal/event" |
| 10 | "reasonix/internal/workspacelease" |
| 11 | ) |
| 12 | |
| 13 | // SessionBackgroundScope owns resources which outlive a controller generation. |
| 14 | // Build candidates acquire a reference before borrowing them; failed candidates |
| 15 | // release only that reference. Jobs do not retain the scope themselves. |
| 16 | type SessionBackgroundScope struct { |
| 17 | Manager *Manager |
| 18 | WorkspaceLease *workspacelease.Owner |
| 19 | mu sync.Mutex |
| 20 | refs int |
| 21 | closed bool |
| 22 | } |
| 23 | |
| 24 | func NewSessionBackgroundScope(manager *Manager, lease *workspacelease.Owner) *SessionBackgroundScope { |
| 25 | return &SessionBackgroundScope{Manager: manager, WorkspaceLease: lease, refs: 1} |
| 26 | } |
| 27 | |
| 28 | func (s *SessionBackgroundScope) Acquire() error { |
| 29 | s.mu.Lock() |
| 30 | defer s.mu.Unlock() |
| 31 | if s.closed { |
| 32 | return errors.New("background scope is closed") |
| 33 | } |
| 34 | s.refs++ |
| 35 | return nil |
| 36 | } |
| 37 | |
| 38 | func (s *SessionBackgroundScope) Release(async bool) { |
| 39 | s.mu.Lock() |
| 40 | if s.refs == 0 { |
| 41 | s.mu.Unlock() |
| 42 | return |
| 43 | } |
| 44 | s.refs-- |
| 45 | closeNow := s.refs == 0 |
| 46 | if closeNow { |
| 47 | s.closed = true |
| 48 | } |
| 49 | s.mu.Unlock() |
| 50 | if closeNow { |
| 51 | if async { |
| 52 | s.Manager.CloseAsync() |
| 53 | } else { |
| 54 | s.Manager.Close() |
| 55 | } |
| 56 | } |
| 57 | } |
| 58 | |
| 59 | // Bind is called only on publication, never while staging a replacement. |
| 60 | func (s *SessionBackgroundScope) Bind(sink event.Sink, recorder TaskRecorder) { |
| 61 | s.Manager.bindingMu.Lock() |
| 62 | s.Manager.sink = sink |
| 63 | if s.Manager.taskRecorder == nil { |
| 64 | s.Manager.taskRecorder = recorder |
| 65 | } |
| 66 | s.Manager.bindingMu.Unlock() |
| 67 | } |
| 68 | |
| 69 | type Lifetime string |
| 70 | |
| 71 | const ( |
| 72 | RuntimeBound Lifetime = "runtime_bound" |
| 73 | SessionProcess Lifetime = "session_process" |
| 74 | ) |
| 75 | |
| 76 | var ErrRebuildInProgress = errors.New("background task admission is paused for model configuration replacement") |
| 77 | |
| 78 | // BeginReplacement checks and seals task registration under the same lock. |
| 79 | // Completion and cancellation remain available throughout the reservation. |
| 80 | func (m *Manager) BeginReplacement(session string) (func(), error) { |
| 81 | started := time.Now() |
| 82 | m.mu.Lock() |
| 83 | defer m.mu.Unlock() |
| 84 | if m.replacing { |
| 85 | return nil, ErrRebuildInProgress |
| 86 | } |
| 87 | if m.root.Err() != nil { |
| 88 | return nil, errors.New("background scope is closed") |
| 89 | } |
| 90 | if len(m.blockingJobsLocked(session)) > 0 { |
| 91 | return nil, errors.New("runtime-dependent background jobs are still running") |
| 92 | } |
| 93 | m.replacing = true |
| 94 | m.eventMu.Lock() |
| 95 | m.eventPaused = true |
| 96 | m.eventMu.Unlock() |
| 97 | return sync.OnceFunc(func() { |
| 98 | m.mu.Lock() |
| 99 | m.replacing = false |
| 100 | m.mu.Unlock() |
| 101 | m.eventMu.Lock() |
| 102 | m.eventPaused = false |
| 103 | m.eventMu.Unlock() |
| 104 | go func() { |
| 105 | m.drainEvents() |
| 106 | // Observers suppress candidate/old-generation callbacks while sealed. |
| 107 | // Resample after publication or rollback even if the final process |
| 108 | // exited during the build and no later job transition will occur. |
| 109 | m.notifyRuntime("", "") |
| 110 | slog.Debug("model replacement background reservation released", "phase", "release", "task_class", SessionProcess, "duration", time.Since(started)) |
| 111 | }() |
| 112 | }), nil |
| 113 | } |
| 114 | |
| 115 | func (m *Manager) BlockingJobs(session string) []View { |
| 116 | m.mu.Lock() |
| 117 | defer m.mu.Unlock() |
| 118 | return m.blockingJobsLocked(session) |
| 119 | } |
| 120 | |
| 121 | func (m *Manager) blockingJobsLocked(session string) []View { |
| 122 | out := []View{} |
| 123 | for _, key := range m.order { |
| 124 | j := m.jobs[key] |
| 125 | if !sessionMatches(session, j.SessionID) || j.lifetime == SessionProcess { |
| 126 | continue |
| 127 | } |
| 128 | select { |
| 129 | case <-j.done: |
| 130 | continue |
| 131 | default: |
| 132 | } |
| 133 | j.mu.Lock() |
| 134 | out = append(out, View{ID: j.ID, Kind: j.Kind, Label: j.Label, Status: string(Running), StartedAt: j.clock.startedAt}) |
| 135 | j.mu.Unlock() |
| 136 | } |
| 137 | return out |
| 138 | } |
| 139 | |
| 140 | func (m *Manager) boundSink() event.Sink { |
| 141 | return m |
| 142 | } |
| 143 | |
| 144 | // Emit queues lifecycle notices across replacement. Only one drainer invokes |
| 145 | // the current sink, always outside registry/binding locks. |
| 146 | func (m *Manager) Emit(e event.Event) { |
| 147 | m.eventMu.Lock() |
| 148 | m.eventQueue = append(m.eventQueue, e) |
| 149 | m.eventMu.Unlock() |
| 150 | m.drainEvents() |
| 151 | } |
| 152 | |
| 153 | func (m *Manager) drainEvents() { |
| 154 | m.eventMu.Lock() |
| 155 | if m.eventPaused || m.eventDraining { |
| 156 | m.eventMu.Unlock() |
| 157 | return |
| 158 | } |
| 159 | m.eventDraining = true |
| 160 | for len(m.eventQueue) > 0 && !m.eventPaused { |
| 161 | e := m.eventQueue[0] |
| 162 | m.eventQueue[0] = event.Event{} |
| 163 | m.eventQueue = m.eventQueue[1:] |
| 164 | m.eventMu.Unlock() |
| 165 | m.bindingMu.RLock() |
| 166 | sink := m.sink |
| 167 | m.bindingMu.RUnlock() |
| 168 | if sink != nil { |
| 169 | sink.Emit(e) |
| 170 | } |
| 171 | m.eventMu.Lock() |
| 172 | } |
| 173 | m.eventDraining = false |
| 174 | m.eventMu.Unlock() |
| 175 | } |
| 176 | |
| 177 | func (m *Manager) boundRecorder() TaskRecorder { |
| 178 | m.bindingMu.RLock() |
| 179 | defer m.bindingMu.RUnlock() |
| 180 | return m.taskRecorder |
| 181 | } |
| 182 | |
| 183 | func (m *Manager) ActiveSessionID() string { |
| 184 | m.mu.Lock() |
| 185 | defer m.mu.Unlock() |
| 186 | return m.active |
| 187 | } |
| 188 | |
| 189 | func (m *Manager) ReplacementInProgress() bool { |
| 190 | m.mu.Lock() |
| 191 | defer m.mu.Unlock() |
| 192 | return m.replacing |
| 193 | } |
| 194 |