Files
soft_quay/core/downloader/queue.go
T

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"
}
}