start to start all worker
()
| 379 | |
| 380 | // start to start all worker |
| 381 | func (q *queue) start() { |
| 382 | tasks := make(chan Task, 1) |
| 383 | |
| 384 | for { |
| 385 | // check worker number |
| 386 | q.schedule() |
| 387 | |
| 388 | select { |
| 389 | // wait worker ready |
| 390 | case <-q.ready: |
| 391 | case <-q.quit: |
| 392 | return |
| 393 | } |
| 394 | |
| 395 | // request Task from queue in background |
| 396 | q.routineGroup.Run(func() { |
| 397 | for { |
| 398 | t, err := q.scheduler.Request() |
| 399 | if t == nil || err != nil { |
| 400 | if err != nil { |
| 401 | select { |
| 402 | case <-q.quit: |
| 403 | if !errors.Is(err, ErrNoTaskInQueue) { |
| 404 | close(tasks) |
| 405 | return |
| 406 | } |
| 407 | case <-time.After(q.taskPullInterval): |
| 408 | // sleep to fetch new Task |
| 409 | } |
| 410 | } |
| 411 | } |
| 412 | if t != nil { |
| 413 | tasks <- t |
| 414 | return |
| 415 | } |
| 416 | |
| 417 | select { |
| 418 | case <-q.quit: |
| 419 | if !errors.Is(err, ErrNoTaskInQueue) { |
| 420 | close(tasks) |
| 421 | return |
| 422 | } |
| 423 | default: |
| 424 | } |
| 425 | } |
| 426 | }) |
| 427 | |
| 428 | t, ok := <-tasks |
| 429 | if !ok { |
| 430 | return |
| 431 | } |
| 432 | |
| 433 | // start new Task |
| 434 | q.metric.IncBusyWorker() |
| 435 | q.routineGroup.Run(func() { |
| 436 | q.work(t) |
| 437 | }) |
| 438 | } |