download.go
⎇
Raw
1package repository
2
3import (
4 "context"
5 "database/sql"
6 "fmt"
7
8 "vidarchive/internal/models"
9)
10
11// downloadColumns is the canonical select list, kept in the order scanDownload
12// expects so the two can't drift apart.
13const downloadColumns = `id, url, status, logs, error_message, preset_id, format_override,
14 custom_flags, output_dir, subscription_id, started_at, completed_at, created_at`
15
16type DownloadRepository struct {
17 db *sql.DB
18}
19
20func NewDownloadRepository(db *sql.DB) *DownloadRepository {
21 return &DownloadRepository{db: db}
22}
23
24func (r *DownloadRepository) Create(d *models.Download) error {
25 result, err := r.db.Exec(
26 `INSERT INTO downloads (url, status, preset_id, format_override, custom_flags, output_dir, subscription_id)
27 VALUES (?, ?, ?, ?, ?, ?, ?)`,
28 d.URL, d.Status, d.PresetID, d.FormatOverride, d.CustomFlags, d.OutputDir, d.SubscriptionID,
29 )
30 if err != nil {
31 return err
32 }
33 id, err := result.LastInsertId()
34 if err != nil {
35 return fmt.Errorf("download last insert id: %w", err)
36 }
37 d.ID = id
38 return nil
39}
40
41func (r *DownloadRepository) GetByID(id int64) (*models.Download, error) {
42 row := r.db.QueryRow(`SELECT `+downloadColumns+` FROM downloads WHERE id = ?`, id)
43 return scanDownload(row)
44}
45
46// GetAll lists one page. created_at is indexed, so the ORDER BY does not sort
47// the whole table to serve it.
48func (r *DownloadRepository) GetAll(status, sortBy string, limit, offset int) ([]*models.Download, error) {
49 query := `SELECT ` + downloadColumns + ` FROM downloads WHERE 1=1`
50 var args []interface{}
51
52 if status != "" && status != "all" {
53 query += ` AND status = ?`
54 args = append(args, status)
55 }
56
57 switch sortBy {
58 case "status":
59 query += ` ORDER BY status, created_at DESC`
60 default:
61 query += ` ORDER BY created_at DESC`
62 }
63
64 query += ` LIMIT ? OFFSET ?`
65 args = append(args, limit, offset)
66
67 return r.queryDownloads(query, args...)
68}
69
70func (r *DownloadRepository) GetQueued(limit int) ([]*models.Download, error) {
71 return r.queryDownloads(
72 `SELECT `+downloadColumns+` FROM downloads WHERE status = 'queued' ORDER BY created_at ASC LIMIT ?`,
73 limit,
74 )
75}
76
77func (r *DownloadRepository) Ping(ctx context.Context) error {
78 var one int
79 return r.db.QueryRowContext(ctx, `SELECT 1`).Scan(&one)
80}
81
82func (r *DownloadRepository) IDsByStatus(status string) ([]int64, error) {
83 rows, err := r.db.Query(`SELECT id FROM downloads WHERE status = ?`, status)
84 if err != nil {
85 return nil, err
86 }
87 defer rows.Close()
88
89 var ids []int64
90 for rows.Next() {
91 var id int64
92 if err := rows.Scan(&id); err != nil {
93 return nil, err
94 }
95 ids = append(ids, id)
96 }
97 return ids, rows.Err()
98}
99
100func (r *DownloadRepository) queryDownloads(query string, args ...interface{}) ([]*models.Download, error) {
101 rows, err := r.db.Query(query, args...)
102 if err != nil {
103 return nil, err
104 }
105 defer rows.Close()
106
107 var downloads []*models.Download
108 for rows.Next() {
109 d, err := scanDownload(rows)
110 if err != nil {
111 return nil, err
112 }
113 downloads = append(downloads, d)
114 }
115 return downloads, rows.Err()
116}
117
118func (r *DownloadRepository) AppendLogs(id int64, logs string) error {
119 _, err := r.db.Exec(
120 `UPDATE downloads SET logs = COALESCE(logs, '') || ? WHERE id = ?`,
121 logs, id,
122 )
123 return err
124}
125
126// MarkStarted claims a download atomically. False means another worker got
127// there first, so the caller must not process it again.
128func (r *DownloadRepository) MarkStarted(id int64) (bool, error) {
129 return r.affected(
130 `UPDATE downloads SET status = 'downloading', started_at = CURRENT_TIMESTAMP WHERE id = ? AND status = 'queued'`,
131 id,
132 )
133}
134
135// Requeue accepts only 'error' and 'cancelled' rows: re-queuing a running or
136// completed one would duplicate work. The bool reports whether a row changed.
137func (r *DownloadRepository) Requeue(id int64) (bool, error) {
138 return r.affected(
139 `UPDATE downloads SET status = 'queued', error_message = NULL, logs = NULL,
140 started_at = NULL, completed_at = NULL
141 WHERE id = ? AND status IN ('error', 'cancelled')`,
142 id,
143 )
144}
145
146// CancelQueued only touches a not-yet-started download. A running one is stopped
147// by cancelling its context, so this must not overwrite a row a worker owns.
148func (r *DownloadRepository) CancelQueued(id int64) (bool, error) {
149 return r.affected(
150 `UPDATE downloads SET status = 'cancelled', completed_at = CURRENT_TIMESTAMP
151 WHERE id = ? AND status = 'queued'`,
152 id,
153 )
154}
155
156func (r *DownloadRepository) affected(query string, args ...interface{}) (bool, error) {
157 res, err := r.db.Exec(query, args...)
158 if err != nil {
159 return false, err
160 }
161 n, err := res.RowsAffected()
162 if err != nil {
163 return false, err
164 }
165 return n > 0, nil
166}
167
168func (r *DownloadRepository) MarkCompleted(id int64, status string) error {
169 _, err := r.db.Exec(
170 `UPDATE downloads SET status = ?, completed_at = CURRENT_TIMESTAMP WHERE id = ?`,
171 status, id,
172 )
173 return err
174}
175
176func (r *DownloadRepository) MarkError(id int64, errMsg string) error {
177 _, err := r.db.Exec(
178 `UPDATE downloads SET status = 'error', error_message = ? WHERE id = ?`,
179 errMsg, id,
180 )
181 return err
182}
183
184// HasActiveForSubscription lets the scheduler avoid stacking a second run on top
185// of one that has not finished.
186func (r *DownloadRepository) HasActiveForSubscription(subID int64) (bool, error) {
187 var n int
188 err := r.db.QueryRow(
189 `SELECT COUNT(*) FROM downloads WHERE subscription_id = ? AND status IN ('queued', 'downloading')`,
190 subID,
191 ).Scan(&n)
192 if err != nil {
193 return false, err
194 }
195 return n > 0, nil
196}
197
198func (r *DownloadRepository) Delete(id int64) error {
199 _, err := r.db.Exec(`DELETE FROM downloads WHERE id = ?`, id)
200 return err
201}
202
203func (r *DownloadRepository) DeleteAll() error {
204 _, err := r.db.Exec(`DELETE FROM downloads`)
205 return err
206}
207
208func (r *DownloadRepository) DeleteByStatus(status string) error {
209 _, err := r.db.Exec(`DELETE FROM downloads WHERE status = ?`, status)
210 return err
211}
212
213func (r *DownloadRepository) UpdateStatusWhere(oldStatus, newStatus string) error {
214 _, err := r.db.Exec(
215 `UPDATE downloads SET status = ? WHERE status = ?`,
216 newStatus, oldStatus,
217 )
218 return err
219}
220
221func scanDownload(row interface{ Scan(...interface{}) error }) (*models.Download, error) {
222 var d models.Download
223 err := row.Scan(
224 &d.ID, &d.URL, &d.Status,
225 &d.Logs, &d.ErrorMessage, &d.PresetID, &d.FormatOverride, &d.CustomFlags, &d.OutputDir, &d.SubscriptionID,
226 &d.StartedAt, &d.CompletedAt, &d.CreatedAt,
227 )
228 if err != nil {
229 return nil, err
230 }
231 return &d, nil
232}
233