package main import ( "context" "database/sql" "encoding/hex" "encoding/json" "fmt" "io" "log" "net/http" "os" "path/filepath" "sort" "strconv" "strings" "sync" "time" _ "modernc.org/sqlite" ) const schedulerInterval = 30 * time.Second const schedulerIdleLimit = time.Hour type schedulerTorrent struct { Hash string `json:"hashString"` Name string `json:"name"` PercentDone float64 `json:"percentDone"` Status int `json:"status"` UploadedEver int64 `json:"uploadedEver"` Error int `json:"error"` } func (t schedulerTorrent) eligible() bool { return t.PercentDone >= 1 && t.Error == 0 && t.Hash != "" } type schedulerRow struct { Hash string Name string LastSeeded int64 SlotStarted int64 LastUpload int64 UploadedEver int64 Active bool } type Scheduler struct { client *Client db *sql.DB mu sync.Mutex now func() time.Time lastError string retryAt map[string]int64 } func NewScheduler(client *Client, path string) (*Scheduler, error) { if path == "" { path = "data/scheduler.db" } if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { return nil, err } db, err := sql.Open("sqlite", path) if err != nil { return nil, err } db.SetMaxOpenConns(1) for _, query := range []string{ "PRAGMA busy_timeout = 5000", "PRAGMA journal_mode = WAL", `CREATE TABLE IF NOT EXISTS scheduler_settings ( id INTEGER PRIMARY KEY CHECK (id = 1), max_active INTEGER NOT NULL CHECK (max_active BETWEEN 1 AND 100) )`, "INSERT OR IGNORE INTO scheduler_settings (id, max_active) VALUES (1, 1)", `CREATE TABLE IF NOT EXISTS scheduler_torrents ( hash TEXT PRIMARY KEY, name TEXT NOT NULL, last_seeded_at INTEGER NOT NULL DEFAULT 0, slot_started_at INTEGER NOT NULL DEFAULT 0, last_upload_at INTEGER NOT NULL DEFAULT 0, uploaded_ever INTEGER NOT NULL DEFAULT 0, active INTEGER NOT NULL DEFAULT 0 )`, `CREATE TABLE IF NOT EXISTS scheduler_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, happened_at INTEGER NOT NULL, from_hash TEXT NOT NULL, from_name TEXT NOT NULL, to_hash TEXT NOT NULL, to_name TEXT NOT NULL, reason TEXT NOT NULL )`, } { if _, err := db.Exec(query); err != nil { db.Close() return nil, err } } return &Scheduler{client: client, db: db, now: time.Now, retryAt: make(map[string]int64)}, nil } func (s *Scheduler) Close() error { return s.db.Close() } func (s *Scheduler) Run(ctx context.Context) { check := func() { checkCtx, cancel := context.WithTimeout(ctx, 25*time.Second) defer cancel() if err := s.Tick(checkCtx); err != nil { log.Printf("scheduler: %v", err) } } check() ticker := time.NewTicker(schedulerInterval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: check() } } } func (s *Scheduler) rows(ctx context.Context) ([]*schedulerRow, error) { r, err := s.db.QueryContext(ctx, `SELECT hash, name, last_seeded_at, slot_started_at, last_upload_at, uploaded_ever, active FROM scheduler_torrents ORDER BY hash`) if err != nil { return nil, err } defer r.Close() var result []*schedulerRow for r.Next() { row := new(schedulerRow) if err := r.Scan(&row.Hash, &row.Name, &row.LastSeeded, &row.SlotStarted, &row.LastUpload, &row.UploadedEver, &row.Active); err != nil { return nil, err } result = append(result, row) } return result, r.Err() } func (s *Scheduler) saveRow(ctx context.Context, row *schedulerRow) error { _, err := s.db.ExecContext(ctx, `UPDATE scheduler_torrents SET name=?, last_seeded_at=?, slot_started_at=?, last_upload_at=?, uploaded_ever=?, active=? WHERE hash=?`, row.Name, row.LastSeeded, row.SlotStarted, row.LastUpload, row.UploadedEver, row.Active, row.Hash) return err } func (s *Scheduler) event(ctx context.Context, at int64, from, to *schedulerRow, reason string) error { fromHash, fromName, toHash, toName := "", "", "", "" if from != nil { fromHash, fromName = from.Hash, from.Name } if to != nil { toHash, toName = to.Hash, to.Name } _, err := s.db.ExecContext(ctx, `INSERT INTO scheduler_events (happened_at, from_hash, from_name, to_hash, to_name, reason) VALUES (?, ?, ?, ?, ?, ?)`, at, fromHash, fromName, toHash, toName, reason) return err } func nextWaiting(rows []*schedulerRow, live map[string]schedulerTorrent, exclude map[string]bool) *schedulerRow { var waiting []*schedulerRow for _, row := range rows { if t, ok := live[row.Hash]; ok && t.eligible() && !row.Active && !exclude[row.Hash] { waiting = append(waiting, row) } } sort.Slice(waiting, func(i, j int) bool { if waiting[i].LastSeeded != waiting[j].LastSeeded { return waiting[i].LastSeeded < waiting[j].LastSeeded } return waiting[i].Hash < waiting[j].Hash }) if len(waiting) == 0 { return nil } return waiting[0] } func (s *Scheduler) activate(ctx context.Context, row *schedulerRow, t schedulerTorrent, at int64, from *schedulerRow, reason string) error { if t.Status != 6 { if err := s.client.TorrentActionHash(ctx, "torrent-start-now", row.Hash); err != nil { return fmt.Errorf("start %s: %w", row.Name, err) } } row.Active = true row.LastSeeded = at row.SlotStarted = at row.LastUpload = at row.UploadedEver = t.UploadedEver if err := s.saveRow(ctx, row); err != nil { return err } return s.event(ctx, at, from, row, reason) } func (s *Scheduler) deactivate(ctx context.Context, row *schedulerRow, t schedulerTorrent) error { if t.Status == 6 || t.Status == 5 { if err := s.client.TorrentActionHash(ctx, "torrent-stop", row.Hash); err != nil { return fmt.Errorf("stop %s: %w", row.Name, err) } } row.Active = false row.SlotStarted = 0 delete(s.retryAt, row.Hash) return s.saveRow(ctx, row) } func (s *Scheduler) Tick(ctx context.Context) error { s.mu.Lock() defer s.mu.Unlock() err := s.tickLocked(ctx) if err != nil { s.lastError = err.Error() } else { s.lastError = "" } return err } func (s *Scheduler) tickLocked(ctx context.Context) error { if s.client == nil { return nil } torrents, err := s.client.SchedulerTorrents(ctx) if err != nil { return err } live := make(map[string]schedulerTorrent, len(torrents)) for _, t := range torrents { live[strings.ToLower(t.Hash)] = t } rows, err := s.rows(ctx) if err != nil { return err } var maxActive int if err := s.db.QueryRowContext(ctx, "SELECT max_active FROM scheduler_settings WHERE id=1").Scan(&maxActive); err != nil { return err } at := s.now().Unix() active := make([]*schedulerRow, 0) for _, row := range rows { t, exists := live[row.Hash] if !exists || !t.eligible() { if row.Active { row.Active = false row.SlotStarted = 0 if err := s.saveRow(ctx, row); err != nil { return err } } continue } row.Name = t.Name if row.Active { if row.SlotStarted == 0 { row.SlotStarted, row.LastUpload = at, at } if t.UploadedEver > row.UploadedEver { row.LastUpload = at } row.UploadedEver = t.UploadedEver active = append(active, row) } else { row.UploadedEver = t.UploadedEver } if err := s.saveRow(ctx, row); err != nil { return err } } // Keep the longest-held slots when the configured limit is lowered. sort.Slice(active, func(i, j int) bool { if active[i].SlotStarted != active[j].SlotStarted { return active[i].SlotStarted < active[j].SlotStarted } return active[i].Hash < active[j].Hash }) for len(active) > maxActive { row := active[len(active)-1] if err := s.deactivate(ctx, row, live[row.Hash]); err != nil { return err } if err := s.event(ctx, at, row, nil, "slot limit"); err != nil { return err } active = active[:len(active)-1] } // Adopt already seeding opted-in torrents if a slot is open. Stop extras. for _, row := range rows { t, ok := live[row.Hash] if !ok || !t.eligible() || row.Active || t.Status != 6 { continue } if len(active) < maxActive { if err := s.activate(ctx, row, t, at, nil, "adopted"); err != nil { return err } active = append(active, row) } else if err := s.deactivate(ctx, row, t); err != nil { return err } } rotatedOut := make(map[string]bool) for _, row := range active { t := live[row.Hash] if t.Status == 6 { delete(s.retryAt, row.Hash) } if at-row.LastUpload >= int64(schedulerIdleLimit/time.Second) { next := nextWaiting(rows, live, rotatedOut) if next != nil { if err := s.deactivate(ctx, row, t); err != nil { return err } rotatedOut[row.Hash] = true if err := s.activate(ctx, next, live[next.Hash], at, row, "idle for 1 hour"); err != nil { return err } continue } } if t.Status != 6 && at >= s.retryAt[row.Hash] { if err := s.client.TorrentActionHash(ctx, "torrent-start-now", row.Hash); err != nil { return fmt.Errorf("restore %s: %w", row.Name, err) } s.retryAt[row.Hash] = at + 5*60 } } count := 0 for _, row := range rows { if row.Active { count++ } } for count < maxActive { next := nextWaiting(rows, live, rotatedOut) if next == nil { break } if err := s.activate(ctx, next, live[next.Hash], at, nil, "slot filled"); err != nil { return err } count++ } return nil } type schedulerItem struct { Hash string `json:"hashString"` Name string `json:"name"` OptedIn bool `json:"optedIn"` Active bool `json:"active"` Status int `json:"status"` Available bool `json:"available"` LastSeeded int64 `json:"lastSeededAt"` LastUpload int64 `json:"lastUploadAt"` } func (s *Scheduler) handleState(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { w.Header().Set("Allow", http.MethodGet) writeError(w, http.StatusMethodNotAllowed, "method not allowed") return } if s.client == nil { writeError(w, http.StatusBadGateway, "transmission client is not configured") return } torrents, err := s.client.SchedulerTorrents(r.Context()) if err != nil { writeError(w, http.StatusBadGateway, err.Error()) return } s.mu.Lock() defer s.mu.Unlock() rows, err := s.rows(r.Context()) if err != nil { writeError(w, http.StatusInternalServerError, err.Error()) return } byHash := make(map[string]*schedulerRow, len(rows)) for _, row := range rows { byHash[row.Hash] = row } items := make([]schedulerItem, 0) seen := make(map[string]bool) for _, t := range torrents { hash := strings.ToLower(t.Hash) if !t.eligible() && byHash[hash] == nil { continue } row := byHash[hash] item := schedulerItem{Hash: hash, Name: t.Name, Status: t.Status, Available: t.eligible()} if row != nil { item.OptedIn, item.Active = true, row.Active item.LastSeeded, item.LastUpload = row.LastSeeded, row.LastUpload } items = append(items, item) seen[hash] = true } for _, row := range rows { if !seen[row.Hash] { items = append(items, schedulerItem{Hash: row.Hash, Name: row.Name, OptedIn: true, LastSeeded: row.LastSeeded, LastUpload: row.LastUpload}) } } sort.Slice(items, func(i, j int) bool { return strings.ToLower(items[i].Name) < strings.ToLower(items[j].Name) }) var maxActive int var latest int64 if err := s.db.QueryRowContext(r.Context(), "SELECT max_active FROM scheduler_settings WHERE id=1").Scan(&maxActive); err != nil { writeError(w, http.StatusInternalServerError, err.Error()) return } if err := s.db.QueryRowContext(r.Context(), "SELECT COALESCE(MAX(id),0) FROM scheduler_events").Scan(&latest); err != nil { writeError(w, http.StatusInternalServerError, err.Error()) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(struct { MaxActive int `json:"maxActive"` Torrents []schedulerItem `json:"torrents"` LatestEventID int64 `json:"latestEventId"` Error string `json:"error"` }{maxActive, items, latest, s.lastError}) } func decodeSchedulerBody(r *http.Request, dest any) error { decoder := json.NewDecoder(io.LimitReader(r.Body, 1024)) decoder.DisallowUnknownFields() if err := decoder.Decode(dest); err != nil { return err } if err := decoder.Decode(&struct{}{}); err != io.EOF { return fmt.Errorf("request body must contain one JSON object") } return nil } func (s *Scheduler) handleSettings(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { w.Header().Set("Allow", http.MethodPost) writeError(w, http.StatusMethodNotAllowed, "method not allowed") return } var body struct { MaxActive *int `json:"maxActive"` } if err := decodeSchedulerBody(r, &body); err != nil || body.MaxActive == nil || *body.MaxActive < 1 || *body.MaxActive > 100 { writeError(w, http.StatusBadRequest, "maxActive must be an integer from 1 to 100") return } s.mu.Lock() _, err := s.db.ExecContext(r.Context(), "UPDATE scheduler_settings SET max_active=? WHERE id=1", *body.MaxActive) s.mu.Unlock() if err != nil { writeError(w, http.StatusInternalServerError, err.Error()) return } if err := s.Tick(r.Context()); err != nil { log.Printf("scheduler settings: %v", err) } writeJSON(w, json.RawMessage(`{"ok":true}`)) } func (s *Scheduler) handleOptIn(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { w.Header().Set("Allow", http.MethodPost) writeError(w, http.StatusMethodNotAllowed, "method not allowed") return } path := strings.TrimPrefix(r.URL.Path, "/api/scheduler/torrents/") parts := strings.Split(path, "/") if len(parts) != 2 || parts[1] != "opt-in" { writeError(w, http.StatusNotFound, "not found") return } hash := strings.ToLower(parts[0]) if len(hash) != 40 { writeError(w, http.StatusBadRequest, "invalid torrent hash") return } if _, err := hex.DecodeString(hash); err != nil { writeError(w, http.StatusBadRequest, "invalid torrent hash") return } var body struct { Enabled *bool `json:"enabled"` } if err := decodeSchedulerBody(r, &body); err != nil || body.Enabled == nil { writeError(w, http.StatusBadRequest, "enabled must be a boolean") return } if s.client == nil { writeError(w, http.StatusBadGateway, "transmission client is not configured") return } if *body.Enabled { torrents, err := s.client.SchedulerTorrents(r.Context()) if err != nil { writeError(w, http.StatusBadGateway, err.Error()) return } var selected *schedulerTorrent for i := range torrents { if strings.EqualFold(torrents[i].Hash, hash) { selected = &torrents[i] break } } if selected == nil { writeError(w, http.StatusNotFound, "torrent not found") return } if !selected.eligible() { writeError(w, http.StatusConflict, "torrent must be complete and error-free") return } s.mu.Lock() _, err = s.db.ExecContext(r.Context(), `INSERT INTO scheduler_torrents (hash, name, uploaded_ever) VALUES (?, ?, ?) ON CONFLICT(hash) DO UPDATE SET name=excluded.name`, hash, selected.Name, selected.UploadedEver) s.mu.Unlock() if err != nil { writeError(w, http.StatusInternalServerError, err.Error()) return } } else { s.mu.Lock() var active bool err := s.db.QueryRowContext(r.Context(), "SELECT active FROM scheduler_torrents WHERE hash=?", hash).Scan(&active) if err == sql.ErrNoRows { s.mu.Unlock() writeError(w, http.StatusNotFound, "torrent is not opted in") return } if err == nil && active { err = s.client.TorrentActionHash(r.Context(), "torrent-stop", hash) } if err == nil { _, err = s.db.ExecContext(r.Context(), "DELETE FROM scheduler_torrents WHERE hash=?", hash) } s.mu.Unlock() if err != nil { writeError(w, http.StatusBadGateway, err.Error()) return } } if err := s.Tick(r.Context()); err != nil { log.Printf("scheduler opt-in: %v", err) } writeJSON(w, json.RawMessage(`{"ok":true}`)) } type schedulerEvent struct { ID int64 `json:"id"` HappenedAt int64 `json:"happenedAt"` FromHash string `json:"fromHash"` FromName string `json:"fromName"` ToHash string `json:"toHash"` ToName string `json:"toName"` Reason string `json:"reason"` } func (s *Scheduler) handleHistory(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { w.Header().Set("Allow", http.MethodGet) writeError(w, http.StatusMethodNotAllowed, "method not allowed") return } before := int64(0) if raw := r.URL.Query().Get("before"); raw != "" { var err error before, err = strconv.ParseInt(raw, 10, 64) if err != nil || before <= 0 { writeError(w, http.StatusBadRequest, "before must be a positive event id") return } } query := `SELECT id, happened_at, from_hash, from_name, to_hash, to_name, reason FROM scheduler_events WHERE (?=0 OR id 50 if hasMore { events = events[:50] } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(struct { Events []schedulerEvent `json:"events"` HasMore bool `json:"hasMore"` }{events, hasMore}) }