package queue import ( "database/sql" "fmt" "os" "path/filepath" "time" _ "modernc.org/sqlite" ) type Status string const ( StatusQueued Status = "queued" StatusRunning Status = "running" StatusOK Status = "ok" StatusFailed Status = "failed" StatusCancel Status = "cancelled" ) type Job struct { ID int64 StreamID int64 StreamName string App string Channel string Reason string Status Status Priority int Attempts int MaxAttempts int DeviceSerial string LastError string NotBefore time.Time CreatedAt time.Time StartedAt time.Time FinishedAt time.Time } type Store struct { DB *sql.DB } func Open(dataDir string) (*Store, error) { if err := os.MkdirAll(dataDir, 0o755); err != nil { return nil, err } path := filepath.Join(dataDir, "agent.db") dsn := fmt.Sprintf("file:%s?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)", filepath.ToSlash(path)) db, err := sql.Open("sqlite", dsn) if err != nil { return nil, err } db.SetMaxOpenConns(1) s := &Store{DB: db} if err := s.migrate(); err != nil { _ = db.Close() return nil, err } return s, nil } func (s *Store) Close() error { return s.DB.Close() } func (s *Store) migrate() error { _, err := s.DB.Exec(` CREATE TABLE IF NOT EXISTS jobs ( id INTEGER PRIMARY KEY AUTOINCREMENT, stream_id INTEGER NOT NULL, stream_name TEXT NOT NULL DEFAULT '', app TEXT NOT NULL DEFAULT '', channel TEXT NOT NULL DEFAULT '', reason TEXT NOT NULL DEFAULT '', status TEXT NOT NULL DEFAULT 'queued', priority INTEGER NOT NULL DEFAULT 100, attempts INTEGER NOT NULL DEFAULT 0, max_attempts INTEGER NOT NULL DEFAULT 5, device_serial TEXT NOT NULL DEFAULT '', last_error TEXT NOT NULL DEFAULT '', not_before TEXT NOT NULL DEFAULT '', created_at TEXT NOT NULL, started_at TEXT NOT NULL DEFAULT '', finished_at TEXT NOT NULL DEFAULT '' ); CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs(status, priority, id); CREATE INDEX IF NOT EXISTS idx_jobs_stream ON jobs(stream_id, status); `) return err } func now() string { return time.Now().UTC().Format(time.RFC3339) } func parseTime(s string) time.Time { if s == "" { return time.Time{} } t, _ := time.Parse(time.RFC3339, s) return t } // EnqueueIfIdle inserts a job unless one is already queued/running for the stream. func (s *Store) EnqueueIfIdle(streamID int64, name, app, channel, reason string, maxAttempts int) (Job, bool, error) { var existing int err := s.DB.QueryRow(` SELECT COUNT(1) FROM jobs WHERE stream_id=? AND status IN ('queued','running')`, streamID).Scan(&existing) if err != nil { return Job{}, false, err } if existing > 0 { return Job{}, false, nil } if maxAttempts <= 0 { maxAttempts = 5 } res, err := s.DB.Exec(` INSERT INTO jobs(stream_id,stream_name,app,channel,reason,status,priority,attempts,max_attempts,created_at,not_before) VALUES(?,?,?,?,?,'queued',100,0,?,?,?)`, streamID, name, app, channel, reason, maxAttempts, now(), now()) if err != nil { return Job{}, false, err } id, _ := res.LastInsertId() j, err := s.Get(id) return j, true, err } func (s *Store) Get(id int64) (Job, error) { row := s.DB.QueryRow(` SELECT id,stream_id,stream_name,app,channel,reason,status,priority,attempts,max_attempts, device_serial,last_error,not_before,created_at,started_at,finished_at FROM jobs WHERE id=?`, id) return scanJob(row) } type scannable interface { Scan(dest ...any) error } func scanJob(row scannable) (Job, error) { var j Job var status, nb, ca, sa, fa string err := row.Scan( &j.ID, &j.StreamID, &j.StreamName, &j.App, &j.Channel, &j.Reason, &status, &j.Priority, &j.Attempts, &j.MaxAttempts, &j.DeviceSerial, &j.LastError, &nb, &ca, &sa, &fa, ) if err != nil { return j, err } j.Status = Status(status) j.NotBefore = parseTime(nb) j.CreatedAt = parseTime(ca) j.StartedAt = parseTime(sa) j.FinishedAt = parseTime(fa) return j, nil } // ClaimNext marks the next ready queued job as running. Returns false if none. func (s *Store) ClaimNext() (Job, bool, error) { tx, err := s.DB.Begin() if err != nil { return Job{}, false, err } defer func() { _ = tx.Rollback() }() nowStr := now() row := tx.QueryRow(` SELECT id FROM jobs WHERE status='queued' AND (not_before='' OR not_before<=?) ORDER BY priority ASC, id ASC LIMIT 1`, nowStr) var id int64 if err := row.Scan(&id); err != nil { if err == sql.ErrNoRows { return Job{}, false, nil } return Job{}, false, err } _, err = tx.Exec(`UPDATE jobs SET status='running', started_at=?, last_error='' WHERE id=? AND status='queued'`, nowStr, id) if err != nil { return Job{}, false, err } if err := tx.Commit(); err != nil { return Job{}, false, err } j, err := s.Get(id) return j, true, err } func (s *Store) SetDevice(id int64, serial string) error { _, err := s.DB.Exec(`UPDATE jobs SET device_serial=? WHERE id=?`, serial, id) return err } func (s *Store) MarkOK(id int64) error { _, err := s.DB.Exec(`UPDATE jobs SET status='ok', finished_at=?, last_error='' WHERE id=?`, now(), id) return err } // MarkFailed increments attempts. Re-queues with backoff unless maxed out. // waitingDevice=true keeps status queued without burning an attempt. func (s *Store) MarkFailed(id int64, msg string, waitingDevice bool) error { j, err := s.Get(id) if err != nil { return err } if waitingDevice { _, err = s.DB.Exec(` UPDATE jobs SET status='queued', device_serial='', last_error=?, started_at='' WHERE id=?`, msg, id) return err } attempts := j.Attempts + 1 if attempts >= j.MaxAttempts { _, err = s.DB.Exec(` UPDATE jobs SET status='failed', attempts=?, last_error=?, finished_at=? WHERE id=?`, attempts, msg, now(), id) return err } backoff := time.Duration(1< 15*time.Minute { backoff = 15 * time.Minute } nb := time.Now().UTC().Add(backoff).Format(time.RFC3339) _, err = s.DB.Exec(` UPDATE jobs SET status='queued', attempts=?, last_error=?, not_before=?, device_serial='', started_at='' WHERE id=?`, attempts, msg, nb, id) return err } func min(a, b int) int { if a < b { return a } return b } func (s *Store) Cancel(id int64) error { _, err := s.DB.Exec(` UPDATE jobs SET status='cancelled', finished_at=? WHERE id=? AND status IN ('queued','running')`, now(), id) return err } func (s *Store) ListRecent(limit int) ([]Job, error) { if limit <= 0 { limit = 30 } rows, err := s.DB.Query(` SELECT id,stream_id,stream_name,app,channel,reason,status,priority,attempts,max_attempts, device_serial,last_error,not_before,created_at,started_at,finished_at FROM jobs ORDER BY id DESC LIMIT ?`, limit) if err != nil { return nil, err } defer rows.Close() var out []Job for rows.Next() { j, err := scanJob(rows) if err != nil { return nil, err } out = append(out, j) } return out, rows.Err() } func (s *Store) ActiveForStream(streamID int64) (bool, error) { var n int err := s.DB.QueryRow(`SELECT COUNT(1) FROM jobs WHERE stream_id=? AND status IN ('queued','running')`, streamID).Scan(&n) return n > 0, err } // RecoverStaleRunning re-queues jobs left in status=running after a crash/kill // so a down stream is not blocked forever by EnqueueIfIdle. func (s *Store) RecoverStaleRunning() (int, error) { nowStr := now() res, err := s.DB.Exec(` UPDATE jobs SET status='queued', device_serial='', started_at='', last_error='recovered after agent restart', not_before=? WHERE status='running'`, nowStr) if err != nil { return 0, err } n, _ := res.RowsAffected() return int(n), nil } // NextDelayed returns the soonest queued job that is still waiting on not_before. func (s *Store) NextDelayed() (Job, time.Duration, bool, error) { nowStr := now() row := s.DB.QueryRow(` SELECT id,stream_id,stream_name,app,channel,reason,status,priority,attempts,max_attempts, device_serial,last_error,not_before,created_at,started_at,finished_at FROM jobs WHERE status='queued' AND not_before!='' AND not_before>? ORDER BY not_before ASC LIMIT 1`, nowStr) j, err := scanJob(row) if err == sql.ErrNoRows { return Job{}, 0, false, nil } if err != nil { return Job{}, 0, false, err } until := time.Until(j.NotBefore) if until < 0 { until = 0 } return j, until, true, nil } // ClearBackoff makes a queued job runnable immediately (e.g. after fixing auth). func (s *Store) ClearBackoff(id int64) error { _, err := s.DB.Exec(`UPDATE jobs SET not_before=? WHERE id=? AND status='queued'`, now(), id) return err } // ActiveJobForStream returns the queued/running job for a stream, if any. func (s *Store) ActiveJobForStream(streamID int64) (Job, bool, error) { row := s.DB.QueryRow(` SELECT id,stream_id,stream_name,app,channel,reason,status,priority,attempts,max_attempts, device_serial,last_error,not_before,created_at,started_at,finished_at FROM jobs WHERE stream_id=? AND status IN ('queued','running') ORDER BY id DESC LIMIT 1`, streamID) j, err := scanJob(row) if err == sql.ErrNoRows { return Job{}, false, nil } if err != nil { return Job{}, false, err } return j, true, nil }