refactor: maintain persistent database connections in instance status and correct wait order for backup processes
Build and Push / build (godump, amd64, linux) (push) Successful in 29s
Build and Push / build (godump, amd64, linux) (push) Successful in 29s
This commit is contained in:
+3
-6
@@ -9,13 +9,10 @@ import (
|
|||||||
_ "github.com/go-sql-driver/mysql"
|
_ "github.com/go-sql-driver/mysql"
|
||||||
)
|
)
|
||||||
|
|
||||||
func discoverDatabases(cfg config.InstanceConfig) ([]string, error) {
|
func discoverDatabases(db *sql.DB, cfg config.InstanceConfig) ([]string, error) {
|
||||||
dsn := fmt.Sprintf("%s:%s@tcp(%s:%d)/?parseTime=true", cfg.User, cfg.Password, cfg.Host, cfg.Port)
|
if db == nil {
|
||||||
db, err := sql.Open("mysql", dsn)
|
return nil, fmt.Errorf("database pool is not initialized")
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
}
|
||||||
defer db.Close()
|
|
||||||
|
|
||||||
// Ensure the connection is actually valid
|
// Ensure the connection is actually valid
|
||||||
if err := db.Ping(); err != nil {
|
if err := db.Ping(); err != nil {
|
||||||
|
|||||||
+14
-2
@@ -1,12 +1,16 @@
|
|||||||
package backup
|
package backup
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"sort"
|
"sort"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
_ "github.com/go-sql-driver/mysql"
|
||||||
|
|
||||||
"godump/config"
|
"godump/config"
|
||||||
"godump/logger"
|
"godump/logger"
|
||||||
"godump/notify"
|
"godump/notify"
|
||||||
@@ -23,6 +27,7 @@ type DBStatus struct {
|
|||||||
|
|
||||||
type InstanceStatus struct {
|
type InstanceStatus struct {
|
||||||
Config config.InstanceConfig
|
Config config.InstanceConfig
|
||||||
|
DB *sql.DB
|
||||||
LastRunTime time.Time
|
LastRunTime time.Time
|
||||||
NextRunTime time.Time
|
NextRunTime time.Time
|
||||||
OverallResult string // success, partial, failed, running
|
OverallResult string // success, partial, failed, running
|
||||||
@@ -97,8 +102,15 @@ func NewManager(cfg *config.Config) *Manager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for _, instCfg := range cfg.Instances {
|
for _, instCfg := range cfg.Instances {
|
||||||
|
dsn := fmt.Sprintf("%s:%s@tcp(%s:%d)/?parseTime=true", instCfg.User, instCfg.Password, instCfg.Host, instCfg.Port)
|
||||||
|
db, err := sql.Open("mysql", dsn)
|
||||||
|
if err != nil {
|
||||||
|
logger.Error(instCfg.Name, "Failed to initialize database pool: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
status := &InstanceStatus{
|
status := &InstanceStatus{
|
||||||
Config: instCfg,
|
Config: instCfg,
|
||||||
|
DB: db,
|
||||||
Databases: make(map[string]*DBStatus),
|
Databases: make(map[string]*DBStatus),
|
||||||
}
|
}
|
||||||
m.instances[instCfg.Name] = status
|
m.instances[instCfg.Name] = status
|
||||||
@@ -148,7 +160,7 @@ func (m *Manager) DiscoverInitial() {
|
|||||||
m.mu.RLock()
|
m.mu.RLock()
|
||||||
defer m.mu.RUnlock()
|
defer m.mu.RUnlock()
|
||||||
for name, inst := range m.instances {
|
for name, inst := range m.instances {
|
||||||
dbs, err := discoverDatabases(inst.Config)
|
dbs, err := discoverDatabases(inst.DB, inst.Config)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error(name, "Initial database discovery failed: %v", err)
|
logger.Error(name, "Initial database discovery failed: %v", err)
|
||||||
inst.mu.Lock()
|
inst.mu.Lock()
|
||||||
@@ -246,7 +258,7 @@ func (m *Manager) RunInstance(name string) {
|
|||||||
logger.Info(name, "Starting backup job")
|
logger.Info(name, "Starting backup job")
|
||||||
|
|
||||||
// 1. Discovery
|
// 1. Discovery
|
||||||
dbs, err := discoverDatabases(inst.Config)
|
dbs, err := discoverDatabases(inst.DB, inst.Config)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error(name, "Database discovery failed: %v", err)
|
logger.Error(name, "Database discovery failed: %v", err)
|
||||||
inst.mu.Lock()
|
inst.mu.Lock()
|
||||||
|
|||||||
+7
-5
@@ -67,13 +67,15 @@ func backupDatabase(cfg config.InstanceConfig, dbName string) (int64, error) {
|
|||||||
return 0, fmt.Errorf("failed to start gzip: %w", err)
|
return 0, fmt.Errorf("failed to start gzip: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := cmdDump.Wait(); err != nil {
|
errGzip := cmdGzip.Wait()
|
||||||
cmdGzip.Process.Kill()
|
errDump := cmdDump.Wait()
|
||||||
return 0, fmt.Errorf("mysqldump failed: %w", err)
|
|
||||||
|
if errDump != nil {
|
||||||
|
return 0, fmt.Errorf("mysqldump failed: %w", errDump)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := cmdGzip.Wait(); err != nil {
|
if errGzip != nil {
|
||||||
return 0, fmt.Errorf("gzip failed: %w", err)
|
return 0, fmt.Errorf("gzip failed: %w", errGzip)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get file size
|
// Get file size
|
||||||
|
|||||||
Reference in New Issue
Block a user