| 71 | } |
| 72 | |
| 73 | func (t *Task) processElement(ctx context.Context, elem TaskElement) error { |
| 74 | logger := log.FromContext(ctx).WithPrefix(fmt.Sprintf("file[%s]", elem.FileInfo.Name)) |
| 75 | |
| 76 | // Check whether the source storage supports reading |
| 77 | readableStorage, ok := elem.SourceStorage.(storage.StorageReadable) |
| 78 | if !ok { |
| 79 | return fmt.Errorf("source storage %s does not support reading", elem.SourceStorage.Name()) |
| 80 | } |
| 81 | |
| 82 | logger.Info("Opening file from source storage") |
| 83 | reader, size, err := readableStorage.OpenFile(ctx, elem.SourcePath) |
| 84 | if err != nil { |
| 85 | return fmt.Errorf("failed to open file: %w", err) |
| 86 | } |
| 87 | defer reader.Close() |
| 88 | |
| 89 | // Build target storage path: /target_path/filename |
| 90 | storagePath := path.Join(elem.TargetPath, elem.FileInfo.Name) |
| 91 | |
| 92 | // Inject file size into context |
| 93 | ctx = context.WithValue(ctx, ctxkey.ContentLength, size) |
| 94 | |
| 95 | if config.C().Stream { |
| 96 | if err := elem.TargetStorage.Save(ctx, reader, storagePath); err != nil { |
| 97 | return fmt.Errorf("failed to upload file to storage: %w", err) |
| 98 | } |
| 99 | } else { |
| 100 | logger.Info("Downloading to temporary file for ReadSeeker support") |
| 101 | tempFile, err := t.downloadToTemp(reader, elem.FileInfo.Name) |
| 102 | if err != nil { |
| 103 | return fmt.Errorf("failed to download to temp: %w", err) |
| 104 | } |
| 105 | defer os.Remove(tempFile.Name()) |
| 106 | defer tempFile.Close() |
| 107 | |
| 108 | if _, err := tempFile.Seek(0, io.SeekStart); err != nil { |
| 109 | return fmt.Errorf("failed to seek temp file: %w", err) |
| 110 | } |
| 111 | |
| 112 | logger.Infof("Uploading file to storage (size: %d bytes)", size) |
| 113 | if err := elem.TargetStorage.Save(ctx, tempFile, storagePath); err != nil { |
| 114 | return fmt.Errorf("failed to upload file to storage: %w", err) |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | t.uploaded.Add(size) |
| 119 | t.Progress.OnProgress(ctx, t) |
| 120 | taskevent.Emit(ctx, taskevent.Event{ |
| 121 | TaskID: t.ID, |
| 122 | Phase: taskevent.PhaseProgress, |
| 123 | TotalBytes: t.totalSize, |
| 124 | DownloadedBytes: t.uploaded.Load(), |
| 125 | }) |
| 126 | |
| 127 | logger.Info("File uploaded successfully") |
| 128 | return nil |
| 129 | } |
| 130 | |