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