From ebfe3540c333ae7fa3effb410a20e5764627e307 Mon Sep 17 00:00:00 2001 From: SayaAndy Date: Mon, 9 Feb 2026 06:33:35 +0700 Subject: feat: add jpeg converter fix: explicitly set content type for b2 feat: extend rewrite settings with "never overwrite" feat: move rewrite setting to converter config feat: allow to skip size config --- config/config-b2.sample.json | 6 +- config/config-local-unix.sample.json | 2 +- config/config.go | 30 ++++-- internal/client/output/b2.go | 12 ++- internal/client/output/local_unix.go | 2 +- internal/client/output/output_client_interface.go | 2 +- internal/converter/converter_interface.go | 1 + internal/converter/jpeg.go | 114 ++++++++++++++++++++++ internal/converter/webp.go | 27 ++--- main.go | 62 ++++++++---- 10 files changed, 216 insertions(+), 42 deletions(-) create mode 100644 internal/converter/jpeg.go diff --git a/config/config-b2.sample.json b/config/config-b2.sample.json index 372972e..60022dd 100644 --- a/config/config-b2.sample.json +++ b/config/config-b2.sample.json @@ -1,5 +1,4 @@ { - "ForceRewrite": false, "MaxProcessThreads": 4, "MaxPreProcessThreads": 12, "LogLevel": "info", @@ -33,6 +32,7 @@ } }, "Output": { + "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -55,6 +55,7 @@ } }, "Output": { + "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -77,6 +78,7 @@ } }, "Output": { + "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -99,6 +101,7 @@ } }, "Output": { + "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { @@ -121,6 +124,7 @@ } }, "Output": { + "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "b2", "Config": { diff --git a/config/config-local-unix.sample.json b/config/config-local-unix.sample.json index 97f69f1..941d995 100644 --- a/config/config-local-unix.sample.json +++ b/config/config-local-unix.sample.json @@ -1,5 +1,4 @@ { - "ForceRewrite": false, "MaxProcessThreads": 4, "MaxPreProcessThreads": 12, "LogLevel": "debug", @@ -27,6 +26,7 @@ } }, "Output": { + "RewriteOn": "UnequalHashInCache", "Storage": { "Type": "local-unix", "Config": { diff --git a/config/config.go b/config/config.go index 1620d75..e55d4b7 100644 --- a/config/config.go +++ b/config/config.go @@ -14,7 +14,6 @@ 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"` } @@ -63,7 +62,7 @@ func (sc *InputStorageConfig) UnmarshalJSON(data []byte) error { } type ConverterConfig struct { - Type string `json:"Type" validate:"required,oneof=webp"` + Type string `json:"Type" validate:"required,oneof=webp jpeg"` Config any `json:"Config" validate:"required"` Output OutputConfig `json:"Output" validate:"required"` } @@ -89,6 +88,12 @@ 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) } @@ -98,12 +103,18 @@ func (pc *ConverterConfig) UnmarshalJSON(data []byte) error { type WebpConfig struct { Quality int `json:"Quality" validate:"required,min=1,max=100"` - Size SizeConfig `json:"Size" validate:"required"` + 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"` } type SizeConfig struct { - MaxWidth int `json:"MaxWidth" validate:"required,min=0"` - MaxHeight int `json:"MaxHeight" validate:"required,min=0"` + MaxWidth int `json:"MaxWidth"` + MaxHeight int `json:"MaxHeight"` } type OutputStorageConfig struct { @@ -157,7 +168,8 @@ type InputLocalUnixConfig struct { } type OutputConfig struct { - Storage OutputStorageConfig `json:"Storage" validate:"required"` + RewriteOn string `json:"RewriteOn" validate:"oneof=Never UnequalHashInCache Always"` + Storage OutputStorageConfig `json:"Storage" validate:"required"` } type OutputLocalUnixConfig struct { @@ -179,6 +191,12 @@ 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 f53532a..0616368 100644 --- a/internal/client/output/b2.go +++ b/internal/client/output/b2.go @@ -2,8 +2,10 @@ package output import ( "context" + "encoding/json" "fmt" "io" + "log/slog" "github.com/Backblaze/blazer/b2" "github.com/SayaAndy/saya-today-thumbnail-generator/config" @@ -37,7 +39,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) (io.WriteCloser, error) { +func (c *B2OutputClient) GetWriter(path string, inputMetadata *input.MetadataStruct, outputContentType string) (io.WriteCloser, error) { obj := c.bucket.Object(c.prefix + path) if obj == nil { return nil, fmt.Errorf("failed to reference object in B2 bucket") @@ -45,6 +47,7 @@ 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 } @@ -99,5 +102,12 @@ 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 668a14b..6f9a7b3 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) (io.WriteCloser, error) { +func (c *LocalUnixOutputClient) GetWriter(path string, inputMetadata *input.MetadataStruct, _ string) (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 c74851d..f9e421b 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) (io.WriteCloser, error) + GetWriter(path string, inputMetadata *input.MetadataStruct, outputContentType string) (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 fda23d1..abff99d 100644 --- a/internal/converter/converter_interface.go +++ b/internal/converter/converter_interface.go @@ -17,4 +17,5 @@ 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 new file mode 100644 index 0000000..a35a47c --- /dev/null +++ b/internal/converter/jpeg.go @@ -0,0 +1,114 @@ +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 8cfe5a7..fbffd01 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) + writer, err := p.outputClient.GetWriter(outputName, inputMetadata, "image/webp") if err != nil { return fmt.Errorf("fail to initialize writer for output: %w", err) } @@ -61,6 +61,11 @@ 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) } @@ -70,25 +75,25 @@ func (p *WebpConverter) Process(inputMetadata *input.MetadataStruct, reader io.R return fmt.Errorf("create webp encoder options: %w", err) } - xCoef := float64(p.maxWidth) / float64(src.Bounds().Max.X) - if p.maxWidth == 0 { - xCoef = 1 + xCoef := 1.0 + if p.maxWidth > 0 { + xCoef = float64(p.maxWidth) / float64(src.Bounds().Max.X) } - yCoef := float64(p.maxHeight) / float64(src.Bounds().Max.Y) - if p.maxHeight == 0 { - yCoef = 1 + 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)) - 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) diff --git a/main.go b/main.go index 521960e..5735e19 100644 --- a/main.go +++ b/main.go @@ -128,26 +128,38 @@ 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)) - threadSigTermChannel := make(chan os.Signal, 1) - signal.Notify(threadSigTermChannel, os.Interrupt, syscall.SIGTERM) + select { + case <-sigTermChan: + generalLogger.Info("exiting due to termination signal") + processTerminating = true + default: + } + go func(index int, inputName string, earlyTerminate bool) { defer func() { <-queueSemaphore; wg.Done() }() - select { - case <-threadSigTermChannel: - fileLogger.Info("exiting due to termination signal") + fileLogger := generalLogger.With(slog.String("input_path", inputName), slog.Int("file_index", index)) + + if earlyTerminate { + fileLogger.Info("skip processing file (process is terminating)") return - default: } var inputMetadata *input.MetadataStruct @@ -185,21 +197,31 @@ func main() { } } - 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())) + switch cfg.Converters[j].Output.RewriteOn { + case "Never": + if !conv.IsMissing(outputName) { + convLogger.Info("skip already existing file") 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 "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 + } } + case "Always": } + convertersToLaunch = append(convertersToLaunch, j) } @@ -239,7 +261,7 @@ func main() { cacheMap[id][converterHashes[convIndex]] = struct{}{} cacheMapMutex.Unlock() } - }(i, file) + }(i, file, processTerminating) } wg.Wait() -- cgit v1.3.1+13