summaryrefslogtreecommitdiff
diff options
from:
to:
context:
space:
mode:
authorGravatar SayaAndy <saya.andy@posteo.com> 2026-02-09 06:33:35 +0700
committerGravatar SayaAndy <saya.andy@posteo.com> 2026-02-09 06:33:35 +0700
commitebfe3540c333ae7fa3effb410a20e5764627e307 (patch)
tree543aa1d73463dcd381f732d67b3361dd53a926a0
parentda672112ebbc25a3fa9f2cf8e37994a7bd1d0e2f (diff)
downloadthumbnail-generator-ebfe3540c333ae7fa3effb410a20e5764627e307.tar.gz
thumbnail-generator-ebfe3540c333ae7fa3effb410a20e5764627e307.zip
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
-rw-r--r--config/config-b2.sample.json6
-rw-r--r--config/config-local-unix.sample.json2
-rw-r--r--config/config.go30
-rw-r--r--internal/client/output/b2.go12
-rw-r--r--internal/client/output/local_unix.go2
-rw-r--r--internal/client/output/output_client_interface.go2
-rw-r--r--internal/converter/converter_interface.go1
-rw-r--r--internal/converter/jpeg.go114
-rw-r--r--internal/converter/webp.go27
-rw-r--r--main.go62
10 files changed, 216 insertions, 42 deletions
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()