From 339299e08ea4fa9cdf55f5429e5712c8f91081a5 Mon Sep 17 00:00:00 2001 From: StarFleetCPTN Date: Sat, 8 Mar 2025 15:55:41 -0800 Subject: [PATCH] feat: Improve S3-compatible storage path handling in rclone operations - Update executeJob method to handle source and destination paths for S3, MinIO, and B2 - Add bucket support for size, list, transfer, and archive operations - Ensure correct path construction for different storage types - Remove duplicate error logging in archive operation --- internal/scheduler/scheduler.go | 64 +++++++++++++++++++++++++++++---- 1 file changed, 57 insertions(+), 7 deletions(-) diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index 5db0f9d..2f11262 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -161,8 +161,21 @@ func (s *Scheduler) executeJob(jobID uint) { "--config", configPath, "size", "--include", job.Config.FilePattern, - fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourcePath), } + + // Add source path with bucket for S3-compatible storage + var sourceSizePath string + if job.Config.SourceType == "s3" || job.Config.SourceType == "minio" || job.Config.SourceType == "b2" { + sourceSizePath = fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourceBucket) + if job.Config.SourcePath != "" && job.Config.SourcePath != "/" { + sourceSizePath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourceBucket, job.Config.SourcePath) + } + } else { + sourceSizePath = fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourcePath) + } + + sizeArgs = append(sizeArgs, sourceSizePath) + // Get the rclone path from the environment variable or use the default path rclonePath := os.Getenv("RCLONE_PATH") if rclonePath == "" { @@ -202,9 +215,21 @@ func (s *Scheduler) executeJob(jobID uint) { "--config", configPath, "lsf", "--include", job.Config.FilePattern, - fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourcePath), } + // Add source path with bucket for S3-compatible storage + var sourceLsPath string + if job.Config.SourceType == "s3" || job.Config.SourceType == "minio" || job.Config.SourceType == "b2" { + sourceLsPath = fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourceBucket) + if job.Config.SourcePath != "" && job.Config.SourcePath != "/" { + sourceLsPath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourceBucket, job.Config.SourcePath) + } + } else { + sourceLsPath = fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourcePath) + } + + listArgs = append(listArgs, sourceLsPath) + fmt.Printf("Listing files for job %d: rclone %s\n", jobID, strings.Join(listArgs, " ")) // Get the rclone path from the environment variable or use the default path rclonePath := os.Getenv("RCLONE_PATH") @@ -247,8 +272,26 @@ func (s *Scheduler) executeJob(jobID uint) { } // Source and destination paths - sourcePath := fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourcePath, file) - destPath := fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestinationPath, file) + var sourcePath, destPath string + + // For S3, MinIO, and B2, include the bucket in the path + if job.Config.SourceType == "s3" || job.Config.SourceType == "minio" || job.Config.SourceType == "b2" { + sourcePath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourceBucket, file) + if job.Config.SourcePath != "" && job.Config.SourcePath != "/" { + sourcePath = fmt.Sprintf("source_%d:%s/%s/%s", job.Config.ID, job.Config.SourceBucket, job.Config.SourcePath, file) + } + } else { + sourcePath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourcePath, file) + } + + if job.Config.DestinationType == "s3" || job.Config.DestinationType == "minio" || job.Config.DestinationType == "b2" { + destPath = fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestBucket, file) + if job.Config.DestinationPath != "" && job.Config.DestinationPath != "/" { + destPath = fmt.Sprintf("dest_%d:%s/%s/%s", job.Config.ID, job.Config.DestBucket, job.Config.DestinationPath, file) + } + } else { + destPath = fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestinationPath, file) + } // Add output filename pattern if specified if job.Config.OutputPattern != "" { @@ -299,9 +342,18 @@ func (s *Scheduler) executeJob(jobID uint) { "--config", configPath, "copyto", sourcePath, - fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.ArchivePath, file), } + // Construct archive path with bucket if needed + var archiveDest string + if job.Config.SourceType == "s3" || job.Config.SourceType == "minio" || job.Config.SourceType == "b2" { + archiveDest = fmt.Sprintf("source_%d:%s/%s/%s", job.Config.ID, job.Config.SourceBucket, job.Config.ArchivePath, file) + } else { + archiveDest = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.ArchivePath, file) + } + + archiveArgs = append(archiveArgs, archiveDest) + fmt.Printf("Executing rclone archive command for job %d, file %s: rclone %s\n", jobID, file, strings.Join(archiveArgs, " ")) // Get the rclone path from the environment variable or use the default path @@ -320,8 +372,6 @@ func (s *Scheduler) executeJob(jobID uint) { fmt.Printf("Warning: Error archiving file %s for job %d: %v\n", file, jobID, archiveErr) transferErrors = append(transferErrors, fmt.Sprintf("Archive error for file %s: %v", file, archiveErr)) - transferErrors = append(transferErrors, - fmt.Sprintf("Archive error for file %s: %v", file, archiveErr)) } } if job.Config.DeleteAfterTransfer {