| package storage |
|
|
| import ( |
| "context" |
| "database/sql" |
| "fmt" |
| "log" |
| "os" |
| "path/filepath" |
| "strconv" |
| "strings" |
| "time" |
|
|
| "ccLoad/internal/config" |
| sqlstore "ccLoad/internal/storage/sql" |
|
|
| _ "github.com/go-sql-driver/mysql" |
| _ "modernc.org/sqlite" |
| ) |
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| func NewStore() (Store, error) { |
| mysqlDSN := os.Getenv("CCLOAD_MYSQL") |
|
|
| |
| if mysqlDSN == "" { |
| dbPath := os.Getenv("SQLITE_PATH") |
| if dbPath == "" { |
| dbPath = resolveSQLitePath() |
| } |
|
|
| store, err := createSQLiteStore(dbPath) |
| if err != nil { |
| return nil, fmt.Errorf("SQLite 初始化失败: %w", err) |
| } |
| log.Printf("使用 SQLite 存储(纯模式): %s", dbPath) |
| return store, nil |
| } |
|
|
| |
| enableHybrid := os.Getenv("CCLOAD_ENABLE_SQLITE_REPLICA") == "1" |
|
|
| |
| if !enableHybrid { |
| mysql, err := createMySQLStore(mysqlDSN) |
| if err != nil { |
| return nil, fmt.Errorf("MySQL 初始化失败: %w", err) |
| } |
| log.Print("使用 MySQL 存储(纯模式)") |
| return mysql, nil |
| } |
|
|
| |
| log.Print("[INFO] 启动混合存储模式(MySQL 主 + SQLite 缓存)") |
|
|
| |
| mysql, err := createMySQLStore(mysqlDSN) |
| if err != nil { |
| return nil, fmt.Errorf("MySQL 初始化失败: %w", err) |
| } |
| log.Print("[INFO] MySQL 主存储已连接") |
|
|
| |
| sqlitePath := os.Getenv("SQLITE_PATH") |
| if sqlitePath == "" { |
| sqlitePath = resolveSQLitePath() |
| } |
| sqlite, err := createSQLiteStore(sqlitePath) |
| if err != nil { |
| _ = mysql.Close() |
| return nil, fmt.Errorf("SQLite 初始化失败: %w", err) |
| } |
| log.Printf("[INFO] SQLite 本地缓存已创建: %s", sqlitePath) |
|
|
| |
| logDays := getLogSyncDays() |
| syncMgr := NewSyncManager(mysql, sqlite) |
|
|
| |
| restoreCtx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) |
| defer cancel() |
|
|
| if err := syncMgr.RestoreOnStartup(restoreCtx, logDays); err != nil { |
| _ = sqlite.Close() |
| _ = mysql.Close() |
| return nil, fmt.Errorf("数据恢复失败: %w", err) |
| } |
|
|
| |
| hybrid := NewHybridStore(sqlite, mysql) |
| log.Printf("[INFO] 混合存储已启用(logs 恢复天数: %d)", logDays) |
| return hybrid, nil |
| } |
|
|
| |
| func createMySQLStore(dsn string) (*sqlstore.SQLStore, error) { |
| |
| if dsn == "" { |
| return nil, fmt.Errorf("MySQL DSN不能为空") |
| } |
|
|
| db, err := sql.Open("mysql", dsn) |
| if err != nil { |
| return nil, fmt.Errorf("打开MySQL连接失败: %w", err) |
| } |
|
|
| |
| db.SetMaxOpenConns(config.SQLiteMaxOpenConnsFile * 2) |
| db.SetMaxIdleConns(config.SQLiteMaxIdleConnsFile * 2) |
| db.SetConnMaxLifetime(config.SQLiteConnMaxLifetime) |
|
|
| |
| pingCtx, pingCancel := context.WithTimeout(context.Background(), config.StartupDBPingTimeout) |
| defer pingCancel() |
| if err := db.PingContext(pingCtx); err != nil { |
| _ = db.Close() |
| return nil, fmt.Errorf("MySQL连接测试失败(超时%v): %w", config.StartupDBPingTimeout, err) |
| } |
|
|
| |
| store := sqlstore.NewSQLStore(db, "mysql") |
|
|
| |
| migrateCtx, migrateCancel := context.WithTimeout(context.Background(), config.StartupMigrationTimeout) |
| defer migrateCancel() |
| if err := migrateMySQL(migrateCtx, db); err != nil { |
| _ = db.Close() |
| return nil, fmt.Errorf("MySQL迁移失败(超时%v): %w", config.StartupMigrationTimeout, err) |
| } |
|
|
| return store, nil |
| } |
|
|
| |
| |
| |
| func CreateSQLiteStore(path string) (Store, error) { |
| s, err := createSQLiteStore(path) |
| if err != nil { |
| return nil, err |
| } |
| return s, nil |
| } |
|
|
| |
| |
| |
| func CreateMySQLStoreForTest(dsn string) (Store, error) { |
| s, err := createMySQLStore(dsn) |
| if err != nil { |
| return nil, err |
| } |
| return s, nil |
| } |
|
|
| |
| func createSQLiteStore(path string) (*sqlstore.SQLStore, error) { |
| |
| if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil { |
| return nil, err |
| } |
|
|
| |
| dsn := buildSQLiteDSN(path) |
| db, err := sql.Open("sqlite", dsn) |
| if err != nil { |
| return nil, fmt.Errorf("打开SQLite失败: %w", err) |
| } |
|
|
| |
| |
| |
| |
| |
| db.SetMaxOpenConns(1) |
| db.SetMaxIdleConns(1) |
| db.SetConnMaxLifetime(config.SQLiteConnMaxLifetime) |
|
|
| |
| store := sqlstore.NewSQLStore(db, "sqlite") |
|
|
| |
| migrateCtx, migrateCancel := context.WithTimeout(context.Background(), config.StartupMigrationTimeout) |
| defer migrateCancel() |
| if err := migrateSQLite(migrateCtx, db); err != nil { |
| _ = db.Close() |
| return nil, fmt.Errorf("SQLite迁移失败(超时%v): %w", config.StartupMigrationTimeout, err) |
| } |
|
|
| if _, err := db.ExecContext(migrateCtx, "PRAGMA optimize"); err != nil { |
| log.Printf("[WARN] SQLite PRAGMA optimize 失败: %v", err) |
| } |
|
|
| return store, nil |
| } |
|
|
| |
| |
| func resolveSQLitePath() string { |
| defaultDir := "data" |
| defaultPath := filepath.Join(defaultDir, "ccload.db") |
|
|
| |
| if isDirWritable(defaultDir) { |
| return defaultPath |
| } |
|
|
| |
| if err := os.MkdirAll(defaultDir, 0o750); err == nil { |
| if isDirWritable(defaultDir) { |
| return defaultPath |
| } |
| } |
|
|
| |
| tmpPath := filepath.Join(os.TempDir(), "ccload", "ccload.db") |
| log.Printf("════════════════════════════════════════════════════════════") |
| log.Printf("[WARN] 警告: 默认路径 %s 不可写", defaultDir) |
| log.Printf("[WARN] 数据将存储在临时目录: %s", tmpPath) |
| log.Printf("[WARN] 临时目录数据可能在系统重启后丢失!") |
| log.Printf("[WARN] 生产环境请设置 SQLITE_PATH 环境变量指定持久化路径") |
| log.Printf("════════════════════════════════════════════════════════════") |
| return tmpPath |
| } |
|
|
| |
| func isDirWritable(dir string) bool { |
| info, err := os.Stat(dir) |
| if err != nil { |
| return false |
| } |
| if !info.IsDir() { |
| return false |
| } |
|
|
| |
| testFile := filepath.Join(dir, ".write_test_"+fmt.Sprintf("%d", os.Getpid())) |
| f, err := os.Create(testFile) |
| if err != nil { |
| return false |
| } |
| _ = f.Close() |
| _ = os.Remove(testFile) |
| return true |
| } |
|
|
| |
| func buildSQLiteDSN(path string) string { |
| journalMode := validateJournalMode(os.Getenv("SQLITE_JOURNAL_MODE")) |
| return fmt.Sprintf("file:%s?_pragma=busy_timeout(5000)&_foreign_keys=on&_pragma=journal_mode=%s&_pragma=wal_autocheckpoint(500)&_loc=Local", path, journalMode) |
| } |
|
|
| |
| func validateJournalMode(mode string) string { |
| if mode == "" { |
| return "WAL" |
| } |
|
|
| validModes := map[string]bool{ |
| "DELETE": true, |
| "TRUNCATE": true, |
| "PERSIST": true, |
| "MEMORY": true, |
| "WAL": true, |
| "OFF": true, |
| } |
|
|
| modeUpper := strings.ToUpper(mode) |
| if !validModes[modeUpper] { |
| log.Fatalf("[FATAL] 安全错误: SQLITE_JOURNAL_MODE 环境变量值非法: %q\n"+ |
| " 允许的值: DELETE, TRUNCATE, PERSIST, MEMORY, WAL, OFF\n"+ |
| " 当前值: %q\n"+ |
| " 修复方法:\n"+ |
| " - 设置合法值: export SQLITE_JOURNAL_MODE=WAL\n"+ |
| " - 或者移除该环境变量,使用默认值 WAL", |
| mode, mode) |
| } |
|
|
| return modeUpper |
| } |
|
|
| |
| |
| |
| |
| |
| func getLogSyncDays() int { |
| daysStr := os.Getenv("CCLOAD_SQLITE_LOG_DAYS") |
| if daysStr == "" { |
| return 7 |
| } |
| days, err := strconv.Atoi(daysStr) |
| if err != nil || days < -1 { |
| log.Printf("[WARN] 无效的 CCLOAD_SQLITE_LOG_DAYS=%s,使用默认值 7", daysStr) |
| return 7 |
| } |
| return days |
| } |
|
|