Files
soft_quay/core/downloader/task.go
T

240 lines
7.5 KiB
Go

package downloader
import (
"errors"
"fmt"
"net/http"
"net/url"
"regexp"
"strings"
"time"
)
const (
TaskSchemaVersion = 1
MaxRequestIDBytes = 128
MaxAppIDBytes = 128
MaxURLBytes = 4096
MaxErrorCodeBytes = 64
MaxValidatorBytes = 512
)
var (
ErrInvalidTask = errors.New("download task is invalid")
ErrTaskConflict = errors.New("download task conflicts with an existing task")
ErrTaskNotFound = errors.New("download task not found")
ErrInvalidCommand = errors.New("download command is invalid for task state")
ErrTaskBusy = errors.New("download task is busy")
ErrQueueClosed = errors.New("download queue is closed")
ErrTaskCorrupt = errors.New("download task state is corrupt")
requestIDPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,127}$`)
appIDPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,127}$`)
errorCodePattern = regexp.MustCompile(`^[a-z0-9][a-z0-9_]{0,63}$`)
strongETagPattern = regexp.MustCompile(`^"[\x21\x23-\x7E\x80-\xFF]*"$`)
)
// TaskStatus is the durable lifecycle of one byte transfer.
type TaskStatus string
const (
StatusQueued TaskStatus = "queued"
StatusDownloading TaskStatus = "downloading"
StatusPaused TaskStatus = "paused"
StatusFailed TaskStatus = "failed"
StatusCompleted TaskStatus = "completed"
)
// EntityValidator binds a resumed range to the same remote representation.
// ETag must be strong. LastModified must be an HTTP-date.
type EntityValidator struct {
ETag string `json:"etag,omitempty"`
LastModified string `json:"last_modified,omitempty"`
}
// Task is the persisted download-task.json v1 protocol.
//
// Local paths are deliberately absent. They are derived from RequestID below
// the configured downloads root so remote metadata cannot choose filesystem
// destinations.
type Task struct {
SchemaVersion int `json:"schema_version"`
RequestID string `json:"request_id"`
AppID string `json:"app_id"`
URL string `json:"url"`
Status TaskStatus `json:"status"`
Attempt uint64 `json:"attempt"`
Done int64 `json:"done"`
TotalKnown bool `json:"total_known"`
Total int64 `json:"total"`
Validator EntityValidator `json:"validator"`
ErrorCode string `json:"error_code"`
CreatedAt string `json:"created_at"`
}
// Terminal reports whether the task no longer occupies the per-app active
// task slot. Completed transfers remain queryable until a later workflow
// consumes or removes them.
func (task Task) Terminal() bool {
return task.Status == StatusCompleted
}
// Validate enforces the durable protocol independently of JSON Schema.
func (task Task) Validate() error {
if task.SchemaVersion != TaskSchemaVersion {
return fmt.Errorf("%w: schema_version=%d", ErrInvalidTask, task.SchemaVersion)
}
if !ValidRequestID(task.RequestID) {
return fmt.Errorf("%w: invalid request_id", ErrInvalidTask)
}
if !appIDPattern.MatchString(task.AppID) || len(task.AppID) > MaxAppIDBytes {
return fmt.Errorf("%w: invalid app_id", ErrInvalidTask)
}
if err := ValidateHTTPSURL(task.URL); err != nil {
return fmt.Errorf("%w: url: %v", ErrInvalidTask, err)
}
switch task.Status {
case StatusQueued, StatusDownloading, StatusPaused, StatusFailed, StatusCompleted:
default:
return fmt.Errorf("%w: invalid status %q", ErrInvalidTask, task.Status)
}
if task.Attempt > uint64(^uint64(0)>>1) {
return fmt.Errorf("%w: attempt is too large", ErrInvalidTask)
}
if task.Done < 0 {
return fmt.Errorf("%w: done=%d", ErrInvalidTask, task.Done)
}
if task.TotalKnown {
if task.Total <= 0 || task.Done > task.Total {
return fmt.Errorf(
"%w: invalid known total done=%d total=%d",
ErrInvalidTask,
task.Done,
task.Total,
)
}
if task.Status == StatusCompleted && task.Done != task.Total {
return fmt.Errorf(
"%w: completed done=%d total=%d",
ErrInvalidTask,
task.Done,
task.Total,
)
}
} else if task.Total != 0 {
return fmt.Errorf("%w: unknown total must be zero", ErrInvalidTask)
}
if err := task.Validator.Validate(); err != nil {
return fmt.Errorf("%w: validator: %v", ErrInvalidTask, err)
}
if task.Status == StatusFailed {
if !errorCodePattern.MatchString(task.ErrorCode) ||
len(task.ErrorCode) > MaxErrorCodeBytes {
return fmt.Errorf("%w: failed task needs a stable error_code", ErrInvalidTask)
}
} else if task.ErrorCode != "" {
return fmt.Errorf("%w: error_code is only valid for failed tasks", ErrInvalidTask)
}
createdAt, err := time.Parse(time.RFC3339Nano, task.CreatedAt)
if err != nil {
return fmt.Errorf("%w: created_at: %v", ErrInvalidTask, err)
}
_, offset := createdAt.Zone()
if offset != 0 {
return fmt.Errorf("%w: created_at must be UTC", ErrInvalidTask)
}
return nil
}
// Validate validates an entity validator.
func (validator EntityValidator) Validate() error {
if validator.ETag != "" && validator.LastModified != "" {
return errors.New("etag and last_modified are mutually exclusive")
}
if validator.ETag != "" {
if len(validator.ETag) > MaxValidatorBytes ||
strings.HasPrefix(strings.ToUpper(validator.ETag), "W/") ||
!strongETagPattern.MatchString(validator.ETag) {
return errors.New("etag is not a strong ETag")
}
}
if validator.LastModified != "" {
if len(validator.LastModified) > MaxValidatorBytes {
return errors.New("last_modified is too long")
}
if _, err := http.ParseTime(validator.LastModified); err != nil {
return errors.New("last_modified is not an HTTP-date")
}
}
return nil
}
// Empty reports whether no reliable resume validator is available.
func (validator EntityValidator) Empty() bool {
return validator.ETag == "" && validator.LastModified == ""
}
// HeaderValue returns the If-Range value.
func (validator EntityValidator) HeaderValue() string {
if validator.ETag != "" {
return validator.ETag
}
return validator.LastModified
}
// ValidRequestID reports whether a request ID can safely derive local names.
func ValidRequestID(requestID string) bool {
return len(requestID) <= MaxRequestIDBytes && requestIDPattern.MatchString(requestID)
}
// ValidateHTTPSURL applies the same transport restrictions as the Catalog.
func ValidateHTTPSURL(value string) error {
if value == "" || len(value) > MaxURLBytes {
return errors.New("URL length is invalid")
}
parsed, err := url.Parse(value)
if err != nil {
return errors.New("URL is malformed")
}
if parsed.Scheme != "https" ||
parsed.Host == "" ||
parsed.User != nil ||
parsed.Fragment != "" ||
!parsed.IsAbs() {
return errors.New("URL must be absolute HTTPS without user info or fragment")
}
return nil
}
func validatorsMatch(expected, actual EntityValidator) bool {
if expected.ETag != "" {
if actual.Empty() {
return true
}
return expected.ETag == actual.ETag
}
if expected.LastModified != "" {
if actual.Empty() {
return true
}
return expected.LastModified == actual.LastModified
}
return true
}
func responseValidator(header http.Header) EntityValidator {
etag := strings.TrimSpace(header.Get("ETag"))
if etag != "" &&
!strings.HasPrefix(strings.ToUpper(etag), "W/") &&
strongETagPattern.MatchString(etag) &&
len(etag) <= MaxValidatorBytes {
return EntityValidator{ETag: etag}
}
lastModified := strings.TrimSpace(header.Get("Last-Modified"))
if lastModified != "" && len(lastModified) <= MaxValidatorBytes {
if _, err := http.ParseTime(lastModified); err == nil {
return EntityValidator{LastModified: lastModified}
}
}
return EntityValidator{}
}