From 115a9cf80d07043903031ab97719782d84de9f60 Mon Sep 17 00:00:00 2001 From: niwaniwa Date: Sun, 31 Aug 2025 03:56:22 +0900 Subject: [PATCH] add: s3 upload and backup --- pkg/bot/bot.go | 5 + pkg/config/config.go | 37 +++++- pkg/container/container.go | 10 +- pkg/handler/processor/upload_handler.go | 82 +++++++++---- pkg/job/database_backup_job.go | 128 ++++++++++++++++++++ pkg/repository/s3_repository.go | 148 ++++++++++++++++++++++++ pkg/utility/utilities.go | 15 +++ 7 files changed, 402 insertions(+), 23 deletions(-) create mode 100644 pkg/job/database_backup_job.go create mode 100644 pkg/repository/s3_repository.go diff --git a/pkg/bot/bot.go b/pkg/bot/bot.go index 61a02ae..725e9e3 100644 --- a/pkg/bot/bot.go +++ b/pkg/bot/bot.go @@ -91,4 +91,9 @@ func (b *Bot) setupHandlers() error { func (b *Bot) startJobs() { b.container.ChannelDeleteJob.Run() b.container.EmojiUpdateInfoJob.Run() + + // データベースバックアップJob(設定で有効化されている場合のみ) + if err := b.container.DatabaseBackupJob.Run(); err != nil { + fmt.Printf("Database backup failed: %v\n", err) + } } diff --git a/pkg/config/config.go b/pkg/config/config.go index fcb15a2..c249732 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -22,6 +22,16 @@ type Config struct { SavePath string DatabasePath string IsDebug bool + + // S3 Configuration + AWSAccessKeyID string + AWSSecretAccessKey string + AWSRegion string + S3Bucket string + S3Endpoint string + S3ForcePathStyle bool + UseS3 bool + EnableDatabaseBackup bool } func LoadConfig() (*Config, error) { @@ -30,11 +40,26 @@ func LoadConfig() (*Config, error) { return nil, errors.Config("failed to load settings.env", err) } - isDebug, err := strconv.ParseBool(os.Getenv("is_debug")) + isDebug, err := strconv.ParseBool(os.Getenv("debug")) if err != nil { isDebug = false } + useS3, err := strconv.ParseBool(os.Getenv("use_s3")) + if err != nil { + useS3 = false + } + + s3ForcePathStyle, err := strconv.ParseBool(os.Getenv("s3_force_path_style")) + if err != nil { + s3ForcePathStyle = false + } + + enableDatabaseBackup, err := strconv.ParseBool(os.Getenv("enable_database_backup")) + if err != nil { + enableDatabaseBackup = false + } + config := &Config{ GuildID: strings.TrimSpace(os.Getenv("guild_id")), BotToken: strings.TrimSpace(os.Getenv("bot_token")), @@ -47,6 +72,16 @@ func LoadConfig() (*Config, error) { SavePath: strings.TrimSpace(os.Getenv("save_path")), DatabasePath: strings.TrimSpace(os.Getenv("database_path")), IsDebug: isDebug, + + // S3 Configuration + AWSAccessKeyID: strings.TrimSpace(os.Getenv("aws_access_key_id")), + AWSSecretAccessKey: strings.TrimSpace(os.Getenv("aws_secret_access_key")), + AWSRegion: strings.TrimSpace(os.Getenv("aws_region")), + S3Bucket: strings.TrimSpace(os.Getenv("s3_bucket")), + S3Endpoint: strings.TrimSpace(os.Getenv("s3_endpoint")), + S3ForcePathStyle: s3ForcePathStyle, + UseS3: useS3, + EnableDatabaseBackup: enableDatabaseBackup, } // Set default values diff --git a/pkg/container/container.go b/pkg/container/container.go index c7ee7e8..4baece3 100644 --- a/pkg/container/container.go +++ b/pkg/container/container.go @@ -33,8 +33,9 @@ type Container struct { ComponentHandler handler.ComponentHandler // Jobs - ChannelDeleteJob job.Job - EmojiUpdateInfoJob job.Job + ChannelDeleteJob job.Job + EmojiUpdateInfoJob job.Job + DatabaseBackupJob job.DatabaseBackupJob // Discord Session Session *discordgo.Session @@ -85,6 +86,10 @@ func NewContainer(cfg *config.Config) (*Container, error) { // Initialize jobs channelDeleteJob := job.NewChannelDeleteJob(emojiRepo, discordRepo) emojiUpdateInfoJob := job.NewEmojiUpdateInfoJob(emojiRepo, misskeyRepo) + databaseBackupJob, err := job.NewDatabaseBackupJob(cfg) + if err != nil { + return nil, errors.Config("failed to initialize database backup job", err) + } container := &Container{ Config: cfg, @@ -98,6 +103,7 @@ func NewContainer(cfg *config.Config) (*Container, error) { ComponentHandler: componentHandler, ChannelDeleteJob: channelDeleteJob, EmojiUpdateInfoJob: emojiUpdateInfoJob, + DatabaseBackupJob: databaseBackupJob, Session: session, Version: string(version), } diff --git a/pkg/handler/processor/upload_handler.go b/pkg/handler/processor/upload_handler.go index 4e43617..aba4703 100644 --- a/pkg/handler/processor/upload_handler.go +++ b/pkg/handler/processor/upload_handler.go @@ -1,6 +1,7 @@ package processor import ( + "bytes" "os" "path/filepath" @@ -9,15 +10,21 @@ import ( "MisskeyEmojiBot/pkg/config" "MisskeyEmojiBot/pkg/entity" "MisskeyEmojiBot/pkg/handler" + "MisskeyEmojiBot/pkg/repository" "MisskeyEmojiBot/pkg/utility" ) type uploadHandler struct { config config.Config + s3Repo repository.S3Repository } func NewUploadHandler(cfg config.Config) handler.EmojiProcessHandler { - return &uploadHandler{config: cfg} + var s3Repo repository.S3Repository + if cfg.UseS3 { + s3Repo, _ = repository.NewS3Repository(&cfg) + } + return &uploadHandler{config: cfg, s3Repo: s3Repo} } func (h *uploadHandler) Request(emoji *entity.Emoji, s *discordgo.Session, cID string) (entity.Response, error) { @@ -42,27 +49,62 @@ func (h *uploadHandler) Response(emoji *entity.Emoji, s *discordgo.Session, m *d "対応ファイルは`.png`,`.jpg`,`.jpeg`,`.gif`です。") return response, nil } - emoji.FilePath = filepath.Join(h.config.SavePath, emoji.ID+ext) - err := utility.EmojiDownload(attachment.URL, emoji.FilePath) - if err != nil { - _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ - "申請中にエラーが発生しました。URLを確認して再アップロードを行うか、管理者へ問い合わせを行ってください。#01a") - return response, nil - } + if h.config.UseS3 && h.s3Repo != nil { + fileData, err := utility.EmojiDownloadToBytes(attachment.URL) + if err != nil { + _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ + "申請中にエラーが発生しました。URLを確認して再アップロードを行うか、管理者へ問い合わせを行ってください。#01a") + return response, nil + } - file, err := os.Open(emoji.FilePath) - if err != nil { - _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ - "申請中にエラーが発生しました。管理者へ問い合わせを行ってください。#01b") - return response, nil - } - defer func() { _ = file.Close() }() + key := "emojis/" + emoji.ID + ext + contentType := repository.GetContentTypeFromExtension(attachment.Filename) + fileURL, err := h.s3Repo.UploadFile(key, fileData, contentType) + if err != nil { + _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ + "申請中にエラーが発生しました。S3へのアップロードに失敗しました。#01sa") + return response, nil + } - _, err = s.ChannelFileSend(m.ChannelID, emoji.FilePath, file) - if err != nil { - _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ - "申請中にエラーが発生しました。管理者へ問い合わせを行ってください。#01d") - return response, nil + emoji.FilePath = fileURL + + _, err = s.ChannelMessageSendComplex(m.ChannelID, &discordgo.MessageSend{ + Content: "アップロード完了: " + fileURL, + Files: []*discordgo.File{ + { + Name: attachment.Filename, + Reader: bytes.NewReader(fileData), + }, + }, + }) + if err != nil { + _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ + "申請中にエラーが発生しました。管理者へ問い合わせを行ってください。#01sb") + return response, nil + } + } else { + emoji.FilePath = filepath.Join(h.config.SavePath, emoji.ID+ext) + err := utility.EmojiDownload(attachment.URL, emoji.FilePath) + if err != nil { + _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ + "申請中にエラーが発生しました。URLを確認して再アップロードを行うか、管理者へ問い合わせを行ってください。#01a") + return response, nil + } + + file, err := os.Open(emoji.FilePath) + if err != nil { + _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ + "申請中にエラーが発生しました。管理者へ問い合わせを行ってください。#01b") + return response, nil + } + defer func() { _ = file.Close() }() + + _, err = s.ChannelFileSend(m.ChannelID, emoji.FilePath, file) + if err != nil { + _, _ = s.ChannelMessageSend(m.ChannelID, ": Error! \n"+ + "申請中にエラーが発生しました。管理者へ問い合わせを行ってください。#01d") + return response, nil + } } response.IsSuccess = true diff --git a/pkg/job/database_backup_job.go b/pkg/job/database_backup_job.go new file mode 100644 index 0000000..6083c9f --- /dev/null +++ b/pkg/job/database_backup_job.go @@ -0,0 +1,128 @@ +package job + +import ( + "fmt" + "io" + "os" + "time" + + "MisskeyEmojiBot/pkg/config" + "MisskeyEmojiBot/pkg/repository" +) + +type DatabaseBackupJob interface { + Run() error +} + +type databaseBackupJob struct { + config *config.Config + s3Repo repository.S3Repository +} + +func NewDatabaseBackupJob(cfg *config.Config) (DatabaseBackupJob, error) { + job := &databaseBackupJob{config: cfg} + + if cfg.UseS3 { + s3Repo, err := repository.NewS3Repository(cfg) + if err != nil { + return nil, fmt.Errorf("failed to create S3 repository: %w", err) + } + job.s3Repo = s3Repo + } + + return job, nil +} + +func (j *databaseBackupJob) Run() error { + + if !j.config.EnableDatabaseBackup { + fmt.Println("Database backup job is disabled in configuration.") + return nil + } + + if j.config.DatabasePath == "" { + return fmt.Errorf("database path is not configured") + } + + cleanRequest := time.NewTicker(12 * time.Hour) + go func() { + for range cleanRequest.C { + if !j.config.UseS3 || j.s3Repo == nil { + // S3が有効でない場合はローカルバックアップのみ + j.createLocalBackup() + } + + // データベースファイルを読み込む + dbData, err := j.readDatabaseFile() + if err != nil { + fmt.Errorf("failed to read database file: %w", err) + } + + // S3にバックアップをアップロード + timestamp := time.Now().Format("2006-01-02_15-04-05") + backupKey := fmt.Sprintf("backups/database_%s.db", timestamp) + + _, err = j.s3Repo.UploadFile(backupKey, dbData, "application/octet-stream") + if err != nil { + fmt.Errorf("failed to upload backup to S3: %w", err) + } + + fmt.Printf("Database backup uploaded to S3: %s\n", backupKey) + + // ローカルバックアップも作成(オプション) + if err := j.createLocalBackup(); err != nil { + fmt.Printf("Warning: Local backup failed: %v\n", err) + } + } + }() + + return nil +} + +func (j *databaseBackupJob) readDatabaseFile() ([]byte, error) { + file, err := os.Open(j.config.DatabasePath) + if err != nil { + return nil, err + } + defer file.Close() + + data, err := io.ReadAll(file) + if err != nil { + return nil, err + } + + return data, nil +} + +func (j *databaseBackupJob) createLocalBackup() error { + timestamp := time.Now().Format("2006-01-02_15-04-05") + backupPath := fmt.Sprintf("%sbackup_%s.db", j.config.SavePath, timestamp) + + // 元のファイルを開く + src, err := os.Open(j.config.DatabasePath) + if err != nil { + return err + } + defer src.Close() + + // バックアップディレクトリを作成 + if err := os.MkdirAll(j.config.SavePath, os.ModePerm); err != nil { + return err + } + + // バックアップファイルを作成 + dst, err := os.Create(backupPath) + if err != nil { + return err + } + defer dst.Close() + + // ファイルをコピー + _, err = io.Copy(dst, src) + if err != nil { + return err + } + + fmt.Printf("Local database backup created: %s\n", backupPath) + return nil +} diff --git a/pkg/repository/s3_repository.go b/pkg/repository/s3_repository.go new file mode 100644 index 0000000..73d742f --- /dev/null +++ b/pkg/repository/s3_repository.go @@ -0,0 +1,148 @@ +package repository + +import ( + "bytes" + "context" + "fmt" + "io" + "path/filepath" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" + + cfg "MisskeyEmojiBot/pkg/config" + "MisskeyEmojiBot/pkg/errors" +) + +type S3Repository interface { + UploadFile(key string, data []byte, contentType string) (string, error) + DownloadFile(key string) ([]byte, error) + DeleteFile(key string) error + GetFileURL(key string) string +} + +type s3Repository struct { + client *s3.Client + bucket string + config *cfg.Config +} + +func NewS3Repository(cfg *cfg.Config) (S3Repository, error) { + if !cfg.UseS3 { + return nil, fmt.Errorf("S3 is not enabled in configuration") + } + + // Create AWS config with custom credentials + awsConfig, err := config.LoadDefaultConfig(context.TODO(), + config.WithRegion(cfg.AWSRegion), + config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider( + cfg.AWSAccessKeyID, + cfg.AWSSecretAccessKey, + "", + )), + ) + if err != nil { + return nil, errors.Config("failed to load AWS config", err) + } + + // Create S3 client with optional custom endpoint (for S3-compatible services) + var client *s3.Client + if cfg.S3Endpoint != "" { + client = s3.NewFromConfig(awsConfig, func(o *s3.Options) { + o.BaseEndpoint = aws.String(cfg.S3Endpoint) + o.UsePathStyle = cfg.S3ForcePathStyle + }) + } else { + client = s3.NewFromConfig(awsConfig) + } + + return &s3Repository{ + client: client, + bucket: cfg.S3Bucket, + config: cfg, + }, nil +} + +func (r *s3Repository) UploadFile(key string, data []byte, contentType string) (string, error) { + if contentType == "" { + contentType = "application/octet-stream" + } + + _, err := r.client.PutObject(context.TODO(), &s3.PutObjectInput{ + Bucket: aws.String(r.bucket), + Key: aws.String(key), + Body: bytes.NewReader(data), + ContentType: aws.String(contentType), + }) + if err != nil { + return "", errors.FileOperation("failed to upload file to S3", err) + } + + return r.GetFileURL(key), nil +} + +func (r *s3Repository) DownloadFile(key string) ([]byte, error) { + result, err := r.client.GetObject(context.TODO(), &s3.GetObjectInput{ + Bucket: aws.String(r.bucket), + Key: aws.String(key), + }) + if err != nil { + return nil, errors.FileOperation("failed to download file from S3", err) + } + defer result.Body.Close() + + data, err := io.ReadAll(result.Body) + if err != nil { + return nil, errors.FileOperation("failed to read file data from S3", err) + } + + return data, nil +} + +func (r *s3Repository) DeleteFile(key string) error { + _, err := r.client.DeleteObject(context.TODO(), &s3.DeleteObjectInput{ + Bucket: aws.String(r.bucket), + Key: aws.String(key), + }) + if err != nil { + return errors.FileOperation("failed to delete file from S3", err) + } + + return nil +} + +func (r *s3Repository) GetFileURL(key string) string { + if r.config.S3Endpoint != "" { + // For custom endpoints (S3-compatible services) + if r.config.S3ForcePathStyle { + return fmt.Sprintf("%s/%s/%s", r.config.S3Endpoint, r.bucket, key) + } + return fmt.Sprintf("https://%s.%s/%s", r.bucket, r.config.S3Endpoint, key) + } + + // For AWS S3 + return fmt.Sprintf("https://%s.s3.%s.amazonaws.com/%s", r.bucket, r.config.AWSRegion, key) +} + +// Helper function to determine content type from file extension +func GetContentTypeFromExtension(filename string) string { + ext := filepath.Ext(filename) + switch ext { + case ".png": + return "image/png" + case ".jpg", ".jpeg": + return "image/jpeg" + case ".gif": + return "image/gif" + case ".webp": + return "image/webp" + case ".json": + return "application/json" + case ".db": + return "application/octet-stream" + default: + return "application/octet-stream" + } +} diff --git a/pkg/utility/utilities.go b/pkg/utility/utilities.go index 171cd1b..2dcf3ca 100644 --- a/pkg/utility/utilities.go +++ b/pkg/utility/utilities.go @@ -35,3 +35,18 @@ func EmojiDownload(url string, filePath string) error { return nil } + +func EmojiDownloadToBytes(url string) ([]byte, error) { + response, err := http.Get(url) + if err != nil { + return nil, err + } + defer func() { _ = response.Body.Close() }() + + data, err := io.ReadAll(response.Body) + if err != nil { + return nil, err + } + + return data, nil +}