NotifyStep 通知步骤变化 (由 Agent 调用)
(stepCount int)
| 154 | |
| 155 | // NotifyStep 通知步骤变化 (由 Agent 调用) |
| 156 | func (s *Scheduler) NotifyStep(stepCount int) { |
| 157 | s.mu.RLock() |
| 158 | |
| 159 | // 复制监听器和任务列表,避免长时间持锁 |
| 160 | listeners := make([]StepCallback, len(s.stepListeners)) |
| 161 | copy(listeners, s.stepListeners) |
| 162 | |
| 163 | tasks := make([]*StepTask, 0, len(s.stepTasks)) |
| 164 | for _, task := range s.stepTasks { |
| 165 | tasks = append(tasks, task) |
| 166 | } |
| 167 | |
| 168 | s.mu.RUnlock() |
| 169 | |
| 170 | // 通知监听器 |
| 171 | for _, listener := range listeners { |
| 172 | if listener == nil { |
| 173 | continue // 跳过已取消的监听器 |
| 174 | } |
| 175 | go func(cb StepCallback) { |
| 176 | if err := cb(s.ctx, stepCount); err != nil { |
| 177 | schedulerLog.Warn(s.ctx, "step callback error", map[string]any{"step": stepCount, "error": err}) |
| 178 | } |
| 179 | }(listener) |
| 180 | } |
| 181 | |
| 182 | // 检查并触发任务 |
| 183 | for _, task := range tasks { |
| 184 | shouldTrigger := stepCount-task.LastTriggered >= task.Every |
| 185 | if !shouldTrigger { |
| 186 | continue |
| 187 | } |
| 188 | |
| 189 | // 更新触发时间 |
| 190 | s.mu.Lock() |
| 191 | task.LastTriggered = stepCount |
| 192 | s.mu.Unlock() |
| 193 | |
| 194 | // 异步执行回调 |
| 195 | go func(t *StepTask) { |
| 196 | if err := t.Callback(s.ctx, stepCount); err != nil { |
| 197 | schedulerLog.Warn(s.ctx, "step task callback error", map[string]any{"task_id": t.ID, "error": err}) |
| 198 | } |
| 199 | |
| 200 | // 通知触发 |
| 201 | if s.opts.OnTrigger != nil { |
| 202 | s.opts.OnTrigger(t.ID, fmt.Sprintf("step:%d", t.Every), TriggerKindStep) |
| 203 | } |
| 204 | }(task) |
| 205 | } |
| 206 | } |
| 207 | |
| 208 | // EveryInterval 每隔一段时间执行 |
| 209 | func (s *Scheduler) EveryInterval(interval time.Duration, callback TaskCallback) (string, error) { |