1029 lines
25 KiB
Go
1029 lines
25 KiB
Go
package downloader
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"context"
|
||
|
|
"errors"
|
||
|
|
"fmt"
|
||
|
|
"log"
|
||
|
|
"os"
|
||
|
|
"sort"
|
||
|
|
"sync"
|
||
|
|
"time"
|
||
|
|
)
|
||
|
|
|
||
|
|
const (
|
||
|
|
DefaultMaxConcurrent = 2
|
||
|
|
DefaultMaxUnknownBytes = int64(4 << 30)
|
||
|
|
DefaultProgressPeriod = 250 * time.Millisecond
|
||
|
|
DefaultEventTimeout = 5 * time.Second
|
||
|
|
)
|
||
|
|
|
||
|
|
// TaskStore persists queue metadata without exposing local transfer paths.
|
||
|
|
type TaskStore interface {
|
||
|
|
Write(Task) error
|
||
|
|
LoadAll() ([]Task, error)
|
||
|
|
Delete(string) error
|
||
|
|
}
|
||
|
|
|
||
|
|
// Clock makes progress speeds and task timestamps deterministic in tests.
|
||
|
|
type Clock interface {
|
||
|
|
Now() time.Time
|
||
|
|
}
|
||
|
|
|
||
|
|
// ObserverErrorHandler records event delivery failures without changing the
|
||
|
|
// durable transfer result. Handlers should return quickly.
|
||
|
|
type ObserverErrorHandler func(Event, error)
|
||
|
|
|
||
|
|
type realClock struct{}
|
||
|
|
|
||
|
|
func (realClock) Now() time.Time {
|
||
|
|
return time.Now()
|
||
|
|
}
|
||
|
|
|
||
|
|
// QueueConfig supplies infrastructure dependencies and resource limits.
|
||
|
|
type QueueConfig struct {
|
||
|
|
DownloadsRoot string
|
||
|
|
Store TaskStore
|
||
|
|
Transport Transport
|
||
|
|
Observer Observer
|
||
|
|
Clock Clock
|
||
|
|
MaxConcurrent int
|
||
|
|
MaxUnknown int64
|
||
|
|
ProgressEvery time.Duration
|
||
|
|
EventTimeout time.Duration
|
||
|
|
OnObserverError ObserverErrorHandler
|
||
|
|
}
|
||
|
|
|
||
|
|
// EnqueueRequest creates one durable task. Total is the signed Catalog package
|
||
|
|
// size when TotalKnown is true.
|
||
|
|
type EnqueueRequest struct {
|
||
|
|
RequestID string
|
||
|
|
AppID string
|
||
|
|
URL string
|
||
|
|
TotalKnown bool
|
||
|
|
Total int64
|
||
|
|
}
|
||
|
|
|
||
|
|
type stopReason uint8
|
||
|
|
|
||
|
|
const (
|
||
|
|
stopNone stopReason = iota
|
||
|
|
stopPause
|
||
|
|
stopCancel
|
||
|
|
stopClose
|
||
|
|
)
|
||
|
|
|
||
|
|
type runtimeTask struct {
|
||
|
|
opMu sync.Mutex
|
||
|
|
|
||
|
|
task Task
|
||
|
|
running bool
|
||
|
|
blocked bool
|
||
|
|
stop stopReason
|
||
|
|
cancel context.CancelFunc
|
||
|
|
done chan struct{}
|
||
|
|
lastResult error
|
||
|
|
}
|
||
|
|
|
||
|
|
// Queue schedules durable resumable byte transfers.
|
||
|
|
type Queue struct {
|
||
|
|
root string
|
||
|
|
store TaskStore
|
||
|
|
transport Transport
|
||
|
|
observer Observer
|
||
|
|
clock Clock
|
||
|
|
maxConcurrent int
|
||
|
|
maxUnknown int64
|
||
|
|
progressEvery time.Duration
|
||
|
|
eventTimeout time.Duration
|
||
|
|
onObserverError ObserverErrorHandler
|
||
|
|
|
||
|
|
ctx context.Context
|
||
|
|
cancel context.CancelFunc
|
||
|
|
|
||
|
|
mu sync.Mutex
|
||
|
|
tasks map[string]*runtimeTask
|
||
|
|
appTasks map[string]string
|
||
|
|
order []string
|
||
|
|
active int
|
||
|
|
closed bool
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewQueue restores disk state, reconciles part/final files and starts queued
|
||
|
|
// work up to the configured concurrency limit.
|
||
|
|
func NewQueue(config QueueConfig) (*Queue, error) {
|
||
|
|
if config.DownloadsRoot == "" || config.Store == nil || config.Transport == nil {
|
||
|
|
return nil, fmt.Errorf("%w: queue dependencies are incomplete", ErrInvalidTask)
|
||
|
|
}
|
||
|
|
if config.MaxConcurrent == 0 {
|
||
|
|
config.MaxConcurrent = DefaultMaxConcurrent
|
||
|
|
}
|
||
|
|
if config.MaxConcurrent < 1 || config.MaxConcurrent > 64 {
|
||
|
|
return nil, fmt.Errorf("%w: invalid concurrency", ErrInvalidTask)
|
||
|
|
}
|
||
|
|
if config.MaxUnknown == 0 {
|
||
|
|
config.MaxUnknown = DefaultMaxUnknownBytes
|
||
|
|
}
|
||
|
|
if config.MaxUnknown < 1 {
|
||
|
|
return nil, fmt.Errorf("%w: invalid unknown-size limit", ErrInvalidTask)
|
||
|
|
}
|
||
|
|
if config.ProgressEvery == 0 {
|
||
|
|
config.ProgressEvery = DefaultProgressPeriod
|
||
|
|
}
|
||
|
|
if config.ProgressEvery < 0 {
|
||
|
|
return nil, fmt.Errorf("%w: invalid progress period", ErrInvalidTask)
|
||
|
|
}
|
||
|
|
if config.EventTimeout == 0 {
|
||
|
|
config.EventTimeout = DefaultEventTimeout
|
||
|
|
}
|
||
|
|
if config.EventTimeout < 0 {
|
||
|
|
return nil, fmt.Errorf("%w: invalid event timeout", ErrInvalidTask)
|
||
|
|
}
|
||
|
|
if config.Observer == nil {
|
||
|
|
config.Observer = discardObserver{}
|
||
|
|
}
|
||
|
|
if config.Clock == nil {
|
||
|
|
config.Clock = realClock{}
|
||
|
|
}
|
||
|
|
if config.OnObserverError == nil {
|
||
|
|
config.OnObserverError = func(event Event, err error) {
|
||
|
|
log.Printf(
|
||
|
|
"downloader observer delivery failed: type=%s request_id=%s app_id=%s: %v",
|
||
|
|
event.Type,
|
||
|
|
event.RequestID,
|
||
|
|
event.AppID,
|
||
|
|
err,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
ctx, cancel := context.WithCancel(context.Background())
|
||
|
|
queue := &Queue{
|
||
|
|
root: config.DownloadsRoot,
|
||
|
|
store: config.Store,
|
||
|
|
transport: config.Transport,
|
||
|
|
observer: config.Observer,
|
||
|
|
clock: config.Clock,
|
||
|
|
maxConcurrent: config.MaxConcurrent,
|
||
|
|
maxUnknown: config.MaxUnknown,
|
||
|
|
progressEvery: config.ProgressEvery,
|
||
|
|
eventTimeout: config.EventTimeout,
|
||
|
|
onObserverError: config.OnObserverError,
|
||
|
|
ctx: ctx,
|
||
|
|
cancel: cancel,
|
||
|
|
tasks: make(map[string]*runtimeTask),
|
||
|
|
appTasks: make(map[string]string),
|
||
|
|
}
|
||
|
|
tasks, err := config.Store.LoadAll()
|
||
|
|
if err != nil {
|
||
|
|
cancel()
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
sort.Slice(tasks, func(left, right int) bool {
|
||
|
|
if tasks[left].CreatedAt == tasks[right].CreatedAt {
|
||
|
|
return tasks[left].RequestID < tasks[right].RequestID
|
||
|
|
}
|
||
|
|
return tasks[left].CreatedAt < tasks[right].CreatedAt
|
||
|
|
})
|
||
|
|
for _, task := range tasks {
|
||
|
|
recovered, err := queue.reconcile(task)
|
||
|
|
if err != nil {
|
||
|
|
cancel()
|
||
|
|
return nil, fmt.Errorf("%w: %s: %v", ErrTaskCorrupt, task.RequestID, err)
|
||
|
|
}
|
||
|
|
if existing, exists := queue.tasks[recovered.RequestID]; exists {
|
||
|
|
cancel()
|
||
|
|
return nil, fmt.Errorf(
|
||
|
|
"%w: duplicate request %s and %s",
|
||
|
|
ErrTaskConflict,
|
||
|
|
existing.task.RequestID,
|
||
|
|
recovered.RequestID,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
if !recovered.Terminal() {
|
||
|
|
if existingID, exists := queue.appTasks[recovered.AppID]; exists {
|
||
|
|
cancel()
|
||
|
|
return nil, fmt.Errorf(
|
||
|
|
"%w: app %s has %s and %s",
|
||
|
|
ErrTaskConflict,
|
||
|
|
recovered.AppID,
|
||
|
|
existingID,
|
||
|
|
recovered.RequestID,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
queue.appTasks[recovered.AppID] = recovered.RequestID
|
||
|
|
}
|
||
|
|
queue.tasks[recovered.RequestID] = &runtimeTask{task: recovered}
|
||
|
|
queue.order = append(queue.order, recovered.RequestID)
|
||
|
|
}
|
||
|
|
queue.schedule()
|
||
|
|
return queue, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// Enqueue durably creates an idempotent task and schedules it.
|
||
|
|
func (queue *Queue) Enqueue(
|
||
|
|
_ context.Context,
|
||
|
|
request EnqueueRequest,
|
||
|
|
) (Task, error) {
|
||
|
|
task := Task{
|
||
|
|
SchemaVersion: TaskSchemaVersion,
|
||
|
|
RequestID: request.RequestID,
|
||
|
|
AppID: request.AppID,
|
||
|
|
URL: request.URL,
|
||
|
|
Status: StatusQueued,
|
||
|
|
TotalKnown: request.TotalKnown,
|
||
|
|
Total: request.Total,
|
||
|
|
CreatedAt: queue.clock.Now().UTC().Format(time.RFC3339Nano),
|
||
|
|
}
|
||
|
|
if err := task.Validate(); err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
|
||
|
|
runtime := &runtimeTask{task: task, blocked: true}
|
||
|
|
queue.mu.Lock()
|
||
|
|
if queue.closed {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return Task{}, ErrQueueClosed
|
||
|
|
}
|
||
|
|
if existing, exists := queue.tasks[task.RequestID]; exists {
|
||
|
|
existingTask := existing.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
if sameTaskIdentity(existingTask, task) {
|
||
|
|
return existingTask, nil
|
||
|
|
}
|
||
|
|
return Task{}, ErrTaskConflict
|
||
|
|
}
|
||
|
|
if _, exists := queue.appTasks[task.AppID]; exists {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return Task{}, ErrTaskConflict
|
||
|
|
}
|
||
|
|
queue.tasks[task.RequestID] = runtime
|
||
|
|
queue.appTasks[task.AppID] = task.RequestID
|
||
|
|
queue.order = append(queue.order, task.RequestID)
|
||
|
|
queue.mu.Unlock()
|
||
|
|
|
||
|
|
if err := queue.store.Write(task); err != nil {
|
||
|
|
queue.mu.Lock()
|
||
|
|
delete(queue.tasks, task.RequestID)
|
||
|
|
delete(queue.appTasks, task.AppID)
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime.blocked = false
|
||
|
|
queue.mu.Unlock()
|
||
|
|
queue.schedule()
|
||
|
|
return task, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// Snapshot returns a copy of one current task.
|
||
|
|
func (queue *Queue) Snapshot(requestID string) (Task, bool) {
|
||
|
|
queue.mu.Lock()
|
||
|
|
defer queue.mu.Unlock()
|
||
|
|
runtime, exists := queue.tasks[requestID]
|
||
|
|
if !exists {
|
||
|
|
return Task{}, false
|
||
|
|
}
|
||
|
|
return runtime.task, true
|
||
|
|
}
|
||
|
|
|
||
|
|
// Tasks returns a stable in-memory snapshot. Consumers should reconcile this
|
||
|
|
// durable state at startup or after observer reconnects because event delivery
|
||
|
|
// is best-effort and reported separately through OnObserverError.
|
||
|
|
func (queue *Queue) Tasks() []Task {
|
||
|
|
queue.mu.Lock()
|
||
|
|
defer queue.mu.Unlock()
|
||
|
|
tasks := make([]Task, 0, len(queue.tasks))
|
||
|
|
for _, requestID := range queue.order {
|
||
|
|
if runtime, exists := queue.tasks[requestID]; exists {
|
||
|
|
tasks = append(tasks, runtime.task)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return tasks
|
||
|
|
}
|
||
|
|
|
||
|
|
// Pause keeps partial bytes and moves a queued/downloading task to paused.
|
||
|
|
func (queue *Queue) Pause(ctx context.Context, requestID string) error {
|
||
|
|
runtime, err := queue.lockRuntime(requestID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
defer runtime.opMu.Unlock()
|
||
|
|
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.runtimeCurrent(requestID, runtime) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return ErrTaskNotFound
|
||
|
|
}
|
||
|
|
switch runtime.task.Status {
|
||
|
|
case StatusQueued:
|
||
|
|
runtime.blocked = true
|
||
|
|
previous := runtime.task
|
||
|
|
runtime.task.Status = StatusPaused
|
||
|
|
snapshot := runtime.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
if err := queue.store.Write(snapshot); err != nil {
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime.task = previous
|
||
|
|
runtime.blocked = false
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime.blocked = false
|
||
|
|
queue.mu.Unlock()
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventPaused,
|
||
|
|
RequestID: snapshot.RequestID,
|
||
|
|
AppID: snapshot.AppID,
|
||
|
|
Attempt: snapshot.Attempt,
|
||
|
|
Done: snapshot.Done,
|
||
|
|
})
|
||
|
|
return nil
|
||
|
|
case StatusDownloading:
|
||
|
|
runtime.stop = stopPause
|
||
|
|
cancel := runtime.cancel
|
||
|
|
done := runtime.done
|
||
|
|
queue.mu.Unlock()
|
||
|
|
cancel()
|
||
|
|
return queue.waitWorker(ctx, runtime, done)
|
||
|
|
default:
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return ErrInvalidCommand
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Resume moves a paused task back to queued.
|
||
|
|
func (queue *Queue) Resume(_ context.Context, requestID string) error {
|
||
|
|
return queue.requeue(requestID, StatusPaused)
|
||
|
|
}
|
||
|
|
|
||
|
|
// Retry preserves partial bytes and retries a failed task.
|
||
|
|
func (queue *Queue) Retry(_ context.Context, requestID string) error {
|
||
|
|
return queue.requeue(requestID, StatusFailed)
|
||
|
|
}
|
||
|
|
|
||
|
|
// Cancel removes a non-completed task, its part and its metadata.
|
||
|
|
func (queue *Queue) Cancel(ctx context.Context, requestID string) error {
|
||
|
|
runtime, err := queue.lockRuntime(requestID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
defer runtime.opMu.Unlock()
|
||
|
|
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.runtimeCurrent(requestID, runtime) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return ErrTaskNotFound
|
||
|
|
}
|
||
|
|
if runtime.task.Status == StatusCompleted {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return ErrInvalidCommand
|
||
|
|
}
|
||
|
|
if runtime.task.Status == StatusDownloading {
|
||
|
|
runtime.stop = stopCancel
|
||
|
|
cancel := runtime.cancel
|
||
|
|
done := runtime.done
|
||
|
|
queue.mu.Unlock()
|
||
|
|
cancel()
|
||
|
|
return queue.waitWorker(ctx, runtime, done)
|
||
|
|
}
|
||
|
|
runtime.blocked = true
|
||
|
|
snapshot := runtime.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
|
||
|
|
if err := queue.cleanupCanceled(snapshot); err != nil {
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime.blocked = false
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
delete(queue.tasks, requestID)
|
||
|
|
delete(queue.appTasks, snapshot.AppID)
|
||
|
|
queue.mu.Unlock()
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventCanceled,
|
||
|
|
RequestID: snapshot.RequestID,
|
||
|
|
AppID: snapshot.AppID,
|
||
|
|
Attempt: snapshot.Attempt,
|
||
|
|
})
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// Close stops workers after persisting their partial state as queued.
|
||
|
|
func (queue *Queue) Close(ctx context.Context) error {
|
||
|
|
queue.mu.Lock()
|
||
|
|
if queue.closed {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
queue.closed = true
|
||
|
|
queue.cancel()
|
||
|
|
var waits []chan struct{}
|
||
|
|
for _, runtime := range queue.tasks {
|
||
|
|
if runtime.running {
|
||
|
|
if runtime.stop == stopNone {
|
||
|
|
runtime.stop = stopClose
|
||
|
|
}
|
||
|
|
if runtime.cancel != nil {
|
||
|
|
runtime.cancel()
|
||
|
|
}
|
||
|
|
waits = append(waits, runtime.done)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
queue.mu.Unlock()
|
||
|
|
|
||
|
|
for _, done := range waits {
|
||
|
|
select {
|
||
|
|
case <-ctx.Done():
|
||
|
|
return ctx.Err()
|
||
|
|
case <-done:
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) requeue(requestID string, expected TaskStatus) error {
|
||
|
|
runtime, err := queue.lockRuntime(requestID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
defer runtime.opMu.Unlock()
|
||
|
|
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.runtimeCurrent(requestID, runtime) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return ErrTaskNotFound
|
||
|
|
}
|
||
|
|
if runtime.task.Status != expected {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return ErrInvalidCommand
|
||
|
|
}
|
||
|
|
runtime.blocked = true
|
||
|
|
previous := runtime.task
|
||
|
|
runtime.task.Status = StatusQueued
|
||
|
|
runtime.task.ErrorCode = ""
|
||
|
|
snapshot := runtime.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
if err := queue.store.Write(snapshot); err != nil {
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime.task = previous
|
||
|
|
runtime.blocked = false
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime.blocked = false
|
||
|
|
queue.mu.Unlock()
|
||
|
|
queue.schedule()
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) lockRuntime(requestID string) (*runtimeTask, error) {
|
||
|
|
queue.mu.Lock()
|
||
|
|
runtime, exists := queue.tasks[requestID]
|
||
|
|
queue.mu.Unlock()
|
||
|
|
if !exists {
|
||
|
|
return nil, ErrTaskNotFound
|
||
|
|
}
|
||
|
|
runtime.opMu.Lock()
|
||
|
|
return runtime, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) waitWorker(
|
||
|
|
ctx context.Context,
|
||
|
|
runtime *runtimeTask,
|
||
|
|
done chan struct{},
|
||
|
|
) error {
|
||
|
|
select {
|
||
|
|
case <-ctx.Done():
|
||
|
|
return ctx.Err()
|
||
|
|
case <-done:
|
||
|
|
queue.mu.Lock()
|
||
|
|
result := runtime.lastResult
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return result
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) schedule() {
|
||
|
|
for {
|
||
|
|
queue.mu.Lock()
|
||
|
|
if queue.closed || queue.active >= queue.maxConcurrent {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return
|
||
|
|
}
|
||
|
|
var selected *runtimeTask
|
||
|
|
for _, requestID := range queue.order {
|
||
|
|
runtime, exists := queue.tasks[requestID]
|
||
|
|
if !exists {
|
||
|
|
continue
|
||
|
|
}
|
||
|
|
if runtime.task.Status == StatusQueued && !runtime.running && !runtime.blocked {
|
||
|
|
selected = runtime
|
||
|
|
break
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if selected == nil {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return
|
||
|
|
}
|
||
|
|
selected.running = true
|
||
|
|
selected.stop = stopNone
|
||
|
|
selected.lastResult = nil
|
||
|
|
selected.task.Status = StatusDownloading
|
||
|
|
selected.task.Attempt++
|
||
|
|
workerCtx, cancel := context.WithCancel(queue.ctx)
|
||
|
|
selected.cancel = cancel
|
||
|
|
selected.done = make(chan struct{})
|
||
|
|
attempt := selected.task.Attempt
|
||
|
|
queue.active++
|
||
|
|
queue.mu.Unlock()
|
||
|
|
go queue.run(selected, workerCtx, attempt)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) run(runtime *runtimeTask, ctx context.Context, attempt uint64) {
|
||
|
|
completedPath, alreadyComplete, err := queue.prepareAttempt(runtime, attempt)
|
||
|
|
if err != nil {
|
||
|
|
if ctx.Err() != nil {
|
||
|
|
reason := queue.stopForAttempt(runtime, attempt)
|
||
|
|
if reason == stopCancel {
|
||
|
|
if !errors.Is(err, ctx.Err()) {
|
||
|
|
task, _ := queue.taskForAttempt(runtime, attempt)
|
||
|
|
log.Printf(
|
||
|
|
"downloader cancel cleanup proceeding after prepare error: request_id=%s: %v",
|
||
|
|
task.RequestID,
|
||
|
|
err,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
queue.finishStopped(runtime, attempt)
|
||
|
|
} else if errors.Is(err, ctx.Err()) {
|
||
|
|
queue.finishStopped(runtime, attempt)
|
||
|
|
} else {
|
||
|
|
queue.finishFailure(runtime, attempt, err)
|
||
|
|
}
|
||
|
|
return
|
||
|
|
}
|
||
|
|
queue.finishFailure(runtime, attempt, err)
|
||
|
|
return
|
||
|
|
}
|
||
|
|
if alreadyComplete {
|
||
|
|
queue.finishSuccess(runtime, attempt, completedPath)
|
||
|
|
return
|
||
|
|
}
|
||
|
|
task, valid := queue.taskForAttempt(runtime, attempt)
|
||
|
|
if !valid {
|
||
|
|
queue.finishStopped(runtime, attempt)
|
||
|
|
return
|
||
|
|
}
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventStarted,
|
||
|
|
RequestID: task.RequestID,
|
||
|
|
AppID: task.AppID,
|
||
|
|
Attempt: task.Attempt,
|
||
|
|
Done: task.Done,
|
||
|
|
TotalKnown: task.TotalKnown,
|
||
|
|
Total: task.Total,
|
||
|
|
})
|
||
|
|
completedPath, currentAttempt, err := queue.transfer(runtime, ctx, attempt)
|
||
|
|
if err != nil {
|
||
|
|
if ctx.Err() != nil {
|
||
|
|
reason := queue.stopForAttempt(runtime, currentAttempt)
|
||
|
|
if reason == stopCancel {
|
||
|
|
if !errors.Is(err, ctx.Err()) {
|
||
|
|
log.Printf(
|
||
|
|
"downloader cancel cleanup proceeding after transfer close error: request_id=%s: %v",
|
||
|
|
task.RequestID,
|
||
|
|
err,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
queue.finishStopped(runtime, currentAttempt)
|
||
|
|
} else if errors.Is(err, ctx.Err()) {
|
||
|
|
queue.finishStopped(runtime, currentAttempt)
|
||
|
|
} else {
|
||
|
|
queue.finishFailure(runtime, currentAttempt, err)
|
||
|
|
}
|
||
|
|
return
|
||
|
|
}
|
||
|
|
queue.finishFailure(runtime, currentAttempt, err)
|
||
|
|
return
|
||
|
|
}
|
||
|
|
queue.finishSuccess(runtime, currentAttempt, completedPath)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) prepareAttempt(
|
||
|
|
runtime *runtimeTask,
|
||
|
|
attempt uint64,
|
||
|
|
) (string, bool, error) {
|
||
|
|
task, valid := queue.taskForAttempt(runtime, attempt)
|
||
|
|
if !valid {
|
||
|
|
return "", false, context.Canceled
|
||
|
|
}
|
||
|
|
paths, err := ensureTransferLayout(queue.root, task.RequestID)
|
||
|
|
if err != nil {
|
||
|
|
return "", false, err
|
||
|
|
}
|
||
|
|
partSize, partExists, err := regularFileSize(paths.Part)
|
||
|
|
if err != nil {
|
||
|
|
return "", false, err
|
||
|
|
}
|
||
|
|
if finalExists, finalErr := regularFileExists(paths.Completed); finalErr != nil {
|
||
|
|
return "", false, finalErr
|
||
|
|
} else if finalExists {
|
||
|
|
return "", false, fmt.Errorf(
|
||
|
|
"%w: final file exists before transfer",
|
||
|
|
ErrTaskCorrupt,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
if task.TotalKnown && partSize > task.Total {
|
||
|
|
return "", false, fmt.Errorf(
|
||
|
|
"%w: part exceeds expected size",
|
||
|
|
ErrTaskCorrupt,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
if task.TotalKnown && partExists && partSize == task.Total {
|
||
|
|
if err := syncRegularFile(paths.Part); err != nil {
|
||
|
|
return "", false, err
|
||
|
|
}
|
||
|
|
if err := activateCompleted(paths, nil); err != nil {
|
||
|
|
return "", false, err
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.attemptCurrent(runtime, attempt) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return "", false, context.Canceled
|
||
|
|
}
|
||
|
|
runtime.task.Done = partSize
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return paths.Completed, true, nil
|
||
|
|
}
|
||
|
|
if partExists && partSize > 0 && task.Validator.Empty() {
|
||
|
|
if err := truncatePart(paths.Part); err != nil {
|
||
|
|
return "", false, err
|
||
|
|
}
|
||
|
|
partSize = 0
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.attemptCurrent(runtime, attempt) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return "", false, context.Canceled
|
||
|
|
}
|
||
|
|
runtime.task.Done = partSize
|
||
|
|
snapshot := runtime.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
if err := queue.store.Write(snapshot); err != nil {
|
||
|
|
return "", false, err
|
||
|
|
}
|
||
|
|
return "", false, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) taskForAttempt(runtime *runtimeTask, attempt uint64) (Task, bool) {
|
||
|
|
queue.mu.Lock()
|
||
|
|
defer queue.mu.Unlock()
|
||
|
|
if !queue.attemptCurrent(runtime, attempt) {
|
||
|
|
return Task{}, false
|
||
|
|
}
|
||
|
|
return runtime.task, true
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) attemptCurrent(runtime *runtimeTask, attempt uint64) bool {
|
||
|
|
return runtime.running && runtime.task.Attempt == attempt
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) stopForAttempt(
|
||
|
|
runtime *runtimeTask,
|
||
|
|
attempt uint64,
|
||
|
|
) stopReason {
|
||
|
|
queue.mu.Lock()
|
||
|
|
defer queue.mu.Unlock()
|
||
|
|
if !queue.attemptCurrent(runtime, attempt) {
|
||
|
|
return stopNone
|
||
|
|
}
|
||
|
|
return runtime.stop
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) runtimeCurrent(requestID string, runtime *runtimeTask) bool {
|
||
|
|
current, exists := queue.tasks[requestID]
|
||
|
|
return exists && current == runtime
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) publish(event Event) error {
|
||
|
|
ctx := queue.ctx
|
||
|
|
cancel := func() {}
|
||
|
|
if queue.eventTimeout > 0 {
|
||
|
|
ctx, cancel = context.WithTimeout(context.Background(), queue.eventTimeout)
|
||
|
|
}
|
||
|
|
defer cancel()
|
||
|
|
return queue.observer.PublishDownload(ctx, event)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) publishBestEffort(event Event) {
|
||
|
|
if err := queue.publish(event); err != nil {
|
||
|
|
queue.onObserverError(event, err)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func sameTaskIdentity(left, right Task) bool {
|
||
|
|
return left.RequestID == right.RequestID &&
|
||
|
|
left.AppID == right.AppID &&
|
||
|
|
left.URL == right.URL &&
|
||
|
|
left.TotalKnown == right.TotalKnown &&
|
||
|
|
left.Total == right.Total
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) cleanupCanceled(task Task) error {
|
||
|
|
paths, err := DeriveTaskPaths(queue.root, task.RequestID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
for _, filePath := range []string{paths.Part, paths.Completed} {
|
||
|
|
if err := removeExactRegularFile(filePath); err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return queue.store.Delete(task.RequestID)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) reconcile(task Task) (Task, error) {
|
||
|
|
if err := task.Validate(); err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
paths, err := ensureTransferLayout(queue.root, task.RequestID)
|
||
|
|
if err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
partSize, partExists, err := regularFileSize(paths.Part)
|
||
|
|
if err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
finalSize, finalExists, err := regularFileSize(paths.Completed)
|
||
|
|
if err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
if partExists && finalExists {
|
||
|
|
return Task{}, errors.New("part and completed files both exist")
|
||
|
|
}
|
||
|
|
if !task.TotalKnown {
|
||
|
|
if partSize > queue.maxUnknown {
|
||
|
|
return Task{}, errors.New("part exceeds unknown-size limit")
|
||
|
|
}
|
||
|
|
if finalSize > queue.maxUnknown {
|
||
|
|
return Task{}, errors.New("completed file exceeds unknown-size limit")
|
||
|
|
}
|
||
|
|
}
|
||
|
|
changed := false
|
||
|
|
if task.Status == StatusCompleted {
|
||
|
|
if !finalExists || partExists {
|
||
|
|
return Task{}, errors.New("completed metadata has no exclusive final file")
|
||
|
|
}
|
||
|
|
if task.TotalKnown && finalSize != task.Total {
|
||
|
|
return Task{}, errors.New("completed file size mismatch")
|
||
|
|
}
|
||
|
|
if task.Done != finalSize {
|
||
|
|
task.Done = finalSize
|
||
|
|
changed = true
|
||
|
|
}
|
||
|
|
} else if finalExists {
|
||
|
|
if task.TotalKnown && finalSize != task.Total {
|
||
|
|
return Task{}, errors.New("unexpected final file size")
|
||
|
|
}
|
||
|
|
task.Status = StatusCompleted
|
||
|
|
task.Done = finalSize
|
||
|
|
task.ErrorCode = ""
|
||
|
|
changed = true
|
||
|
|
} else {
|
||
|
|
if task.TotalKnown && partSize > task.Total {
|
||
|
|
return Task{}, errors.New("part exceeds expected size")
|
||
|
|
}
|
||
|
|
if task.Done != partSize {
|
||
|
|
task.Done = partSize
|
||
|
|
changed = true
|
||
|
|
}
|
||
|
|
if task.TotalKnown && partExists && partSize == task.Total {
|
||
|
|
if err := syncRegularFile(paths.Part); err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
if err := activateCompleted(paths, nil); err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
task.Status = StatusCompleted
|
||
|
|
task.ErrorCode = ""
|
||
|
|
changed = true
|
||
|
|
} else if task.Status == StatusDownloading {
|
||
|
|
task.Status = StatusQueued
|
||
|
|
changed = true
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if changed {
|
||
|
|
if err := queue.store.Write(task); err != nil {
|
||
|
|
return Task{}, err
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return task, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) finishSuccess(
|
||
|
|
runtime *runtimeTask,
|
||
|
|
attempt uint64,
|
||
|
|
completedPath string,
|
||
|
|
) {
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.attemptCurrent(runtime, attempt) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return
|
||
|
|
}
|
||
|
|
if runtime.stop == stopCancel {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
queue.finishStopped(runtime, attempt)
|
||
|
|
return
|
||
|
|
}
|
||
|
|
runtime.task.Status = StatusCompleted
|
||
|
|
runtime.task.ErrorCode = ""
|
||
|
|
snapshot := runtime.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
err := queue.store.Write(snapshot)
|
||
|
|
if err == nil {
|
||
|
|
queue.mu.Lock()
|
||
|
|
if queue.attemptCurrent(runtime, attempt) {
|
||
|
|
delete(queue.appTasks, snapshot.AppID)
|
||
|
|
}
|
||
|
|
queue.mu.Unlock()
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventCompleted,
|
||
|
|
RequestID: snapshot.RequestID,
|
||
|
|
AppID: snapshot.AppID,
|
||
|
|
Attempt: snapshot.Attempt,
|
||
|
|
Done: snapshot.Done,
|
||
|
|
TotalKnown: snapshot.TotalKnown,
|
||
|
|
Total: snapshot.Total,
|
||
|
|
CompletedPath: completedPath,
|
||
|
|
})
|
||
|
|
} else {
|
||
|
|
log.Printf(
|
||
|
|
"downloader completed metadata persistence failed: request_id=%s: %v",
|
||
|
|
snapshot.RequestID,
|
||
|
|
err,
|
||
|
|
)
|
||
|
|
}
|
||
|
|
queue.completeWorker(runtime, attempt, err)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) finishFailure(runtime *runtimeTask, attempt uint64, failure error) {
|
||
|
|
task, valid := queue.taskForAttempt(runtime, attempt)
|
||
|
|
if !valid {
|
||
|
|
return
|
||
|
|
}
|
||
|
|
paths, _ := DeriveTaskPaths(queue.root, task.RequestID)
|
||
|
|
done, _, statErr := regularFileSize(paths.Part)
|
||
|
|
if statErr != nil {
|
||
|
|
failure = statErr
|
||
|
|
}
|
||
|
|
errorCode := transferErrorCode(failure)
|
||
|
|
queue.mu.Lock()
|
||
|
|
if !queue.attemptCurrent(runtime, attempt) {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return
|
||
|
|
}
|
||
|
|
runtime.task.Status = StatusFailed
|
||
|
|
runtime.task.Done = done
|
||
|
|
runtime.task.ErrorCode = errorCode
|
||
|
|
snapshot := runtime.task
|
||
|
|
queue.mu.Unlock()
|
||
|
|
if err := queue.store.Write(snapshot); err != nil {
|
||
|
|
failure = err
|
||
|
|
}
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventFailed,
|
||
|
|
RequestID: snapshot.RequestID,
|
||
|
|
AppID: snapshot.AppID,
|
||
|
|
Attempt: snapshot.Attempt,
|
||
|
|
Done: snapshot.Done,
|
||
|
|
ErrorCode: errorCode,
|
||
|
|
})
|
||
|
|
queue.completeWorker(runtime, attempt, failure)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) finishStopped(runtime *runtimeTask, attempt uint64) {
|
||
|
|
task, valid := queue.taskForAttempt(runtime, attempt)
|
||
|
|
if !valid {
|
||
|
|
return
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
reason := runtime.stop
|
||
|
|
queue.mu.Unlock()
|
||
|
|
|
||
|
|
var result error
|
||
|
|
switch reason {
|
||
|
|
case stopPause, stopClose:
|
||
|
|
paths, _ := DeriveTaskPaths(queue.root, task.RequestID)
|
||
|
|
done, _, err := regularFileSize(paths.Part)
|
||
|
|
if err != nil {
|
||
|
|
result = err
|
||
|
|
break
|
||
|
|
}
|
||
|
|
queue.mu.Lock()
|
||
|
|
if queue.attemptCurrent(runtime, attempt) {
|
||
|
|
runtime.task.Done = done
|
||
|
|
runtime.task.ErrorCode = ""
|
||
|
|
if reason == stopPause {
|
||
|
|
runtime.task.Status = StatusPaused
|
||
|
|
} else {
|
||
|
|
runtime.task.Status = StatusQueued
|
||
|
|
}
|
||
|
|
task = runtime.task
|
||
|
|
}
|
||
|
|
queue.mu.Unlock()
|
||
|
|
result = queue.store.Write(task)
|
||
|
|
if result == nil && reason == stopPause {
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventPaused,
|
||
|
|
RequestID: task.RequestID,
|
||
|
|
AppID: task.AppID,
|
||
|
|
Attempt: task.Attempt,
|
||
|
|
Done: task.Done,
|
||
|
|
})
|
||
|
|
}
|
||
|
|
case stopCancel:
|
||
|
|
result = queue.cleanupCanceled(task)
|
||
|
|
if result == nil {
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventCanceled,
|
||
|
|
RequestID: task.RequestID,
|
||
|
|
AppID: task.AppID,
|
||
|
|
Attempt: task.Attempt,
|
||
|
|
})
|
||
|
|
queue.mu.Lock()
|
||
|
|
delete(queue.tasks, task.RequestID)
|
||
|
|
delete(queue.appTasks, task.AppID)
|
||
|
|
queue.mu.Unlock()
|
||
|
|
} else {
|
||
|
|
queue.mu.Lock()
|
||
|
|
if queue.attemptCurrent(runtime, attempt) {
|
||
|
|
runtime.task.Status = StatusFailed
|
||
|
|
runtime.task.ErrorCode = transferErrorCode(result)
|
||
|
|
task = runtime.task
|
||
|
|
}
|
||
|
|
queue.mu.Unlock()
|
||
|
|
_ = queue.store.Write(task)
|
||
|
|
queue.publishBestEffort(Event{
|
||
|
|
Type: EventFailed,
|
||
|
|
RequestID: task.RequestID,
|
||
|
|
AppID: task.AppID,
|
||
|
|
Attempt: task.Attempt,
|
||
|
|
Done: task.Done,
|
||
|
|
ErrorCode: task.ErrorCode,
|
||
|
|
})
|
||
|
|
}
|
||
|
|
default:
|
||
|
|
result = context.Canceled
|
||
|
|
}
|
||
|
|
queue.completeWorker(runtime, attempt, result)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (queue *Queue) completeWorker(
|
||
|
|
runtime *runtimeTask,
|
||
|
|
attempt uint64,
|
||
|
|
result error,
|
||
|
|
) {
|
||
|
|
queue.mu.Lock()
|
||
|
|
if runtime.task.Attempt != attempt || !runtime.running {
|
||
|
|
queue.mu.Unlock()
|
||
|
|
return
|
||
|
|
}
|
||
|
|
runtime.running = false
|
||
|
|
runtime.cancel = nil
|
||
|
|
runtime.stop = stopNone
|
||
|
|
runtime.lastResult = result
|
||
|
|
done := runtime.done
|
||
|
|
runtime.done = nil
|
||
|
|
queue.active--
|
||
|
|
queue.mu.Unlock()
|
||
|
|
close(done)
|
||
|
|
queue.schedule()
|
||
|
|
}
|
||
|
|
|
||
|
|
func transferErrorCode(err error) string {
|
||
|
|
switch {
|
||
|
|
case errors.Is(err, ErrRangeMismatch):
|
||
|
|
return "range_mismatch"
|
||
|
|
case errors.Is(err, ErrRangeEntityChanged):
|
||
|
|
return "range_entity_changed"
|
||
|
|
case errors.Is(err, ErrResponseEncoding):
|
||
|
|
return "response_encoding"
|
||
|
|
case errors.Is(err, ErrHTTPStatus):
|
||
|
|
return "http_status"
|
||
|
|
case errors.Is(err, ErrTransferTooLarge):
|
||
|
|
return "transfer_too_large"
|
||
|
|
case errors.Is(err, ErrTransferIncomplete):
|
||
|
|
return "transfer_incomplete"
|
||
|
|
case errors.Is(err, ErrTaskCorrupt):
|
||
|
|
return "storage_corrupt"
|
||
|
|
case errors.Is(err, os.ErrPermission):
|
||
|
|
return "storage_permission"
|
||
|
|
default:
|
||
|
|
return "download_failed"
|
||
|
|
}
|
||
|
|
}
|