mirror of
https://github.com/StarFleetCPTN/GoMFT.git
synced 2026-09-08 15:41:20 +02:00
feat: Enhance file listing and metadata extraction for rclone transfers
- Replace separate rclone size and list commands with single lsjson operation - Improve file metadata extraction from JSON output - Add support for hash and file size retrieval from remote sources - Optimize file processing logic with more robust metadata handling - Enhance error handling and logging for file listing and transfer processes
This commit is contained in:
@@ -162,7 +162,9 @@ templ JobRunDetails(ctx context.Context, data JobRunDetailsData) {
|
||||
<dt class="text-sm font-medium text-secondary-500 dark:text-secondary-400">Source Path</dt>
|
||||
<dd class="mt-1 text-sm text-secondary-900 dark:text-secondary-100 font-mono bg-secondary-50 dark:bg-secondary-900 p-2 rounded">
|
||||
if data.Config.SourceType == "sftp" {
|
||||
{ data.Config.SourceUser } `@` { data.Config.SourceHost }:{ data.Config.SourcePath }
|
||||
{ data.Config.SourceUser }{`@`}{ data.Config.SourceHost }{`:`}{ data.Config.SourcePath }
|
||||
} else if data.Config.SourceType == "s3" || data.Config.SourceType == "minio" || data.Config.SourceType == "b2" {
|
||||
{ data.Config.SourceBucket }{`:`}{ data.Config.SourcePath }
|
||||
} else {
|
||||
{ data.Config.SourcePath }
|
||||
}
|
||||
@@ -173,7 +175,9 @@ templ JobRunDetails(ctx context.Context, data JobRunDetailsData) {
|
||||
<dt class="text-sm font-medium text-secondary-500 dark:text-secondary-400">Destination Path</dt>
|
||||
<dd class="mt-1 text-sm text-secondary-900 dark:text-secondary-100 font-mono bg-secondary-50 dark:bg-secondary-900 p-2 rounded">
|
||||
if data.Config.DestinationType == "sftp" {
|
||||
data.Config.DestUser@data.Config.DestHost:data.Config.DestinationPath
|
||||
{ data.Config.DestUser }{`@`}{ data.Config.DestHost }{`:`}{ data.Config.DestinationPath }
|
||||
} else if data.Config.DestinationType == "s3" || data.Config.DestinationType == "minio" || data.Config.DestinationType == "b2" {
|
||||
{ data.Config.DestBucket }{`:`}{ data.Config.DestinationPath }
|
||||
} else {
|
||||
{ data.Config.DestinationPath }
|
||||
}
|
||||
|
||||
+165
-203
@@ -11,7 +11,6 @@ import (
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -164,231 +163,194 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
// Get rclone config path
|
||||
configPath := s.db.GetConfigRclonePath(&job.Config)
|
||||
|
||||
// Size of transfer using rclone size
|
||||
sizeArgs := []string{
|
||||
// Use lsjson to get file list and metadata in one operation instead of separate size and ls commands
|
||||
listArgs := []string{
|
||||
"--config", configPath,
|
||||
"size",
|
||||
"--include", job.Config.FilePattern,
|
||||
"lsjson",
|
||||
"--hash",
|
||||
"--recursive",
|
||||
}
|
||||
|
||||
// Add file pattern filter if specified
|
||||
if job.Config.FilePattern != "" && job.Config.FilePattern != "*" {
|
||||
// Create a temporary filter file for complex patterns
|
||||
filterFile, err := createRcloneFilterFile(job.Config.FilePattern)
|
||||
if err != nil {
|
||||
fmt.Printf("Error creating filter file for job %d: %v\n", jobID, err)
|
||||
history.Status = "failed"
|
||||
history.ErrorMessage = fmt.Sprintf("Filter Creation Error: %v", err)
|
||||
endTime := time.Now()
|
||||
history.EndTime = &endTime
|
||||
if err := s.db.UpdateJobHistory(history); err != nil {
|
||||
fmt.Printf("Error updating job history for job %d: %v\n", jobID, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
defer os.Remove(filterFile)
|
||||
listArgs = append(listArgs, "--filter-from", filterFile)
|
||||
}
|
||||
|
||||
// Add source path with bucket for S3-compatible storage
|
||||
var sourceSizePath string
|
||||
var sourceListPath 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)
|
||||
sourceListPath = 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)
|
||||
sourceListPath = 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)
|
||||
sourceListPath = fmt.Sprintf("source_%d:%s", job.Config.ID, job.Config.SourcePath)
|
||||
}
|
||||
|
||||
sizeArgs = append(sizeArgs, sourceSizePath)
|
||||
listArgs = append(listArgs, sourceListPath)
|
||||
|
||||
// Get the rclone path from the environment variable or use the default path
|
||||
// Execute lsjson command
|
||||
fmt.Printf("Listing files with metadata for job %d: rclone %s\n", jobID, strings.Join(listArgs, " "))
|
||||
rclonePath := os.Getenv("RCLONE_PATH")
|
||||
if rclonePath == "" {
|
||||
rclonePath = "rclone"
|
||||
}
|
||||
output, err := exec.Command(rclonePath, sizeArgs...).CombinedOutput()
|
||||
fmt.Printf("Running rclone size: %s %s\nOutput: %s\n", rclonePath, strings.Join(sizeArgs, " "), output)
|
||||
if err != nil {
|
||||
fmt.Printf("Error running rclone size: %v\nOutput: %s\n", err, output)
|
||||
// Update job history with error
|
||||
listCmd := exec.Command(rclonePath, listArgs...)
|
||||
listOutput, listErr := listCmd.CombinedOutput()
|
||||
|
||||
if listErr != nil {
|
||||
fmt.Printf("Error listing files for job %d: %v\n", jobID, listErr)
|
||||
history.Status = "failed"
|
||||
history.ErrorMessage = fmt.Sprintf("Size calculation error: %v\nOutput: %s", err, string(output))
|
||||
history.EndTime = &startTime // Use start time as end time for a quick failure
|
||||
history.ErrorMessage = fmt.Sprintf("File Listing Error: %v\nOutput: %s", listErr, string(listOutput))
|
||||
endTime := time.Now()
|
||||
history.EndTime = &endTime
|
||||
if err := s.db.UpdateJobHistory(history); err != nil {
|
||||
fmt.Printf("Error updating job history for job %d: %v\n", jobID, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
// Parse rclone size output "Total objects: 1 Total size: 10 B (10 Byte)"
|
||||
outputStr := string(output)
|
||||
|
||||
totalObjects := strings.TrimSpace(strings.Split(outputStr, "\n")[0])
|
||||
totalObjects = strings.TrimSpace(strings.Split(totalObjects, ":")[1])
|
||||
totalSize := strings.TrimSpace(strings.Split(outputStr, "Total size:")[1])
|
||||
if strings.Contains(totalSize, "(") {
|
||||
totalSize = strings.TrimSpace(strings.Split(totalSize, "(")[1])
|
||||
totalSize = strings.TrimSpace(strings.Split(totalSize, " ")[0])
|
||||
} else {
|
||||
totalSize = "0"
|
||||
// Parse JSON output to get file information
|
||||
var fileEntries []map[string]interface{}
|
||||
if err := json.Unmarshal(listOutput, &fileEntries); err != nil {
|
||||
fmt.Printf("Error parsing file list JSON for job %d: %v\n", jobID, err)
|
||||
history.Status = "failed"
|
||||
history.ErrorMessage = fmt.Sprintf("JSON Parsing Error: %v", err)
|
||||
endTime := time.Now()
|
||||
history.EndTime = &endTime
|
||||
if err := s.db.UpdateJobHistory(history); err != nil {
|
||||
fmt.Printf("Error updating job history for job %d: %v\n", jobID, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
bytesTransferred, _ := strconv.ParseInt(totalSize, 10, 64)
|
||||
history.BytesTransferred = bytesTransferred
|
||||
// Calculate total size and filter out directories
|
||||
var files []map[string]interface{}
|
||||
var totalSize int64
|
||||
for _, entry := range fileEntries {
|
||||
// Skip directories
|
||||
if isDir, ok := entry["IsDir"].(bool); ok && isDir {
|
||||
continue
|
||||
}
|
||||
|
||||
if totalObjects == "0" {
|
||||
// Add to files list
|
||||
files = append(files, entry)
|
||||
|
||||
// Add to total size
|
||||
if size, ok := entry["Size"].(float64); ok {
|
||||
totalSize += int64(size)
|
||||
}
|
||||
}
|
||||
|
||||
fmt.Printf("Found %d files totaling %d bytes to transfer for job %d\n", len(files), totalSize, jobID)
|
||||
|
||||
// Update history with size information
|
||||
history.BytesTransferred = totalSize
|
||||
|
||||
if len(files) == 0 {
|
||||
fmt.Printf("No files to transfer for job %d\n", jobID)
|
||||
history.Status = "completed"
|
||||
history.ErrorMessage = ""
|
||||
history.FilesTransferred = 0
|
||||
}
|
||||
|
||||
if totalObjects != "0" {
|
||||
// First, list all files that match the pattern
|
||||
listArgs := []string{
|
||||
"--config", configPath,
|
||||
"lsf",
|
||||
"--include", job.Config.FilePattern,
|
||||
}
|
||||
|
||||
// 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")
|
||||
if rclonePath == "" {
|
||||
rclonePath = "rclone"
|
||||
}
|
||||
listCmd := exec.Command(rclonePath, listArgs...)
|
||||
listOutput, listErr := listCmd.CombinedOutput()
|
||||
|
||||
if listErr != nil {
|
||||
fmt.Printf("Error listing files for job %d: %v\n", jobID, listErr)
|
||||
history.Status = "failed"
|
||||
history.ErrorMessage = fmt.Sprintf("File Listing Error: %v\nOutput: %s", listErr, string(listOutput))
|
||||
return
|
||||
}
|
||||
|
||||
// Split the output by newlines to get individual files
|
||||
fileLines := strings.Split(strings.TrimSpace(string(listOutput)), "\n")
|
||||
|
||||
// Clean up the file list - remove empty lines and ensure uniqueness
|
||||
var files []string
|
||||
uniqueFiles := make(map[string]bool)
|
||||
|
||||
for _, line := range fileLines {
|
||||
// Skip empty lines
|
||||
trimmedLine := strings.TrimSpace(line)
|
||||
if trimmedLine == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
// Only add each file once
|
||||
if !uniqueFiles[trimmedLine] {
|
||||
uniqueFiles[trimmedLine] = true
|
||||
files = append(files, trimmedLine)
|
||||
}
|
||||
}
|
||||
|
||||
fmt.Printf("Found %d unique files to transfer for job %d\n", len(files), jobID)
|
||||
|
||||
} else {
|
||||
var transferErrors []string
|
||||
filesTransferred := 0
|
||||
|
||||
// Process each file individually
|
||||
for _, file := range files {
|
||||
// Skip files that have already been processed in this execution
|
||||
if processedFiles[file] {
|
||||
fmt.Printf("Skipping duplicate file entry: %s (already processed in this execution)\n", file)
|
||||
for _, fileEntry := range files {
|
||||
fileName, ok := fileEntry["Path"].(string)
|
||||
if !ok || fileName == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
fmt.Printf("Processing file: %s for job %d\n", file, jobID)
|
||||
// Skip files that have already been processed in this execution
|
||||
if processedFiles[fileName] {
|
||||
fmt.Printf("Skipping duplicate file entry: %s (already processed in this execution)\n", fileName)
|
||||
continue
|
||||
}
|
||||
|
||||
// Get detailed file info if this is a local source
|
||||
// Extract file metadata from the JSON entry
|
||||
var fileSize int64
|
||||
var createTime, modTime time.Time
|
||||
if size, ok := fileEntry["Size"].(float64); ok {
|
||||
fileSize = int64(size)
|
||||
}
|
||||
|
||||
// Extract modification time
|
||||
modTime := time.Now()
|
||||
if modTimeStr, ok := fileEntry["ModTime"].(string); ok {
|
||||
if parsedTime, err := time.Parse(time.RFC3339, modTimeStr); err == nil {
|
||||
modTime = parsedTime
|
||||
}
|
||||
}
|
||||
|
||||
// Create time is usually not available for remote files, so we'll use modTime
|
||||
createTime := modTime
|
||||
|
||||
// Extract hash if available
|
||||
var fileHash string
|
||||
var fileMetadataErr error
|
||||
|
||||
// Get file metadata based on source type
|
||||
if job.Config.SourceType == "local" {
|
||||
// Construct the full local path
|
||||
localFilePath := filepath.Join(job.Config.SourcePath, file)
|
||||
|
||||
// Get file info (size, creation time, modification time)
|
||||
fileSize, createTime, modTime, fileMetadataErr = getFileInfo(localFilePath)
|
||||
if fileMetadataErr != nil {
|
||||
fmt.Printf("Warning: Could not get file info for %s: %v\n", file, fileMetadataErr)
|
||||
if hashes, ok := fileEntry["Hashes"].(map[string]interface{}); ok {
|
||||
if md5, ok := hashes["md5"].(string); ok {
|
||||
fileHash = md5
|
||||
}
|
||||
}
|
||||
|
||||
// Calculate file hash (MD5)
|
||||
fileHash, fileMetadataErr = calculateFileHash(localFilePath)
|
||||
if fileMetadataErr != nil {
|
||||
fmt.Printf("Warning: Could not calculate hash for %s: %v\n", file, fileMetadataErr)
|
||||
// For local files, calculate hash if not available
|
||||
if fileHash == "" && job.Config.SourceType == "local" {
|
||||
localFilePath := filepath.Join(job.Config.SourcePath, fileName)
|
||||
calculatedHash, hashErr := calculateFileHash(localFilePath)
|
||||
if hashErr == nil {
|
||||
fileHash = calculatedHash
|
||||
}
|
||||
}
|
||||
|
||||
// Check if this file has been processed before (by hash)
|
||||
if fileHash != "" {
|
||||
processed, prevMetadata, _ := s.hasFileBeenProcessed(jobID, fileHash)
|
||||
if processed {
|
||||
fmt.Printf("File %s has been processed before (hash: %s, previous file: %s)\n",
|
||||
file, fileHash, prevMetadata.FileName)
|
||||
// Check if this file has been processed before (by hash)
|
||||
if fileHash != "" {
|
||||
processed, prevMetadata, _ := s.hasFileBeenProcessed(jobID, fileHash)
|
||||
if processed {
|
||||
fmt.Printf("File %s has been processed before (hash: %s, previous file: %s)\n",
|
||||
fileName, fileHash, prevMetadata.FileName)
|
||||
|
||||
// If configured to skip previously processed files,
|
||||
// we could add that logic here
|
||||
// For now, we'll just log it and continue
|
||||
}
|
||||
}
|
||||
|
||||
// Also check the processing history for this specific file name
|
||||
prevMetadata, histErr := s.checkFileProcessingHistory(jobID, file)
|
||||
if histErr == nil {
|
||||
fmt.Printf("File %s was previously processed on %s with status: %s\n",
|
||||
file, prevMetadata.ProcessedTime.Format(time.RFC3339), prevMetadata.Status)
|
||||
|
||||
// If the file was previously processed successfully and the hash hasn't changed,
|
||||
// we could skip processing
|
||||
// Skip previously processed files with the same hash if they were processed successfully
|
||||
if prevMetadata.Status == "processed" ||
|
||||
prevMetadata.Status == "archived" ||
|
||||
prevMetadata.Status == "deleted" ||
|
||||
prevMetadata.Status == "archived_and_deleted" {
|
||||
if fileHash != "" && fileHash == prevMetadata.FileHash {
|
||||
fmt.Printf("Skipping unchanged file %s (hash matches previous processing)\n", file)
|
||||
// Skip this file and continue to the next one
|
||||
continue
|
||||
}
|
||||
fmt.Printf("Skipping unchanged file %s (hash matches previous processing)\n", fileName)
|
||||
continue
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// For non-local sources, use rclone lsjson to get metadata
|
||||
fileSize, createTime, modTime, fileHash, fileMetadataErr = s.getRemoteFileInfo(&job.Config, file)
|
||||
if fileMetadataErr != nil {
|
||||
fmt.Printf("Warning: Could not get remote file info for %s: %v\n", file, fileMetadataErr)
|
||||
// We'll continue with placeholder values
|
||||
fileSize = 0
|
||||
createTime = time.Now()
|
||||
modTime = time.Now()
|
||||
fileHash = ""
|
||||
} else {
|
||||
// If we got a hash, check for previous processing
|
||||
if fileHash != "" {
|
||||
processed, prevMetadata, _ := s.hasFileBeenProcessed(jobID, fileHash)
|
||||
if processed {
|
||||
fmt.Printf("Remote file %s has been processed before (hash: %s, previous file: %s)\n",
|
||||
file, fileHash, prevMetadata.FileName)
|
||||
}
|
||||
|
||||
// Skip previously processed files with the same hash if they were processed successfully
|
||||
if prevMetadata.Status == "processed" ||
|
||||
prevMetadata.Status == "archived" ||
|
||||
prevMetadata.Status == "deleted" ||
|
||||
prevMetadata.Status == "archived_and_deleted" {
|
||||
fmt.Printf("Skipping unchanged remote file %s (hash matches previous processing)\n", file)
|
||||
// Skip this file and continue to the next one
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
// Also check the processing history for this specific file name
|
||||
prevMetadata, histErr := s.checkFileProcessingHistory(jobID, fileName)
|
||||
if histErr == nil {
|
||||
fmt.Printf("File %s was previously processed on %s with status: %s\n",
|
||||
fileName, prevMetadata.ProcessedTime.Format(time.RFC3339), prevMetadata.Status)
|
||||
|
||||
// Also check by filename
|
||||
prevMetadata, histErr := s.checkFileProcessingHistory(jobID, file)
|
||||
if histErr == nil {
|
||||
fmt.Printf("Remote file %s was previously processed on %s with status: %s\n",
|
||||
file, prevMetadata.ProcessedTime.Format(time.RFC3339), prevMetadata.Status)
|
||||
// If the file was previously processed successfully and the hash hasn't changed,
|
||||
// we could skip processing
|
||||
if prevMetadata.Status == "processed" ||
|
||||
prevMetadata.Status == "archived" ||
|
||||
prevMetadata.Status == "deleted" ||
|
||||
prevMetadata.Status == "archived_and_deleted" {
|
||||
if fileHash != "" && fileHash == prevMetadata.FileHash {
|
||||
fmt.Printf("Skipping unchanged file %s (hash matches previous processing)\n", fileName)
|
||||
// Skip this file and continue to the next one
|
||||
continue
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -408,29 +370,29 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
|
||||
// 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)
|
||||
sourcePath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourceBucket, fileName)
|
||||
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)
|
||||
sourcePath = fmt.Sprintf("source_%d:%s/%s/%s", job.Config.ID, job.Config.SourceBucket, job.Config.SourcePath, fileName)
|
||||
}
|
||||
} else {
|
||||
sourcePath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourcePath, file)
|
||||
sourcePath = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.SourcePath, fileName)
|
||||
}
|
||||
|
||||
var destFile string = file
|
||||
var destFile string = fileName
|
||||
|
||||
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)
|
||||
destPath = fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestBucket, fileName)
|
||||
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)
|
||||
destPath = fmt.Sprintf("dest_%d:%s/%s/%s", job.Config.ID, job.Config.DestBucket, job.Config.DestinationPath, fileName)
|
||||
}
|
||||
} else {
|
||||
destPath = fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestinationPath, file)
|
||||
destPath = fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestinationPath, fileName)
|
||||
}
|
||||
|
||||
// Add output filename pattern if specified
|
||||
if job.Config.OutputPattern != "" {
|
||||
// Process the output pattern for this specific file
|
||||
destFile = ProcessOutputPattern(job.Config.OutputPattern, file)
|
||||
destFile = ProcessOutputPattern(job.Config.OutputPattern, fileName)
|
||||
|
||||
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, destFile)
|
||||
@@ -441,7 +403,7 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
destPath = fmt.Sprintf("dest_%d:%s/%s", job.Config.ID, job.Config.DestinationPath, destFile)
|
||||
}
|
||||
|
||||
fmt.Printf("Renaming file from %s to %s for job %d\n", file, destFile, jobID)
|
||||
fmt.Printf("Renaming file from %s to %s for job %d\n", fileName, destFile, jobID)
|
||||
}
|
||||
|
||||
// Add custom flags if specified
|
||||
@@ -456,7 +418,7 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
|
||||
// Execute transfer for this file
|
||||
fmt.Printf("Executing rclone transfer command for job %d, file %s: rclone %s\n",
|
||||
jobID, file, strings.Join(transferArgs, " "))
|
||||
jobID, fileName, strings.Join(transferArgs, " "))
|
||||
// Get the rclone path from the environment variable or use the default path
|
||||
rclonePath := os.Getenv("RCLONE_PATH")
|
||||
if rclonePath == "" {
|
||||
@@ -466,7 +428,7 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
fileOutput, fileErr := cmd.CombinedOutput()
|
||||
|
||||
// Print the output
|
||||
fmt.Printf("Output for file %s: %s\n", file, string(fileOutput))
|
||||
fmt.Printf("Output for file %s: %s\n", fileName, string(fileOutput))
|
||||
|
||||
// Create file metadata record
|
||||
fileStatus := "processed"
|
||||
@@ -475,13 +437,13 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
|
||||
// Check if file was successfully transferred
|
||||
if fileErr != nil {
|
||||
fmt.Printf("Error transferring file %s for job %d: %v\n", file, jobID, fileErr)
|
||||
transferErrors = append(transferErrors, fmt.Sprintf("File %s: %v", file, fileErr))
|
||||
fmt.Printf("Error transferring file %s for job %d: %v\n", fileName, jobID, fileErr)
|
||||
transferErrors = append(transferErrors, fmt.Sprintf("File %s: %v", fileName, fileErr))
|
||||
fileStatus = "error"
|
||||
fileErrorMsg = fileErr.Error()
|
||||
} else {
|
||||
filesTransferred++
|
||||
fmt.Printf("Successfully transferred file %s for job %d\n", file, jobID)
|
||||
fmt.Printf("Successfully transferred file %s for job %d\n", fileName, jobID)
|
||||
|
||||
// Extract the actual destination path (without rclone remote prefix)
|
||||
if job.Config.DestinationType == "local" {
|
||||
@@ -501,7 +463,7 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
|
||||
// If archiving is enabled and transfer was successful, move files to archive
|
||||
if job.Config.ArchiveEnabled && job.Config.ArchivePath != "" {
|
||||
fmt.Printf("Archiving file %s for job %d\n", file, jobID)
|
||||
fmt.Printf("Archiving file %s for job %d\n", fileName, jobID)
|
||||
|
||||
// We don't need to move the file since we used moveto, but we can copy it to archive
|
||||
archiveArgs := []string{
|
||||
@@ -513,15 +475,15 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
// 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)
|
||||
archiveDest = fmt.Sprintf("source_%d:%s/%s/%s", job.Config.ID, job.Config.SourceBucket, job.Config.ArchivePath, fileName)
|
||||
} else {
|
||||
archiveDest = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.ArchivePath, file)
|
||||
archiveDest = fmt.Sprintf("source_%d:%s/%s", job.Config.ID, job.Config.ArchivePath, fileName)
|
||||
}
|
||||
|
||||
archiveArgs = append(archiveArgs, archiveDest)
|
||||
|
||||
fmt.Printf("Executing rclone archive command for job %d, file %s: rclone %s\n",
|
||||
jobID, file, strings.Join(archiveArgs, " "))
|
||||
jobID, fileName, strings.Join(archiveArgs, " "))
|
||||
// Get the rclone path from the environment variable or use the default path
|
||||
rclonePath := os.Getenv("RCLONE_PATH")
|
||||
if rclonePath == "" {
|
||||
@@ -531,31 +493,31 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
archiveOutput, archiveErr := archiveCmd.CombinedOutput()
|
||||
|
||||
// Print the output
|
||||
fmt.Printf("Output for file %s: %s\n", file, string(archiveOutput))
|
||||
fmt.Printf("Output for file %s: %s\n", fileName, string(archiveOutput))
|
||||
|
||||
// Check if file was successfully transferred
|
||||
if archiveErr != nil {
|
||||
fmt.Printf("Warning: Error archiving file %s for job %d: %v\n", file, jobID, archiveErr)
|
||||
fmt.Printf("Warning: Error archiving file %s for job %d: %v\n", fileName, jobID, archiveErr)
|
||||
transferErrors = append(transferErrors,
|
||||
fmt.Sprintf("Archive error for file %s: %v", file, archiveErr))
|
||||
fmt.Sprintf("Archive error for file %s: %v", fileName, archiveErr))
|
||||
} else {
|
||||
fileStatus = "archived"
|
||||
}
|
||||
}
|
||||
|
||||
if job.Config.DeleteAfterTransfer {
|
||||
fmt.Printf("Deleting file %s for job %d\n", file, jobID)
|
||||
fmt.Printf("Deleting file %s for job %d\n", fileName, jobID)
|
||||
deleteArgs := []string{
|
||||
"--config", configPath,
|
||||
"deletefile",
|
||||
sourcePath}
|
||||
deleteCmd := exec.Command(rclonePath, deleteArgs...)
|
||||
deleteOutput, deleteErr := deleteCmd.CombinedOutput()
|
||||
fmt.Printf("Output for file %s: %s\n", file, string(deleteOutput))
|
||||
fmt.Printf("Output for file %s: %s\n", fileName, string(deleteOutput))
|
||||
if deleteErr != nil {
|
||||
fmt.Printf("Error deleting file %s for job %d: %v\n", file, jobID, deleteErr)
|
||||
fmt.Printf("Error deleting file %s for job %d: %v\n", fileName, jobID, deleteErr)
|
||||
transferErrors = append(transferErrors,
|
||||
fmt.Sprintf("Delete error for file %s: %v", file, deleteErr))
|
||||
fmt.Sprintf("Delete error for file %s: %v", fileName, deleteErr))
|
||||
} else {
|
||||
if fileStatus == "archived" {
|
||||
fileStatus = "archived_and_deleted"
|
||||
@@ -567,12 +529,12 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
}
|
||||
|
||||
// Mark this file as processed for this execution
|
||||
processedFiles[file] = true
|
||||
processedFiles[fileName] = true
|
||||
|
||||
// Create and save file metadata
|
||||
metadata := &db.FileMetadata{
|
||||
JobID: jobID,
|
||||
FileName: file,
|
||||
FileName: fileName,
|
||||
OriginalPath: job.Config.SourcePath,
|
||||
FileSize: fileSize,
|
||||
FileHash: fileHash,
|
||||
@@ -585,9 +547,9 @@ func (s *Scheduler) executeJob(jobID uint) {
|
||||
}
|
||||
|
||||
if err := s.db.CreateFileMetadata(metadata); err != nil {
|
||||
fmt.Printf("Error creating file metadata for %s: %v\n", file, err)
|
||||
fmt.Printf("Error creating file metadata for %s: %v\n", fileName, err)
|
||||
} else {
|
||||
fmt.Printf("Created file metadata record for %s (ID: %d)\n", file, metadata.ID)
|
||||
fmt.Printf("Created file metadata record for %s (ID: %d)\n", fileName, metadata.ID)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user