summaryrefslogtreecommitdiff
path: root/main.go
diff refs
from: back
to: back
| flip
diff options
context:
space:
mode:
Diffstat (limited to 'main.go')
-rw-r--r--main.go74
1 files changed, 24 insertions, 50 deletions
diff --git a/main.go b/main.go
index f3c1196..521960e 100644
--- a/main.go
+++ b/main.go
@@ -128,48 +128,36 @@ func main() {
}
}
- select {
- case <-sigTermChan:
- generalLogger.Info("exiting due to termination signal")
- os.Exit(130)
- default:
- }
-
processSemaphore := make(chan struct{}, cfg.MaxProcessThreads)
queueSemaphore := make(chan struct{}, cfg.MaxPreProcessThreads)
var wg sync.WaitGroup
wg.Add(fileCount)
- processTerminating := false
-
for i, file := range files {
queueSemaphore <- struct{}{}
+ go func(index int, inputName string) {
+ fileLogger := generalLogger.With(slog.String("input_path", inputName), slog.Int("file_index", index))
- select {
- case <-sigTermChan:
- generalLogger.Info("exiting due to termination signal")
- processTerminating = true
- default:
- }
+ threadSigTermChannel := make(chan os.Signal, 1)
+ signal.Notify(threadSigTermChannel, os.Interrupt, syscall.SIGTERM)
- go func(index int, inputName string, earlyTerminate bool) {
defer func() { <-queueSemaphore; wg.Done() }()
- fileLogger := generalLogger.With(slog.String("input_path", inputName), slog.Int("file_index", index))
-
- if earlyTerminate {
- fileLogger.Info("skip processing file (process is terminating)")
+ select {
+ case <-threadSigTermChannel:
+ fileLogger.Info("exiting due to termination signal")
return
+ default:
}
var inputMetadata *input.MetadataStruct
id := inputClient.ID(file)
- cacheMapMutex.Lock()
if _, ok := cacheMap[id]; !ok {
+ cacheMapMutex.Lock()
cacheMap[id] = make(map[uint32]struct{})
+ cacheMapMutex.Unlock()
}
- cacheMapMutex.Unlock()
convertersToLaunch := []int{}
for j, conv := range converters {
@@ -197,31 +185,21 @@ func main() {
}
}
- switch cfg.Converters[j].Output.RewriteOn {
- case "Never":
- if !conv.IsMissing(outputName) {
- convLogger.Info("skip already existing file")
+ if !cfg.ForceRewrite && !conv.IsMissing(outputName) {
+ outputMetadata, err := conv.ReadMetadata(outputName)
+ if err != nil {
+ convLogger.Warn("fail to read metadata of (supposedly existing) output file", slog.String("error", err.Error()))
continue
}
- case "UnequalHashInCache":
- if !conv.IsMissing(outputName) {
- outputMetadata, err := conv.ReadMetadata(outputName)
- if err != nil {
- convLogger.Warn("fail to read metadata of (supposedly existing) output file", slog.String("error", err.Error()))
- continue
- }
- originalInputHash = outputMetadata.HashOriginal
- if inputMetadata.Hash == originalInputHash {
- convLogger.Info("skip already processed file (based on equal hash)", slog.String("input_hash", inputMetadata.Hash))
- cacheMapMutex.Lock()
- cacheMap[id][converterHashes[j]] = struct{}{}
- cacheMapMutex.Unlock()
- continue
- }
+ originalInputHash = outputMetadata.HashOriginal
+ if inputMetadata.Hash == originalInputHash {
+ convLogger.Info("skip already processed file (based on equal hash)", slog.String("input_hash", inputMetadata.Hash))
+ cacheMapMutex.Lock()
+ cacheMap[id][converterHashes[j]] = struct{}{}
+ cacheMapMutex.Unlock()
+ continue
}
- case "Always":
}
-
convertersToLaunch = append(convertersToLaunch, j)
}
@@ -239,14 +217,10 @@ func main() {
return
}
- fileContent, err := io.ReadAll(reader)
- if err != nil {
+ fileContent := make([]byte, inputMetadata.Size)
+ if _, err = reader.Read(fileContent); err != nil {
fileLogger.Warn("fail to read content of input file", slog.String("error", err.Error()))
return
- } else if int64(len(fileContent)) != inputMetadata.Size {
- fileLogger.Warn("the downloaded file seems to be missing some of its content",
- slog.Int("actual_size_bytes", len(fileContent)),
- slog.Int64("expected_size_bytes", inputMetadata.Size))
}
reader.Close()
@@ -265,7 +239,7 @@ func main() {
cacheMap[id][converterHashes[convIndex]] = struct{}{}
cacheMapMutex.Unlock()
}
- }(i, file, processTerminating)
+ }(i, file)
}
wg.Wait()