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 }