pool.go
⎇
Raw
1package worker
2
3import (
4 "context"
5 "fmt"
6 "log"
7 "time"
8
9 "vidarchive/internal/models"
10 "vidarchive/internal/service"
11)
12
13type Pool struct {
14 downloadSvc *service.DownloadService
15 workers int
16 queue chan *models.Download
17 ctx context.Context
18 cancel context.CancelFunc
19}
20
21func New(downloadSvc *service.DownloadService, workers int) *Pool {
22 ctx, cancel := context.WithCancel(context.Background())
23 return &Pool{
24 downloadSvc: downloadSvc,
25 workers: workers,
26 queue: make(chan *models.Download, 100),
27 ctx: ctx,
28 cancel: cancel,
29 }
30}
31
32func (p *Pool) Start() {
33 for i := 0; i < p.workers; i++ {
34 go p.worker(i)
35 }
36
37 // Queue checker - polls DB for queued downloads
38 go p.queueChecker()
39}
40
41func (p *Pool) Stop() {
42 p.cancel()
43}
44
45func (p *Pool) Submit(d *models.Download) {
46 select {
47 case p.queue <- d:
48 log.Printf("Download %d queued", d.ID)
49 case <-p.ctx.Done():
50 }
51}
52
53func (p *Pool) worker(id int) {
54 log.Printf("Worker %d started", id)
55 for {
56 select {
57 case d := <-p.queue:
58 if d == nil {
59 continue
60 }
61 log.Printf("Worker %d processing download %d", id, d.ID)
62 if err := p.downloadSvc.ExecuteDownload(d); err != nil {
63 log.Printf("Worker %d download %d failed: %v", id, d.ID, err)
64 } else {
65 log.Printf("Worker %d download %d completed", id, d.ID)
66 }
67 case <-p.ctx.Done():
68 log.Printf("Worker %d stopped", id)
69 return
70 }
71 }
72}
73
74func (p *Pool) queueChecker() {
75 ticker := time.NewTicker(2 * time.Second)
76 defer ticker.Stop()
77
78 for {
79 select {
80 case <-ticker.C:
81 p.checkQueue()
82 case <-p.ctx.Done():
83 return
84 }
85 }
86}
87
88func (p *Pool) checkQueue() {
89 // Check how many items are in queue vs active
90 downloads, err := p.downloadSvc.GetAll("queued", "date")
91 if err != nil {
92 log.Printf("Queue check error: %v", err)
93 return
94 }
95
96 for _, d := range downloads {
97 // Try to submit - if queue is full, it'll block briefly
98 select {
99 case p.queue <- d:
100 default:
101 // Queue is full, skip for now
102 return
103 }
104 }
105}
106
107func (p *Pool) GetQueueStatus() (queued, active, completed, failed int, err error) {
108 all, err := p.downloadSvc.GetAll("all", "date")
109 if err != nil {
110 return 0, 0, 0, 0, fmt.Errorf("get all downloads: %w", err)
111 }
112
113 for _, d := range all {
114 switch d.Status {
115 case "queued":
116 queued++
117 case "downloading":
118 active++
119 case "completed":
120 completed++
121 case "error":
122 failed++
123 }
124 }
125
126 return queued, active, completed, failed, nil
127}
128