mirror of
https://github.com/krau/SaveAny-Bot.git
synced 2026-08-19 11:23:57 +08:00
Use the resource fingerprint as the dedup key on insert and delete. Check and set processing entries atomically instead of TOCTOU. Count only successful downloads.
88 lines
2.1 KiB
Go
88 lines
2.1 KiB
Go
package parsed
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net/http"
|
|
"sync"
|
|
"sync/atomic"
|
|
|
|
"github.com/krau/SaveAny-Bot/common/utils/netutil"
|
|
"github.com/krau/SaveAny-Bot/config"
|
|
"github.com/krau/SaveAny-Bot/core"
|
|
"github.com/krau/SaveAny-Bot/pkg/enums/tasktype"
|
|
"github.com/krau/SaveAny-Bot/pkg/parser"
|
|
"github.com/krau/SaveAny-Bot/storage"
|
|
)
|
|
|
|
var _ core.Executable = (*Task)(nil)
|
|
|
|
type Task struct {
|
|
ID string
|
|
Ctx context.Context
|
|
Stor storage.Storage
|
|
StorPath string
|
|
item *parser.Item
|
|
httpClient *http.Client // [TODO] btorrent support?
|
|
progress ProgressTracker
|
|
stream bool
|
|
|
|
totalResources int64
|
|
downloaded atomic.Int64 // downloaded resources count
|
|
totalBytes int64 // total bytes to download
|
|
downloadedBytes atomic.Int64 // downloaded bytes count
|
|
processing map[string]ResourceInfo
|
|
processingMu sync.RWMutex
|
|
}
|
|
|
|
// Title implements core.Exectable.
|
|
func (t *Task) Title() string {
|
|
return fmt.Sprintf("[%s](%s->%s:%s)", t.Type(), t.item.Title, t.Stor.Name(), t.StorPath)
|
|
}
|
|
|
|
func (t *Task) Type() tasktype.TaskType {
|
|
return tasktype.TaskTypeParseditem
|
|
}
|
|
|
|
func (t *Task) TaskID() string {
|
|
return t.ID
|
|
}
|
|
|
|
func NewTask(
|
|
id string,
|
|
ctx context.Context,
|
|
stor storage.Storage,
|
|
storPath string,
|
|
item *parser.Item,
|
|
progressTracker ProgressTracker,
|
|
) *Task {
|
|
client := netutil.DefaultParserHTTPClient()
|
|
_, ok := stor.(storage.StorageCannotStream)
|
|
stream := config.C().Stream && !ok
|
|
return &Task{
|
|
ID: id,
|
|
Ctx: ctx,
|
|
Stor: stor,
|
|
StorPath: storPath,
|
|
item: item,
|
|
totalResources: int64(len(item.Resources)),
|
|
downloaded: atomic.Int64{},
|
|
totalBytes: func() int64 {
|
|
var total int64
|
|
for _, res := range item.Resources {
|
|
if res.Size < 0 {
|
|
continue // skip resources with unknown size
|
|
}
|
|
total += res.Size
|
|
}
|
|
return total
|
|
}(),
|
|
stream: stream,
|
|
downloadedBytes: atomic.Int64{},
|
|
httpClient: client,
|
|
progress: progressTracker,
|
|
processing: make(map[string]ResourceInfo),
|
|
processingMu: sync.RWMutex{},
|
|
}
|
|
}
|