Живые обновления прогресса и раздел «Раздача» (live-status)
Воркер ведёт in-memory снимок телеметрии раздач (прогресс, скорость, ETA, рейтинг, сиды/пиры, отдано) под отдельным RWMutex, обновляя его на каждом тике поллинга сразу после построения byHash — без лишних вызовов qBittorrent и без хранения в БД (волатильно). qbt.Torrent дополнен полями телеметрии. Веб-UI читает снимок через узкий контракт LiveStatus: карточки активных загрузок показывают живой прогресс-бар (htmx-поллинг фрагмента every 3s, точечно — без сброса фильтров), на странице загрузки появилась секция «Раздача» для сидирующих задач. Начальный кадр рендерится сразу со значениями; при отсутствии данных UI деградирует штатно. Капабилити live-status (OpenSpec), web-ui дополнен. Change заархивирован. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -36,6 +36,9 @@ type downloadDetailView struct {
|
||||
Confidence string
|
||||
Files []fileRow
|
||||
|
||||
// Живая статистика раздачи (заполняется из снимка воркера).
|
||||
Seeding seedingView
|
||||
|
||||
// Действия по состоянию (как на главной).
|
||||
Terminal bool
|
||||
Reviewable bool
|
||||
@@ -104,5 +107,10 @@ func (s *server) handleDownload(w http.ResponseWriter, r *http.Request) {
|
||||
view.Files = buildFileRows(rd.Plan, rd.Preview)
|
||||
}
|
||||
|
||||
// Живая статистика раздачи — со значениями уже в первом кадре; секция
|
||||
// деградирует (пустой контейнер), если торрент не сидирует/данных нет.
|
||||
l, ok := s.deps.Live.Live(d.Infohash.String)
|
||||
view.Seeding = buildSeeding(id, l, ok)
|
||||
|
||||
s.render(w, "download.html", view)
|
||||
}
|
||||
|
||||
@@ -52,6 +52,7 @@ type Deps struct {
|
||||
Commander Commander
|
||||
Reader Reader
|
||||
Reviewer Reviewer
|
||||
Live LiveStatus
|
||||
}
|
||||
|
||||
type server struct {
|
||||
@@ -80,6 +81,9 @@ func NewRouter(d Deps) (http.Handler, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if d.Live == nil {
|
||||
d.Live = noLive{} // источник телеметрии не подключён — деградируем штатно
|
||||
}
|
||||
s := &server{deps: d, tmpl: tmpl, assetVer: assetVer}
|
||||
|
||||
r := chi.NewRouter()
|
||||
@@ -95,6 +99,10 @@ func NewRouter(d Deps) (http.Handler, error) {
|
||||
// Веб-UI.
|
||||
r.Get("/", s.handleIndex)
|
||||
r.Get("/download/{id}", s.handleDownload)
|
||||
|
||||
// Живые фрагменты телеметрии (htmx-поллинг; читают снимок воркера).
|
||||
r.Get("/fragments/downloads/{id}/progress", s.handleFragProgress)
|
||||
r.Get("/fragments/downloads/{id}/seeding", s.handleFragSeeding)
|
||||
r.Post("/ui/downloads", s.handleUIAdd)
|
||||
r.Post("/ui/downloads/{id}/cancel", s.handleUICancel)
|
||||
r.Post("/ui/downloads/{id}/retry", s.handleUIRetry)
|
||||
@@ -148,12 +156,14 @@ type downloadView struct {
|
||||
SearchText string // haystack для клиентского поиска (lowercase)
|
||||
Error string
|
||||
Terminal bool
|
||||
Deleted bool // скрыт по умолчанию на главной
|
||||
Reviewable bool // review/deferred — есть экран ревью
|
||||
Undoable bool // done — можно откатить раскладку
|
||||
Relinkable bool // reverted/cancelled/target_missing — можно перепривязать заново
|
||||
Retriable bool // failed/stuck — можно повторить попытку
|
||||
Note string // пояснение рассинхрона (target_missing/orphaned/deleted)
|
||||
IsDownloading bool // активная загрузка → живой прогресс-бар + поллинг
|
||||
Progress progressView // живой прогресс (заполняется в handleIndex из снимка)
|
||||
Deleted bool // скрыт по умолчанию на главной
|
||||
Reviewable bool // review/deferred — есть экран ревью
|
||||
Undoable bool // done — можно откатить раскладку
|
||||
Relinkable bool // reverted/cancelled/target_missing — можно перепривязать заново
|
||||
Retriable bool // failed/stuck — можно повторить попытку
|
||||
Note string // пояснение рассинхрона (target_missing/orphaned/deleted)
|
||||
}
|
||||
|
||||
func (s *server) handleIndex(w http.ResponseWriter, r *http.Request) {
|
||||
@@ -165,7 +175,14 @@ func (s *server) handleIndex(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
view := indexView{Error: r.URL.Query().Get("err")}
|
||||
for _, d := range downloads {
|
||||
view.Downloads = append(view.Downloads, toView(d))
|
||||
v := toView(d)
|
||||
// Живой прогресс активных загрузок — со значениями уже в первом кадре
|
||||
// (без мигания); дальше карточка дозапрашивает фрагмент поллингом.
|
||||
if v.IsDownloading {
|
||||
l, ok := s.deps.Live.Live(d.Infohash.String)
|
||||
v.Progress = buildProgress(d.ID, true, l, ok)
|
||||
}
|
||||
view.Downloads = append(view.Downloads, v)
|
||||
}
|
||||
s.render(w, "index.html", view)
|
||||
}
|
||||
@@ -351,6 +368,7 @@ func toView(d store.Download) downloadView {
|
||||
SearchText: strings.ToLower(d.SourceRef + " " + d.Infohash.String + " " + d.Context),
|
||||
Error: d.ErrorMsg.String,
|
||||
Terminal: d.State.IsTerminal(),
|
||||
IsDownloading: d.State == store.StateDownloading,
|
||||
Deleted: d.State == store.StateDeleted,
|
||||
Reviewable: d.State == store.StateReview || d.State == store.StateDeferred,
|
||||
Undoable: d.State == store.StateDone,
|
||||
|
||||
@@ -0,0 +1,186 @@
|
||||
package httpapi
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
|
||||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||||
"git.vakhrushev.me/av/jellybit/internal/worker"
|
||||
)
|
||||
|
||||
// LiveStatus — источник живой телеметрии загрузок (снимок воркера). Контракт
|
||||
// узкий и не зависит от способа доставки в браузер (поллинг сейчас, SSE позже).
|
||||
type LiveStatus interface {
|
||||
// Live возвращает телеметрию по infohash; ok=false — данных нет (нет
|
||||
// торрента в последнем тике), читатель деградирует без живых значений.
|
||||
Live(infohash string) (worker.Live, bool)
|
||||
}
|
||||
|
||||
// noLive — заглушка на случай, когда источник телеметрии не подключён
|
||||
// (Deps.Live == nil): живых данных нет, UI деградирует штатно.
|
||||
type noLive struct{}
|
||||
|
||||
func (noLive) Live(string) (worker.Live, bool) { return worker.Live{}, false }
|
||||
|
||||
// progressView — живой прогресс активной загрузки (для карточки и фрагмента
|
||||
// /progress). Active управляется store-состоянием (downloading), а не qbt:
|
||||
// когда задача покидает downloading, фрагмент возвращается без поллинга.
|
||||
type progressView struct {
|
||||
ID int64
|
||||
Active bool // store-состояние downloading → показываем бар и поллим
|
||||
Has bool // есть данные снимка
|
||||
Percent int
|
||||
DlSpeed string
|
||||
ETA string
|
||||
}
|
||||
|
||||
// seedingView — живая статистика раздачи (для страницы и фрагмента /seeding).
|
||||
// Has истинно только если торрент сидирует и данные есть — иначе секция
|
||||
// деградирует (пустой контейнер, поллинг прекращается).
|
||||
type seedingView struct {
|
||||
ID int64
|
||||
Has bool
|
||||
Percent int
|
||||
Ratio string
|
||||
Uploaded string
|
||||
Seeds int
|
||||
Peers int
|
||||
UpSpeed string
|
||||
}
|
||||
|
||||
func buildProgress(id int64, active bool, l worker.Live, ok bool) progressView {
|
||||
v := progressView{ID: id, Active: active}
|
||||
if ok {
|
||||
v.Has = true
|
||||
v.Percent = pct(l.Progress)
|
||||
v.DlSpeed = fmtSpeed(l.DlSpeed)
|
||||
// ETA опускаем при неизвестном/sentinel (stalled) — «осталось —» уродливо;
|
||||
// шаблонный {{if .ETA}} тогда скрывает хвост строки.
|
||||
if e := fmtETA(l.ETA); e != "—" {
|
||||
v.ETA = e
|
||||
}
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
func buildSeeding(id int64, l worker.Live, ok bool) seedingView {
|
||||
v := seedingView{ID: id}
|
||||
if ok && l.Seeding {
|
||||
v.Has = true
|
||||
v.Percent = pct(l.Progress)
|
||||
v.Ratio = fmtRatio(l.Ratio)
|
||||
v.Uploaded = fmtBytes(l.Uploaded)
|
||||
v.Seeds = l.Seeds
|
||||
v.Peers = l.Peers
|
||||
v.UpSpeed = fmtSpeed(l.UpSpeed)
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
// handleFragProgress отдаёт партиал живого прогресса карточки (htmx-поллинг).
|
||||
func (s *server) handleFragProgress(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := pathID(r)
|
||||
if err != nil {
|
||||
http.Error(w, "некорректный id", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
d, err := s.deps.Reader.GetDownload(r.Context(), id)
|
||||
if err != nil {
|
||||
s.fragErr(w, err, id)
|
||||
return
|
||||
}
|
||||
active := d.State == store.StateDownloading
|
||||
l, ok := s.deps.Live.Live(d.Infohash.String)
|
||||
s.render(w, "progress", buildProgress(id, active, l, ok))
|
||||
}
|
||||
|
||||
// handleFragSeeding отдаёт партиал секции «Раздача» (htmx-поллинг).
|
||||
func (s *server) handleFragSeeding(w http.ResponseWriter, r *http.Request) {
|
||||
id, err := pathID(r)
|
||||
if err != nil {
|
||||
http.Error(w, "некорректный id", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
d, err := s.deps.Reader.GetDownload(r.Context(), id)
|
||||
if err != nil {
|
||||
s.fragErr(w, err, id)
|
||||
return
|
||||
}
|
||||
l, ok := s.deps.Live.Live(d.Infohash.String)
|
||||
s.render(w, "seeding", buildSeeding(id, l, ok))
|
||||
}
|
||||
|
||||
// fragErr транслирует ошибку чтения задачи для фрагмент-роутов: ErrNotFound →
|
||||
// 404, прочее → 500 (полная ошибка уже залогирована на доменной границе).
|
||||
func (s *server) fragErr(w http.ResponseWriter, err error, id int64) {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
http.Error(w, "не найдено", http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
s.deps.Logger.Error("live fragment", "download_id", id, "error", err)
|
||||
http.Error(w, "внутренняя ошибка", http.StatusInternalServerError)
|
||||
}
|
||||
|
||||
// --- форматирование телеметрии ---
|
||||
|
||||
// etaInfinity — sentinel qBittorrent для неизвестного/бесконечного ETA.
|
||||
const etaInfinity = 8640000
|
||||
|
||||
func pct(progress float64) int {
|
||||
if progress < 0 {
|
||||
return 0
|
||||
}
|
||||
if progress > 1 {
|
||||
return 100
|
||||
}
|
||||
return int(progress*100 + 0.5)
|
||||
}
|
||||
|
||||
// fmtBytes переводит байты в человекочитаемые единицы (двоичные, IEC).
|
||||
func fmtBytes(n int64) string {
|
||||
if n < 1024 {
|
||||
return fmt.Sprintf("%d Б", n)
|
||||
}
|
||||
const unit = 1024
|
||||
div, exp := int64(unit), 0
|
||||
units := []string{"КиБ", "МиБ", "ГиБ", "ТиБ", "ПиБ", "ЭиБ"}
|
||||
for v := n / unit; v >= unit && exp < len(units)-1; v /= unit {
|
||||
div *= unit
|
||||
exp++
|
||||
}
|
||||
return fmt.Sprintf("%.1f %s", float64(n)/float64(div), units[exp])
|
||||
}
|
||||
|
||||
// fmtSpeed форматирует скорость (байт/с).
|
||||
func fmtSpeed(n int64) string {
|
||||
if n <= 0 {
|
||||
return "0 Б/с"
|
||||
}
|
||||
return fmtBytes(n) + "/с"
|
||||
}
|
||||
|
||||
// fmtETA форматирует оценку времени; sentinel/отрицательное → «—».
|
||||
func fmtETA(sec int64) string {
|
||||
if sec < 0 || sec >= etaInfinity {
|
||||
return "—"
|
||||
}
|
||||
switch {
|
||||
case sec < 60:
|
||||
return fmt.Sprintf("%d с", sec)
|
||||
case sec < 3600:
|
||||
return fmt.Sprintf("%d мин", sec/60)
|
||||
case sec < 86400:
|
||||
return fmt.Sprintf("%d ч %d мин", sec/3600, (sec%3600)/60)
|
||||
default:
|
||||
return fmt.Sprintf("%d дн", sec/86400)
|
||||
}
|
||||
}
|
||||
|
||||
// fmtRatio форматирует рейтинг отдачи; отрицательный (sentinel) → «—».
|
||||
func fmtRatio(r float64) string {
|
||||
if r < 0 {
|
||||
return "—"
|
||||
}
|
||||
return fmt.Sprintf("%.2f", r)
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package httpapi
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||||
"git.vakhrushev.me/av/jellybit/internal/worker"
|
||||
)
|
||||
|
||||
// TestFragProgressDownloading: активная задача → фрагмент с прогрессом,
|
||||
// значениями снимка и атрибутами htmx-поллинга.
|
||||
func TestFragProgressDownloading(t *testing.T) {
|
||||
dl := store.Download{ID: 5, Infohash: store.NullString("ih5"), State: store.StateDownloading}
|
||||
lv := stubLive{m: map[string]worker.Live{"ih5": {Progress: 0.42, DlSpeed: 6400000, ETA: 720}}}
|
||||
h := testRouterLive(t, stubReader{one: &dl}, stubReviewer{}, lv)
|
||||
|
||||
rr := get(t, h, "/fragments/downloads/5/progress")
|
||||
if rr.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rr.Code)
|
||||
}
|
||||
body := rr.Body.String()
|
||||
for _, want := range []string{`hx-trigger="every 3s"`, "/fragments/downloads/5/progress", "width:42%", "42%"} {
|
||||
if !strings.Contains(body, want) {
|
||||
t.Errorf("фрагмент прогресса не содержит %q\n%s", want, body)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestFragProgressStopsWhenNotDownloading: когда задача покинула downloading,
|
||||
// фрагмент отдаётся без атрибутов поллинга (поллинг прекращается).
|
||||
func TestFragProgressStopsWhenNotDownloading(t *testing.T) {
|
||||
dl := store.Download{ID: 5, Infohash: store.NullString("ih5"), State: store.StateDone}
|
||||
h := testRouterLive(t, stubReader{one: &dl}, stubReviewer{}, stubLive{})
|
||||
|
||||
rr := get(t, h, "/fragments/downloads/5/progress")
|
||||
if rr.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rr.Code)
|
||||
}
|
||||
if body := rr.Body.String(); strings.Contains(body, "hx-trigger") {
|
||||
t.Errorf("завершённая задача всё ещё поллит:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFragSeeding: сидирующая задача → секция «Раздача» со статистикой и
|
||||
// поллингом.
|
||||
func TestFragSeeding(t *testing.T) {
|
||||
dl := store.Download{ID: 9, Infohash: store.NullString("ih9"), State: store.StateDone}
|
||||
lv := stubLive{m: map[string]worker.Live{"ih9": {
|
||||
Seeding: true, Progress: 1, Ratio: 2.41, Seeds: 38, Peers: 14, Uploaded: 1 << 30, UpSpeed: 1153433,
|
||||
}}}
|
||||
h := testRouterLive(t, stubReader{one: &dl}, stubReviewer{}, lv)
|
||||
|
||||
rr := get(t, h, "/fragments/downloads/9/seeding")
|
||||
if rr.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rr.Code)
|
||||
}
|
||||
body := rr.Body.String()
|
||||
for _, want := range []string{"Раздача", "2.41", "38 / 14", `hx-trigger="every 3s"`} {
|
||||
if !strings.Contains(body, want) {
|
||||
t.Errorf("фрагмент раздачи не содержит %q\n%s", want, body)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestFragSeedingDegrades: нет живых данных → секция отсутствует, поллинга нет.
|
||||
func TestFragSeedingDegrades(t *testing.T) {
|
||||
dl := store.Download{ID: 9, Infohash: store.NullString("ih9"), State: store.StateDone}
|
||||
h := testRouterLive(t, stubReader{one: &dl}, stubReviewer{}, stubLive{})
|
||||
|
||||
rr := get(t, h, "/fragments/downloads/9/seeding")
|
||||
if rr.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rr.Code)
|
||||
}
|
||||
body := rr.Body.String()
|
||||
if strings.Contains(body, "Раздача") || strings.Contains(body, "hx-trigger") {
|
||||
t.Errorf("секция раздачи не деградировала:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestIndexCardShowsLiveProgress: активная карточка в списке несёт прогресс уже
|
||||
// в первом кадре (значения снимка) и атрибуты поллинга.
|
||||
func TestIndexCardShowsLiveProgress(t *testing.T) {
|
||||
dl := store.Download{ID: 3, SourceRef: "The.Bear.S03", Infohash: store.NullString("ih3"), State: store.StateDownloading}
|
||||
lv := stubLive{m: map[string]worker.Live{"ih3": {Progress: 0.46, DlSpeed: 6400000, ETA: 720}}}
|
||||
h := testRouterLive(t, stubReader{list: []store.Download{dl}}, stubReviewer{}, lv)
|
||||
|
||||
rr := get(t, h, "/")
|
||||
if rr.Code != http.StatusOK {
|
||||
t.Fatalf("status = %d, want 200", rr.Code)
|
||||
}
|
||||
body := rr.Body.String()
|
||||
for _, want := range []string{`class="progress"`, "width:46%", "/fragments/downloads/3/progress"} {
|
||||
if !strings.Contains(body, want) {
|
||||
t.Errorf("карточка без живого прогресса: нет %q", want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestFragNotFound: фрагмент несуществующей задачи → 404.
|
||||
func TestFragNotFound(t *testing.T) {
|
||||
h := testRouterLive(t, stubReader{}, stubReviewer{}, stubLive{})
|
||||
if rr := get(t, h, "/fragments/downloads/404/progress"); rr.Code != http.StatusNotFound {
|
||||
t.Fatalf("status = %d, want 404", rr.Code)
|
||||
}
|
||||
}
|
||||
@@ -48,12 +48,26 @@ func (stubReviewer) ChooseCandidate(context.Context, int64, int64) error
|
||||
func (stubReviewer) SetProviderID(context.Context, int64, string, string) error { return nil }
|
||||
func (stubReviewer) ClearProvider(context.Context, int64) error { return nil }
|
||||
|
||||
// stubLive — заглушка источника живой телеметрии.
|
||||
type stubLive struct{ m map[string]worker.Live }
|
||||
|
||||
func (s stubLive) Live(infohash string) (worker.Live, bool) {
|
||||
l, ok := s.m[infohash]
|
||||
return l, ok
|
||||
}
|
||||
|
||||
func testRouter(t *testing.T, r stubReader, rv stubReviewer) http.Handler {
|
||||
t.Helper()
|
||||
return testRouterLive(t, r, rv, stubLive{})
|
||||
}
|
||||
|
||||
func testRouterLive(t *testing.T, r stubReader, rv stubReviewer, lv stubLive) http.Handler {
|
||||
t.Helper()
|
||||
h, err := NewRouter(Deps{
|
||||
Logger: slog.New(slog.NewTextHandler(io.Discard, nil)),
|
||||
Reader: r,
|
||||
Reviewer: rv,
|
||||
Live: lv,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("NewRouter: %v", err)
|
||||
|
||||
+13
-1
@@ -44,7 +44,10 @@ type Client struct {
|
||||
mu sync.Mutex // сериализует логин
|
||||
}
|
||||
|
||||
// Torrent — подмножество полей /torrents/info, нужное jellybit.
|
||||
// Torrent — подмножество полей /torrents/info, нужное jellybit. Помимо полей
|
||||
// для машины состояний и сопоставления несёт живую телеметрию (скорости, ETA,
|
||||
// статистика раздачи) — её собирает снимок воркера для веб-UI; все эти поля
|
||||
// приходят в том же ответе, отдельного вызова не нужно.
|
||||
type Torrent struct {
|
||||
Hash string `json:"hash"`
|
||||
Name string `json:"name"`
|
||||
@@ -58,6 +61,15 @@ type Torrent struct {
|
||||
AddedOn int64 `json:"added_on"`
|
||||
InfohashV1 string `json:"infohash_v1"`
|
||||
InfohashV2 string `json:"infohash_v2"`
|
||||
|
||||
// Живая телеметрия (для снимка воркера и веб-UI).
|
||||
Dlspeed int64 `json:"dlspeed"` // скорость загрузки, байт/с
|
||||
Upspeed int64 `json:"upspeed"` // скорость отдачи, байт/с
|
||||
Eta int64 `json:"eta"` // оценка до завершения, с (8640000 ≈ ∞)
|
||||
Ratio float64 `json:"ratio"` // рейтинг отдачи (может быть -1 = ∞/н/д)
|
||||
NumSeeds int `json:"num_seeds"` // подключённые сиды
|
||||
NumLeechs int `json:"num_leechs"` // подключённые личи (пиры)
|
||||
Uploaded int64 `json:"uploaded"` // отдано всего, байт
|
||||
}
|
||||
|
||||
// File — элемент /torrents/files: путь файла относительно save_path
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"git.vakhrushev.me/av/jellybit/internal/qbt"
|
||||
"git.vakhrushev.me/av/jellybit/internal/store"
|
||||
)
|
||||
|
||||
// TestPollBuildsLiveSnapshot: после Poll снимок несёт телеметрию качающейся и
|
||||
// сидирующей задач (с верным Seeding), доступную по любому из трёх хэшей;
|
||||
// неизвестный/пустой infohash → ok=false.
|
||||
func TestPollBuildsLiveSnapshot(t *testing.T) {
|
||||
qb := &fakeQbt{torrents: []qbt.Torrent{
|
||||
// Качается (Category пуст → discover не усыновляет, store не мешает).
|
||||
{Hash: "aaa", State: "downloading", Progress: 0.5, Dlspeed: 1000, Eta: 120},
|
||||
// Сидирует, торрент v2 (три ключа).
|
||||
{
|
||||
Hash: "bbb", InfohashV1: "bbb1", InfohashV2: "BBB2",
|
||||
State: "uploading", Progress: 1.0,
|
||||
Ratio: 2.5, NumSeeds: 3, NumLeechs: 1, Uploaded: 999, Upspeed: 50,
|
||||
},
|
||||
}}
|
||||
w := newTestWorker(&fakeStore{downloads: map[int64]*store.Download{}}, qb)
|
||||
|
||||
if err := w.Poll(context.Background()); err != nil {
|
||||
t.Fatalf("Poll: %v", err)
|
||||
}
|
||||
|
||||
dl, ok := w.Live("aaa")
|
||||
if !ok {
|
||||
t.Fatal("нет телеметрии качающейся задачи")
|
||||
}
|
||||
if dl.Seeding {
|
||||
t.Error("качающаяся задача помечена Seeding")
|
||||
}
|
||||
if dl.Progress != 0.5 || dl.DlSpeed != 1000 || dl.ETA != 120 {
|
||||
t.Errorf("телеметрия качания неверна: %+v", dl)
|
||||
}
|
||||
|
||||
// Сидирующая задача находится по любому из трёх хэшей (lowercase).
|
||||
for _, h := range []string{"bbb", "bbb1", "BBB2", "bbb2"} {
|
||||
sd, ok := w.Live(h)
|
||||
if !ok {
|
||||
t.Fatalf("нет телеметрии раздачи по ключу %q", h)
|
||||
}
|
||||
if !sd.Seeding {
|
||||
t.Errorf("ключ %q: Seeding=false для uploading", h)
|
||||
}
|
||||
if sd.Ratio != 2.5 || sd.Seeds != 3 || sd.Peers != 1 || sd.Uploaded != 999 || sd.UpSpeed != 50 {
|
||||
t.Errorf("ключ %q: статистика раздачи неверна: %+v", h, sd)
|
||||
}
|
||||
}
|
||||
|
||||
if _, ok := w.Live("unknownhash"); ok {
|
||||
t.Error("неизвестный infohash вернул ok=true")
|
||||
}
|
||||
if _, ok := w.Live(""); ok {
|
||||
t.Error("пустой infohash вернул ok=true")
|
||||
}
|
||||
}
|
||||
@@ -134,6 +134,40 @@ type Config struct {
|
||||
SourceMissingThreshold int
|
||||
}
|
||||
|
||||
// Live — живая телеметрия одной раздачи из снимка воркера. Курированный срез
|
||||
// qbt.Torrent: читатели (httpapi) не зависят от пакета qbt, контракт чтения
|
||||
// узкий. Seeding вычисляется воркером (classify), чтобы трактовка завершённости
|
||||
// не дублировалась в транспорте.
|
||||
type Live struct {
|
||||
Progress float64 // доля 0..1
|
||||
DlSpeed int64 // скорость загрузки, байт/с
|
||||
ETA int64 // оценка до завершения, с (8640000 ≈ ∞)
|
||||
State string // сырое состояние qBittorrent
|
||||
Seeding bool // торрент завершён и раздаётся
|
||||
|
||||
Ratio float64 // рейтинг отдачи (может быть <0 = ∞/н/д)
|
||||
Seeds int // подключённые сиды
|
||||
Peers int // подключённые личи
|
||||
Uploaded int64 // отдано всего, байт
|
||||
UpSpeed int64 // скорость отдачи, байт/с
|
||||
}
|
||||
|
||||
// liveFrom собирает Live из торрента qBittorrent (Seeding — через classify).
|
||||
func liveFrom(t qbt.Torrent) Live {
|
||||
return Live{
|
||||
Progress: t.Progress,
|
||||
DlSpeed: t.Dlspeed,
|
||||
ETA: t.Eta,
|
||||
State: t.State,
|
||||
Seeding: classify(t.State) == classReady,
|
||||
Ratio: t.Ratio,
|
||||
Seeds: t.NumSeeds,
|
||||
Peers: t.NumLeechs,
|
||||
Uploaded: t.Uploaded,
|
||||
UpSpeed: t.Upspeed,
|
||||
}
|
||||
}
|
||||
|
||||
// Worker — поллер и владелец переходов.
|
||||
type Worker struct {
|
||||
store Store
|
||||
@@ -149,6 +183,14 @@ type Worker struct {
|
||||
notifier Notifier // опц. исходящие пинги
|
||||
scanner Scanner // опц. пересканирование Jellyfin
|
||||
|
||||
// live — снимок живой телеметрии раздач (ключ — lowercase infohash, по три
|
||||
// ключа на торрент, как byHash). Обновляется атомарным свопом карты на
|
||||
// каждом тике Poll. Отдельный RWMutex (не w.mu): UI читает телеметрию часто,
|
||||
// смешивать частые чтения с замком переходов — лишняя конкуренция. Снимок
|
||||
// волатилен, в БД не хранится.
|
||||
liveMu sync.RWMutex
|
||||
live map[string]Live
|
||||
|
||||
// failNotified — дебаунс повторных EventFailed по задаче (download_id →
|
||||
// время последнего пинга). Мерцающий stalled-торрент колеблется
|
||||
// stuck↔downloading; без дебаунса каждый цикл слал бы уведомление. Память
|
||||
@@ -179,9 +221,30 @@ func New(st Store, qb QBittorrent, rec Recognizer, lay Layouter, cfg Config, log
|
||||
now: time.Now,
|
||||
newID: defaultBatchID,
|
||||
failNotified: map[int64]time.Time{},
|
||||
live: map[string]Live{},
|
||||
}
|
||||
}
|
||||
|
||||
// Live возвращает живую телеметрию раздачи по infohash (любому из v1/v2/hash).
|
||||
// ok=false, если infohash пуст или раздачи не было в последнем тике поллинга —
|
||||
// тогда читатель деградирует без живых значений. Чтение под RLock.
|
||||
func (w *Worker) Live(infohash string) (Live, bool) {
|
||||
if infohash == "" {
|
||||
return Live{}, false
|
||||
}
|
||||
w.liveMu.RLock()
|
||||
defer w.liveMu.RUnlock()
|
||||
l, ok := w.live[strings.ToLower(infohash)]
|
||||
return l, ok
|
||||
}
|
||||
|
||||
// setLive атомарно подменяет снимок телеметрии готовой картой.
|
||||
func (w *Worker) setLive(snap map[string]Live) {
|
||||
w.liveMu.Lock()
|
||||
w.live = snap
|
||||
w.liveMu.Unlock()
|
||||
}
|
||||
|
||||
// defaultBatchID — уникальный идентификатор батча раскладки.
|
||||
func defaultBatchID() string {
|
||||
return fmt.Sprintf("b-%d", time.Now().UnixNano())
|
||||
@@ -235,13 +298,20 @@ func (w *Worker) Poll(ctx context.Context) error {
|
||||
return fmt.Errorf("poll: list torrents: %w", err)
|
||||
}
|
||||
byHash := make(map[string]qbt.Torrent, len(torrents)*2)
|
||||
live := make(map[string]Live, len(torrents)*2)
|
||||
for _, t := range torrents {
|
||||
l := liveFrom(t)
|
||||
for _, h := range []string{t.Hash, t.InfohashV1, t.InfohashV2} {
|
||||
if h != "" {
|
||||
byHash[strings.ToLower(h)] = t
|
||||
key := strings.ToLower(h)
|
||||
byHash[key] = t
|
||||
live[key] = l
|
||||
}
|
||||
}
|
||||
}
|
||||
// Снимок зависит только от torrents — свопаем сразу, до store-операций
|
||||
// ниже (их ранний return по ошибке не должен лишать UI свежей телеметрии).
|
||||
w.setLive(live)
|
||||
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
Reference in New Issue
Block a user