func (this *TaskPool) PutTasks(items map[string]cache.Item) error { if len(this.tasks) != 0 { this.clean() } for _, item := range items { task, _ := item.Object.(*types.AlarmTask) for { err := this.putTask(task) if err == nil { break } if err == ErrTaskPoolFull { expireT
(task *types.AlarmTask)
| 57 | // } |
| 58 | |
| 59 | func (tp *TaskPool) putTask(task *types.AlarmTask) error { |
| 60 | select { |
| 61 | case tp.tasks <- task: |
| 62 | lg.Info("put new task into task pool, taskid:%s strategy:%s hostname:%s ip:%s", |
| 63 | task.ID, task.Strategy.Name, task.Host.Hostname, task.Host.IP) |
| 64 | return nil |
| 65 | default: |
| 66 | return ErrTaskPoolFull |
| 67 | } |
| 68 | } |
| 69 | |
| 70 | func (tp *TaskPool) getTasks(batchSize int) []*types.AlarmTask { |
| 71 | tasks := make([]*types.AlarmTask, 0) |
no outgoing calls
no test coverage detected