download.go
⎇
Raw
1package repository
2
3import (
4 "database/sql"
5 "fmt"
6
7 "vidarchive/internal/models"
8)
9
10type DownloadRepository struct {
11 db *sql.DB
12}
13
14func NewDownloadRepository(db *sql.DB) *DownloadRepository {
15 return &DownloadRepository{db: db}
16}
17
18func (r *DownloadRepository) Create(d *models.Download) error {
19 result, err := r.db.Exec(
20 `INSERT INTO downloads (url, status, preset_id, format_override, custom_flags, output_dir, subscription_id)
21 VALUES (?, ?, ?, ?, ?, ?, ?)`,
22 d.URL, d.Status, d.PresetID, d.FormatOverride, d.CustomFlags, d.OutputDir, d.SubscriptionID,
23 )
24 if err != nil {
25 return err
26 }
27 id, err := result.LastInsertId()
28 if err != nil {
29 return fmt.Errorf("download last insert id: %w", err)
30 }
31 d.ID = id
32 return nil
33}
34
35func (r *DownloadRepository) GetByID(id int64) (*models.Download, error) {
36 row := r.db.QueryRow(
37 `SELECT id, url, status, logs, error_message, preset_id, format_override, custom_flags, output_dir, subscription_id, started_at, completed_at, created_at
38 FROM downloads WHERE id = ?`, id,
39 )
40 return scanDownload(row)
41}
42
43func (r *DownloadRepository) GetAll(status, sortBy string) ([]*models.Download, error) {
44 query := `SELECT id, url, status, logs, error_message, preset_id, format_override, custom_flags, output_dir, subscription_id, started_at, completed_at, created_at
45 FROM downloads WHERE 1=1`
46 var args []interface{}
47
48 if status != "" && status != "all" {
49 query += ` AND status = ?`
50 args = append(args, status)
51 }
52
53 switch sortBy {
54 case "date":
55 query += ` ORDER BY created_at DESC`
56 case "status":
57 query += ` ORDER BY status, created_at DESC`
58 default:
59 query += ` ORDER BY created_at DESC`
60 }
61
62 rows, err := r.db.Query(query, args...)
63 if err != nil {
64 return nil, err
65 }
66 defer rows.Close()
67
68 var downloads []*models.Download
69 for rows.Next() {
70 d, err := scanDownload(rows)
71 if err != nil {
72 return nil, err
73 }
74 downloads = append(downloads, d)
75 }
76 return downloads, rows.Err()
77}
78
79func (r *DownloadRepository) GetQueued(limit int) ([]*models.Download, error) {
80 rows, err := r.db.Query(
81 `SELECT id, url, status, logs, error_message, preset_id, format_override, custom_flags, output_dir, subscription_id, started_at, completed_at, created_at
82 FROM downloads WHERE status = 'queued' ORDER BY created_at ASC LIMIT ?`,
83 limit,
84 )
85 if err != nil {
86 return nil, err
87 }
88 defer rows.Close()
89
90 var downloads []*models.Download
91 for rows.Next() {
92 d, err := scanDownload(rows)
93 if err != nil {
94 return nil, err
95 }
96 downloads = append(downloads, d)
97 }
98 return downloads, rows.Err()
99}
100
101func (r *DownloadRepository) UpdateStatus(id int64, status string) error {
102 _, err := r.db.Exec(`UPDATE downloads SET status = ? WHERE id = ?`, status, id)
103 return err
104}
105
106func (r *DownloadRepository) AppendLogs(id int64, logs string) error {
107 _, err := r.db.Exec(
108 `UPDATE downloads SET logs = COALESCE(logs, '') || ? WHERE id = ?`,
109 logs, id,
110 )
111 return err
112}
113
114// MarkStarted atomically transitions a download from 'queued' to 'downloading'.
115// It reports whether this call actually claimed it: false means another worker
116// already started it, so the caller must not process it again.
117func (r *DownloadRepository) MarkStarted(id int64) (bool, error) {
118 res, err := r.db.Exec(
119 `UPDATE downloads SET status = 'downloading', started_at = CURRENT_TIMESTAMP WHERE id = ? AND status = 'queued'`,
120 id,
121 )
122 if err != nil {
123 return false, err
124 }
125 n, err := res.RowsAffected()
126 if err != nil {
127 return false, err
128 }
129 return n > 0, nil
130}
131
132func (r *DownloadRepository) MarkCompleted(id int64, status string) error {
133 _, err := r.db.Exec(
134 `UPDATE downloads SET status = ?, completed_at = CURRENT_TIMESTAMP WHERE id = ?`,
135 status, id,
136 )
137 return err
138}
139
140func (r *DownloadRepository) MarkError(id int64, errMsg string) error {
141 _, err := r.db.Exec(
142 `UPDATE downloads SET status = 'error', error_message = ? WHERE id = ?`,
143 errMsg, id,
144 )
145 return err
146}
147
148// HasActiveForSubscription reports whether the given subscription already has a
149// download that is queued or in progress, so the scheduler can avoid stacking a
150// second run on top of one that hasn't finished.
151func (r *DownloadRepository) HasActiveForSubscription(subID int64) (bool, error) {
152 var n int
153 err := r.db.QueryRow(
154 `SELECT COUNT(*) FROM downloads WHERE subscription_id = ? AND status IN ('queued', 'downloading')`,
155 subID,
156 ).Scan(&n)
157 if err != nil {
158 return false, err
159 }
160 return n > 0, nil
161}
162
163func (r *DownloadRepository) Delete(id int64) error {
164 _, err := r.db.Exec(`DELETE FROM downloads WHERE id = ?`, id)
165 return err
166}
167
168func (r *DownloadRepository) DeleteAll() error {
169 _, err := r.db.Exec(`DELETE FROM downloads`)
170 return err
171}
172
173func (r *DownloadRepository) UpdateStatusWhere(oldStatus, newStatus string) error {
174 _, err := r.db.Exec(
175 `UPDATE downloads SET status = ? WHERE status = ?`,
176 newStatus, oldStatus,
177 )
178 return err
179}
180
181func scanDownload(row interface{ Scan(...interface{}) error }) (*models.Download, error) {
182 var d models.Download
183 err := row.Scan(
184 &d.ID, &d.URL, &d.Status,
185 &d.Logs, &d.ErrorMessage, &d.PresetID, &d.FormatOverride, &d.CustomFlags, &d.OutputDir, &d.SubscriptionID,
186 &d.StartedAt, &d.CompletedAt, &d.CreatedAt,
187 )
188 if err != nil {
189 return nil, err
190 }
191 return &d, nil
192}
193