adding scheduler
This commit is contained in:
+612
@@ -0,0 +1,612 @@
|
||||
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})
|
||||
}
|
||||
Reference in New Issue
Block a user