createTransferTask 创建存储间传输任务
(taskID string, createdAt time.Time, req *CreateTaskRequest)
| 321 | |
| 322 | // createTransferTask 创建存储间传输任务 |
| 323 | func (f *TaskFactory) createTransferTask(taskID string, createdAt time.Time, req *CreateTaskRequest) (*CreateTaskResponse, error) { |
| 324 | var params TransferParams |
| 325 | if err := json.Unmarshal(req.Params, ¶ms); err != nil { |
| 326 | return nil, fmt.Errorf("invalid params: %w", err) |
| 327 | } |
| 328 | |
| 329 | // 验证源存储和目标存储 |
| 330 | sourceStor, ok := storage.Storages[params.SourceStorage] |
| 331 | if !ok { |
| 332 | return nil, fmt.Errorf("source storage not found: %s", params.SourceStorage) |
| 333 | } |
| 334 | |
| 335 | targetStor, ok := storage.Storages[params.TargetStorage] |
| 336 | if !ok { |
| 337 | return nil, fmt.Errorf("target storage not found: %s", params.TargetStorage) |
| 338 | } |
| 339 | |
| 340 | // 检查源存储是否可读 |
| 341 | sourceReadable, ok := sourceStor.(storage.StorageReadable) |
| 342 | if !ok { |
| 343 | return nil, fmt.Errorf("source storage does not support reading: %s", params.SourceStorage) |
| 344 | } |
| 345 | |
| 346 | // 检查源存储是否可列 |
| 347 | sourceListable, ok := sourceStor.(storage.StorageListable) |
| 348 | if !ok { |
| 349 | return nil, fmt.Errorf("source storage does not support listing: %s", params.SourceStorage) |
| 350 | } |
| 351 | |
| 352 | // 列出源文件 |
| 353 | files, err := sourceListable.ListFiles(f.ctx, params.SourcePath) |
| 354 | if err != nil { |
| 355 | return nil, fmt.Errorf("failed to list source files: %w", err) |
| 356 | } |
| 357 | |
| 358 | if len(files) == 0 { |
| 359 | return nil, fmt.Errorf("no files found at source path: %s", params.SourcePath) |
| 360 | } |
| 361 | |
| 362 | // 创建传输元素 |
| 363 | elems := make([]transfer.TaskElement, 0, len(files)) |
| 364 | for _, file := range files { |
| 365 | elem := transfer.NewTaskElement(sourceReadable, file, targetStor, params.TargetPath) |
| 366 | elems = append(elems, *elem) |
| 367 | } |
| 368 | |
| 369 | task := transfer.NewTransferTask(taskID, f.ctx, elems, nil, true) |
| 370 | |
| 371 | err = f.registerAndEnqueueTask(task, tasktype.TaskTypeTransfer, params.TargetStorage, params.TargetPath, req.Webhook) |
| 372 | if err != nil { |
| 373 | return nil, err |
| 374 | } |
| 375 | |
| 376 | return &CreateTaskResponse{ |
| 377 | TaskID: taskID, |
| 378 | Type: tasktype.TaskTypeTransfer, |
| 379 | Status: TaskStatusQueued, |
| 380 | CreatedAt: createdAt, |
no test coverage detected