summaryrefslogtreecommitdiff
diff refs
from: back
to: back
| flip
diff options
context:
space:
mode:
-rw-r--r--.gitignore4
-rw-r--r--config/config-b2.sample.json3
-rw-r--r--config/config-local-unix.sample.json3
-rw-r--r--config/config.go13
-rw-r--r--main.go12
5 files changed, 14 insertions, 21 deletions
diff --git a/.gitignore b/.gitignore
index 46c2db2..f44834d 100644
--- a/.gitignore
+++ b/.gitignore
@@ -28,8 +28,8 @@ go.work.sum
.env
# Editor/IDE
-.idea/
-.vscode/
+# .idea/
+# .vscode/
config*.json
!config*.sample.json
diff --git a/config/config-b2.sample.json b/config/config-b2.sample.json
index c48b1b2..d3d772e 100644
--- a/config/config-b2.sample.json
+++ b/config/config-b2.sample.json
@@ -1,7 +1,6 @@
{
"ForceRewrite": false,
- "MaxProcessThreads": 4,
- "MaxPreProcessThreads": 12,
+ "MaxConcurrentJobs": 4,
"LogLevel": "info",
"Input": {
"Storage": {
diff --git a/config/config-local-unix.sample.json b/config/config-local-unix.sample.json
index 36d8114..637a05a 100644
--- a/config/config-local-unix.sample.json
+++ b/config/config-local-unix.sample.json
@@ -1,7 +1,6 @@
{
"ForceRewrite": false,
- "MaxProcessThreads": 4,
- "MaxPreProcessThreads": 12,
+ "MaxConcurrentJobs": 4,
"LogLevel": "debug",
"Input": {
"Storage": {
diff --git a/config/config.go b/config/config.go
index 5b05ce9..d2b7072 100644
--- a/config/config.go
+++ b/config/config.go
@@ -10,13 +10,12 @@ import (
)
type Config struct {
- Input InputConfig `json:"Input" validate:"required"`
- Converter ConverterConfig `json:"Converter" validate:"required"`
- Output OutputConfig `json:"Output" validate:"required"`
- MaxProcessThreads int `json:"MaxProcessThreads" validate:"required,min=1"`
- MaxPreProcessThreads int `json:"MaxPreProcessThreads" validate:"min=1;gtefield=MaxProcessThreads"`
- ForceRewrite bool `json:"ForceRewrite" validate:"required"`
- LogLevel slog.Level `json:"LogLevel" validate:"required"`
+ Input InputConfig `json:"Input" validate:"required"`
+ Converter ConverterConfig `json:"Converter" validate:"required"`
+ Output OutputConfig `json:"Output" validate:"required"`
+ MaxConcurrentJobs int `json:"MaxConcurrentJobs" validate:"required,min=1"`
+ ForceRewrite bool `json:"ForceRewrite" validate:"required"`
+ LogLevel slog.Level `json:"LogLevel" validate:"required"`
}
type InputConfig struct {
diff --git a/main.go b/main.go
index bb3e208..d41baa9 100644
--- a/main.go
+++ b/main.go
@@ -86,13 +86,11 @@ func main() {
default:
}
- processSemaphore := make(chan struct{}, cfg.MaxProcessThreads)
- queueSemaphore := make(chan struct{}, cfg.MaxPreProcessThreads)
+ semaphore := make(chan struct{}, cfg.MaxConcurrentJobs)
var wg sync.WaitGroup
wg.Add(len(files))
for i, file := range files {
- queueSemaphore <- struct{}{}
go func(index int, inputName string) {
outputName := conv.DeductOutputPath(inputName)
fileLogger := generalLogger.With(slog.String("input_path", inputName), slog.String("output_path", outputName), slog.Int("file_index", index))
@@ -100,8 +98,9 @@ func main() {
threadSigTermChannel := make(chan os.Signal, 1)
signal.Notify(threadSigTermChannel, os.Interrupt, syscall.SIGTERM)
- defer func() { <-queueSemaphore; wg.Done() }()
-
+ defer func() { <-semaphore }()
+ defer wg.Done()
+ semaphore <- struct{}{}
select {
case <-threadSigTermChannel:
fileLogger.Info("exiting due to termination signal")
@@ -129,9 +128,6 @@ func main() {
}
}
- processSemaphore <- struct{}{}
- defer func() { <-processSemaphore }()
-
fileLogger.Info("start to process file", slog.String("input_hash", inputMetadata.Hash), slog.String("original_input_hash", originalInputHash))
reader, err := inputClient.GetReader(inputName)
if err != nil {