Files
ivanch c98af34202
Build and Deploy (internal) / Build Transmission Manager Image (push) Successful in 3m13s
Build and Deploy (internal) / Deploy Transmission Manager (internal) (push) Successful in 3s
adding scheduler
2026-09-23 17:32:39 -03:00

613 lines
17 KiB
Go

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<?) ORDER BY id DESC LIMIT 51`
s.mu.Lock()
defer s.mu.Unlock()
rows, err := s.db.QueryContext(r.Context(), query, before, before)
if err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
defer rows.Close()
events := make([]schedulerEvent, 0)
for rows.Next() {
var event schedulerEvent
if err := rows.Scan(&event.ID, &event.HappenedAt, &event.FromHash, &event.FromName,
&event.ToHash, &event.ToName, &event.Reason); err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
events = append(events, event)
}
if err := rows.Err(); err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
hasMore := len(events) > 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})
}