summaryrefslogtreecommitdiff
diff refs
from: back
to: back
| flip
diff options
context:
space:
mode:
-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, 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)
diff --git a/main.go b/main.go
index 5735e19..521960e 100644
--- a/main.go
+++ b/main.go
@@ -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()