| -rw-r--r-- | config/config-b2.sample.json | 6 | ||||
| -rw-r--r-- | config/config-local-unix.sample.json | 2 | ||||
| -rw-r--r-- | config/config.go | 30 | ||||
| -rw-r--r-- | internal/client/output/b2.go | 12 | ||||
| -rw-r--r-- | internal/client/output/local_unix.go | 2 | ||||
| -rw-r--r-- | internal/client/output/output_client_interface.go | 2 | ||||
| -rw-r--r-- | internal/converter/converter_interface.go | 1 | ||||
| -rw-r--r-- | internal/converter/jpeg.go | 114 | ||||
| -rw-r--r-- | internal/converter/webp.go | 27 | ||||
| -rw-r--r-- | main.go | 62 |
10 files changed, 42 insertions, 216 deletions
diff --git a/config/config-b2.sample.json b/config/config-b2.sample.json index 60022dd..372972e 100644 --- a/config/config-b2.sample.json +++ b/config/config-b2.sample.json @@ -1,4 +1,5 @@ { + "ForceRewrite": false, "MaxProcessThreads": 4, "MaxPreProcessThreads": 12, "LogLevel": "info", @@ -32,7 +33,6 @@ } }, "Output": { - "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -55,7 +55,6 @@ } }, "Output": { - "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -78,7 +77,6 @@ } }, "Output": { - "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -101,7 +99,6 @@ } }, "Output": { - "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -124,7 +121,6 @@ } }, "Output": { - "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { diff --git a/config/config-local-unix.sample.json b/config/config-local-unix.sample.json index 941d995..97f69f1 100644 --- a/config/config-local-unix.sample.json +++ b/config/config-local-unix.sample.json @@ -1,4 +1,5 @@ { + "ForceRewrite": false, "MaxProcessThreads": 4, "MaxPreProcessThreads": 12, "LogLevel": "debug", @@ -26,7 +27,6 @@ } }, "Output": { - "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "local-unix", "Config": { diff --git a/config/config.go b/config/config.go index e55d4b7..1620d75 100644 --- a/config/config.go +++ b/config/config.go @@ -14,6 +14,7 @@ type Config struct { Converters []ConverterConfig `json:"Converters" 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"` } @@ -62,7 +63,7 @@ func (sc *InputStorageConfig) UnmarshalJSON(data []byte) error { } type ConverterConfig struct { - Type string `json:"Type" validate:"required,oneof=webp jpeg"` + Type string `json:"Type" validate:"required,oneof=webp"` Config any `json:"Config" validate:"required"` Output OutputConfig `json:"Output" validate:"required"` } @@ -88,12 +89,6 @@ func (pc *ConverterConfig) UnmarshalJSON(data []byte) error { return fmt.Errorf("unmarshal WebpConfig: %w", err) } pc.Config = &webpConfig - case "jpeg": - var jpegConfig JpegConfig - if err := json.Unmarshal(tmp.Config, &jpegConfig); err != nil { - return fmt.Errorf("unmarshal JpegConfig: %w", err) - } - pc.Config = &jpegConfig default: return fmt.Errorf("unsupported storage type: %s", tmp.Type) } @@ -103,18 +98,12 @@ func (pc *ConverterConfig) UnmarshalJSON(data []byte) error { type WebpConfig struct { Quality int `json:"Quality" validate:"required,min=1,max=100"` - Size SizeConfig `json:"Size"` -} - -type JpegConfig struct { - ExtensionName string `json:"ExtensionName" validate:"alpha"` - Quality int `json:"Quality" validate:"required,min=1,max=100"` - Size SizeConfig `json:"Size"` + Size SizeConfig `json:"Size" validate:"required"` } type SizeConfig struct { - MaxWidth int `json:"MaxWidth"` - MaxHeight int `json:"MaxHeight"` + MaxWidth int `json:"MaxWidth" validate:"required,min=0"` + MaxHeight int `json:"MaxHeight" validate:"required,min=0"` } type OutputStorageConfig struct { @@ -168,8 +157,7 @@ type InputLocalUnixConfig struct { } type OutputConfig struct { - RewriteOn string `json:"RewriteOn" validate:"oneof=Never UnequalHashInCache Always"` - Storage OutputStorageConfig `json:"Storage" validate:"required"` + Storage OutputStorageConfig `json:"Storage" validate:"required"` } type OutputLocalUnixConfig struct { @@ -191,12 +179,6 @@ func LoadConfig(path string, config *Config) error { return err } - for i := range config.Converters { - if config.Converters[i].Output.RewriteOn == "" { - config.Converters[i].Output.RewriteOn = "UnequalHashInCache" - } - } - return nil } diff --git a/internal/client/output/b2.go b/internal/client/output/b2.go index 0616368..f53532a 100644 --- a/internal/client/output/b2.go +++ b/internal/client/output/b2.go @@ -2,10 +2,8 @@ package output import ( "context" - "encoding/json" "fmt" "io" - "log/slog" "github.com/Backblaze/blazer/b2" "github.com/SayaAndy/saya-today-thumbnail-generator/config" @@ -39,7 +37,7 @@ func NewB2OutputClient(cfg *config.OutputConfig) (OutputClient, error) { return &B2OutputClient{b2cl: b2cl, bucket: bucket, prefix: b2cfg.Prefix}, nil } -func (c *B2OutputClient) GetWriter(path string, inputMetadata *input.MetadataStruct, outputContentType string) (io.WriteCloser, error) { +func (c *B2OutputClient) GetWriter(path string, inputMetadata *input.MetadataStruct) (io.WriteCloser, error) { obj := c.bucket.Object(c.prefix + path) if obj == nil { return nil, fmt.Errorf("failed to reference object in B2 bucket") @@ -47,7 +45,6 @@ func (c *B2OutputClient) GetWriter(path string, inputMetadata *input.MetadataStr attrs := &b2.Attrs{Info: make(map[string]string)} attrs.Info["sha1-original"] = inputMetadata.Hash - attrs.ContentType = outputContentType return obj.NewWriter(context.Background(), b2.WithAttrsOption(attrs)), nil } @@ -102,12 +99,5 @@ func (c *B2OutputClient) IsMissing(path string) bool { return true } - attrsJson, _ := json.Marshal(attrs) - slog.Debug("got object attrs", slog.String("path", path), slog.String("attrs", string(attrsJson))) - - if attrs.Size == 0 { - return true - } - return attrs.Status == b2.Hider } diff --git a/internal/client/output/local_unix.go b/internal/client/output/local_unix.go index 6f9a7b3..668a14b 100644 --- a/internal/client/output/local_unix.go +++ b/internal/client/output/local_unix.go @@ -44,7 +44,7 @@ func NewLocalUnixOutputClient(cfg *config.OutputConfig) (OutputClient, error) { return &LocalUnixOutputClient{localCfg.Path, uint32(fpm), uint32(dpm), localCfg.AttributesImplementation}, nil } -func (c *LocalUnixOutputClient) GetWriter(path string, inputMetadata *input.MetadataStruct, _ string) (io.WriteCloser, error) { +func (c *LocalUnixOutputClient) GetWriter(path string, inputMetadata *input.MetadataStruct) (io.WriteCloser, error) { pathSegments := strings.Split(path, "/") dirpath := strings.Join(pathSegments[0:len(pathSegments)-1], "/") if err := os.MkdirAll(c.path+dirpath, os.FileMode(c.dirMode)); err != nil { diff --git a/internal/client/output/output_client_interface.go b/internal/client/output/output_client_interface.go index f9e421b..c74851d 100644 --- a/internal/client/output/output_client_interface.go +++ b/internal/client/output/output_client_interface.go @@ -9,7 +9,7 @@ import ( ) type OutputClient interface { - GetWriter(path string, inputMetadata *input.MetadataStruct, outputContentType string) (io.WriteCloser, error) + GetWriter(path string, inputMetadata *input.MetadataStruct) (io.WriteCloser, error) ReadMetadata(path string) (*MetadataStruct, error) IsMissing(path string) bool } diff --git a/internal/converter/converter_interface.go b/internal/converter/converter_interface.go index abff99d..fda23d1 100644 --- a/internal/converter/converter_interface.go +++ b/internal/converter/converter_interface.go @@ -17,5 +17,4 @@ type Converter interface { var NewConverterMap = map[string]func(cfg *config.ConverterConfig) (Converter, error){ "webp": NewWebpConverter, - "jpeg": NewJpegConverter, } diff --git a/internal/converter/jpeg.go b/internal/converter/jpeg.go deleted file mode 100644 index a35a47c..0000000 --- a/internal/converter/jpeg.go +++ /dev/null @@ -1,114 +0,0 @@ -package converter - -import ( - "fmt" - "image" - "image/jpeg" - "image/png" - "io" - "log/slog" - "path/filepath" - "strings" - - "github.com/SayaAndy/saya-today-thumbnail-generator/config" - "github.com/SayaAndy/saya-today-thumbnail-generator/internal/client/input" - "github.com/SayaAndy/saya-today-thumbnail-generator/internal/client/output" - "github.com/kolesa-team/go-webp/webp" - "golang.org/x/image/draw" -) - -var _ Converter = (*JpegConverter)(nil) - -type JpegConverter struct { - maxWidth int - maxHeight int - extensionName string - quality int - outputClient output.OutputClient -} - -func NewJpegConverter(cfg *config.ConverterConfig) (Converter, error) { - if cfg.Type != "jpeg" { - return nil, fmt.Errorf("invalid storage type for JpegConverter") - } - jpegCfg := cfg.Config.(*config.JpegConfig) - - outputClient, err := output.NewOutputClientMap[cfg.Output.Storage.Type](&cfg.Output) - if err != nil { - return nil, fmt.Errorf("fail to initialize output client: %w", err) - } - - extensionName := ".jpg" - if jpegCfg.ExtensionName != "" { - extensionName = "." + jpegCfg.ExtensionName - } - - return &JpegConverter{jpegCfg.Size.MaxWidth, jpegCfg.Size.MaxHeight, extensionName, jpegCfg.Quality, outputClient}, nil -} - -func (p *JpegConverter) Process(inputMetadata *input.MetadataStruct, reader io.Reader, outputName string) error { - var src image.Image - - writer, err := p.outputClient.GetWriter(outputName, inputMetadata, "image/jpeg") - if err != nil { - return fmt.Errorf("fail to initialize writer for output: %w", err) - } - defer writer.Close() - - switch inputMetadata.ContentType { - case "image/jpeg": - src, err = jpeg.Decode(reader) - if err != nil { - return fmt.Errorf("decode jpeg: %w", err) - } - case "image/png": - src, err = png.Decode(reader) - if err != nil { - return fmt.Errorf("decode png: %w", err) - } - case "image/webp": - src, err = webp.Decode(reader, nil) - if err != nil { - return fmt.Errorf("decode webp: %w", err) - } - default: - return fmt.Errorf("unsupported content type: %s", inputMetadata.ContentType) - } - - xCoef := 1.0 - if p.maxWidth > 0 { - xCoef = float64(p.maxWidth) / float64(src.Bounds().Max.X) - } - yCoef := 1.0 - if p.maxHeight > 0 { - yCoef = float64(p.maxHeight) / float64(src.Bounds().Max.Y) - } - slog.Debug("calculated coefficients", slog.Float64("x_coef", xCoef), slog.Float64("y_coef", yCoef)) - - minCoef := xCoef - if yCoef < minCoef { - minCoef = yCoef - } - - if minCoef >= 1.0 { - return jpeg.Encode(writer, src, &jpeg.Options{Quality: p.quality}) - } - - dst := image.NewRGBA(image.Rect(0, 0, int(float64(src.Bounds().Max.X)*minCoef+0.5), int(float64(src.Bounds().Max.Y)*minCoef+0.5))) - draw.CatmullRom.Scale(dst, dst.Rect, src, src.Bounds(), draw.Over, nil) - - return jpeg.Encode(writer, dst, &jpeg.Options{Quality: p.quality}) -} - -func (p *JpegConverter) DeductOutputPath(inputPath string) string { - withoutExt, _ := strings.CutSuffix(inputPath, filepath.Ext(inputPath)) - return withoutExt + p.extensionName -} - -func (p *JpegConverter) ReadMetadata(path string) (*output.MetadataStruct, error) { - return p.outputClient.ReadMetadata(path) -} - -func (p *JpegConverter) IsMissing(path string) bool { - return p.outputClient.IsMissing(path) -} diff --git a/internal/converter/webp.go b/internal/converter/webp.go index fbffd01..8cfe5a7 100644 --- a/internal/converter/webp.go +++ b/internal/converter/webp.go @@ -44,7 +44,7 @@ func NewWebpConverter(cfg *config.ConverterConfig) (Converter, error) { func (p *WebpConverter) Process(inputMetadata *input.MetadataStruct, reader io.Reader, outputName string) error { var src image.Image - writer, err := p.outputClient.GetWriter(outputName, inputMetadata, "image/webp") + writer, err := p.outputClient.GetWriter(outputName, inputMetadata) if err != nil { return fmt.Errorf("fail to initialize writer for output: %w", err) } @@ -61,11 +61,6 @@ func (p *WebpConverter) Process(inputMetadata *input.MetadataStruct, reader io.R if err != nil { return fmt.Errorf("decode png: %w", err) } - case "image/webp": - src, err = webp.Decode(reader, nil) - if err != nil { - return fmt.Errorf("decode webp: %w", err) - } default: return fmt.Errorf("unsupported content type: %s", inputMetadata.ContentType) } @@ -75,25 +70,25 @@ func (p *WebpConverter) Process(inputMetadata *input.MetadataStruct, reader io.R return fmt.Errorf("create webp encoder options: %w", err) } - xCoef := 1.0 - if p.maxWidth > 0 { - xCoef = float64(p.maxWidth) / float64(src.Bounds().Max.X) + xCoef := float64(p.maxWidth) / float64(src.Bounds().Max.X) + if p.maxWidth == 0 { + xCoef = 1 } - yCoef := 1.0 - if p.maxHeight > 0 { - yCoef = float64(p.maxHeight) / float64(src.Bounds().Max.Y) + yCoef := float64(p.maxHeight) / float64(src.Bounds().Max.Y) + if p.maxHeight == 0 { + yCoef = 1 } slog.Debug("calculated coefficients", slog.Float64("x_coef", xCoef), slog.Float64("y_coef", yCoef)) + if xCoef > 1 && yCoef > 1 { + return webp.Encode(writer, src, opts) + } + minCoef := xCoef if yCoef < minCoef { minCoef = yCoef } - if minCoef >= 1.0 { - return webp.Encode(writer, src, opts) - } - dst := image.NewRGBA(image.Rect(0, 0, int(float64(src.Bounds().Max.X)*minCoef+0.5), int(float64(src.Bounds().Max.Y)*minCoef+0.5))) draw.CatmullRom.Scale(dst, dst.Rect, src, src.Bounds(), draw.Over, nil) @@ -128,38 +128,26 @@ 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 @@ -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) } @@ -261,7 +239,7 @@ func main() { cacheMap[id][converterHashes[convIndex]] = struct{}{} cacheMapMutex.Unlock() } - }(i, file, processTerminating) + }(i, file) } wg.Wait() |