129 lines
3.5 KiB
Go
129 lines
3.5 KiB
Go
package repository
|
|
|
|
import (
|
|
"database/sql"
|
|
"errors"
|
|
"time"
|
|
)
|
|
|
|
// OcrJobRow maps ocr_jobs.
|
|
type OcrJobRow struct {
|
|
ID string
|
|
FileID string
|
|
DeviceID string
|
|
Status string
|
|
Text *string
|
|
Error *string
|
|
}
|
|
|
|
// ErrJobNotFound marks an OCR job absent or owned by another device.
|
|
var ErrJobNotFound = errors.New("ocr job not found")
|
|
|
|
// ErrNoQueuedJobs marks an OCR queue drained (no job to claim).
|
|
var ErrNoQueuedJobs = errors.New("no queued ocr jobs")
|
|
|
|
type OcrJobs struct{ DB *sql.DB }
|
|
|
|
func (o *OcrJobs) Create(jobID, deviceID, fileID string) error {
|
|
_, err := o.DB.Exec(
|
|
`INSERT INTO ocr_jobs (job_id, device_id, file_id) VALUES ($1, $2, $3)`,
|
|
jobID, deviceID, fileID,
|
|
)
|
|
return err
|
|
}
|
|
|
|
// Get returns a job scoped by device (no-rows → ErrJobNotFound).
|
|
func (o *OcrJobs) Get(deviceID, jobID string) (OcrJobRow, error) {
|
|
var row OcrJobRow
|
|
var text, errMsg sql.NullString
|
|
err := o.DB.QueryRow(
|
|
`SELECT job_id, file_id, device_id, status,
|
|
NULLIF(text, ''), NULLIF(error, '')
|
|
FROM ocr_jobs WHERE job_id = $1 AND device_id = $2`,
|
|
jobID, deviceID,
|
|
).Scan(&row.ID, &row.FileID, &row.DeviceID, &row.Status, &text, &errMsg)
|
|
if err == sql.ErrNoRows {
|
|
return OcrJobRow{}, ErrJobNotFound
|
|
}
|
|
if err != nil {
|
|
return OcrJobRow{}, err
|
|
}
|
|
if text.Valid {
|
|
row.Text = &text.String
|
|
}
|
|
if errMsg.Valid {
|
|
row.Error = &errMsg.String
|
|
}
|
|
return row, nil
|
|
}
|
|
|
|
func (o *OcrJobs) TouchProcessing(deviceID, jobID string) error {
|
|
_, err := o.DB.Exec(
|
|
`UPDATE ocr_jobs SET status = 'processing', started_at = NOW()
|
|
WHERE job_id = $1 AND device_id = $2 AND status = 'queued'`,
|
|
jobID, deviceID,
|
|
)
|
|
return err
|
|
}
|
|
|
|
// ClaimNext atomically picks the oldest queued job (FIFO) and moves it to
|
|
// `processing`. Safe for concurrent workers: `FOR UPDATE SKIP LOCKED` blocks
|
|
// the row as part of the same statement. Empty queue → ErrNoQueuedJobs.
|
|
func (o *OcrJobs) ClaimNext() (OcrJobRow, error) {
|
|
var row OcrJobRow
|
|
err := o.DB.QueryRow(
|
|
`WITH next AS (
|
|
SELECT job_id FROM ocr_jobs
|
|
WHERE status = 'queued'
|
|
ORDER BY created_at, job_id
|
|
LIMIT 1
|
|
FOR UPDATE SKIP LOCKED
|
|
)
|
|
UPDATE ocr_jobs
|
|
SET status = 'processing', started_at = NOW()
|
|
FROM next
|
|
WHERE ocr_jobs.job_id = next.job_id
|
|
RETURNING ocr_jobs.job_id, ocr_jobs.file_id, ocr_jobs.device_id, ocr_jobs.status`,
|
|
).Scan(&row.ID, &row.FileID, &row.DeviceID, &row.Status)
|
|
if err == sql.ErrNoRows {
|
|
return OcrJobRow{}, ErrNoQueuedJobs
|
|
}
|
|
if err != nil {
|
|
return OcrJobRow{}, err
|
|
}
|
|
return row, nil
|
|
}
|
|
|
|
// ResetStaleProcessing requeues jobs stuck in `processing` (worker crash /
|
|
// processus redémarré) et plus vieux que `olderThan`. Retourne le nb de jobs
|
|
// requeued.
|
|
func (o *OcrJobs) ResetStaleProcessing(olderThan time.Duration) (int64, error) {
|
|
res, err := o.DB.Exec(
|
|
`UPDATE ocr_jobs SET status = 'queued', started_at = NULL
|
|
WHERE status = 'processing' AND started_at < $1`,
|
|
time.Now().Add(-olderThan),
|
|
)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return res.RowsAffected()
|
|
}
|
|
|
|
func (o *OcrJobs) Complete(deviceID, jobID, text string) error {
|
|
_, err := o.DB.Exec(
|
|
`UPDATE ocr_jobs SET status = 'done', text = NULLIF($3, ''), started_at = COALESCE(started_at, NOW()), completed_at = NOW()
|
|
WHERE job_id = $1 AND device_id = $2`,
|
|
jobID, deviceID, text,
|
|
)
|
|
return err
|
|
}
|
|
|
|
func (o *OcrJobs) Fail(deviceID, jobID, message string) error {
|
|
_, err := o.DB.Exec(
|
|
`UPDATE ocr_jobs SET status = 'failed', error = NULLIF($3, ''), started_at = COALESCE(started_at, NOW()), completed_at = NOW()
|
|
WHERE job_id = $1 AND device_id = $2`,
|
|
jobID, deviceID, message,
|
|
)
|
|
return err
|
|
}
|