mirror of
https://github.com/Portabase/agent.git
synced 2026-09-10 01:57:10 +00:00
feat: add job logs for restoration process
This commit is contained in:
@@ -19,7 +19,7 @@ pub trait Database: Send + Sync {
|
||||
fn file_extension(&self) -> &'static str;
|
||||
async fn ping(&self) -> Result<bool>;
|
||||
async fn backup(&self, backup_dir: &Path, logger: Arc<JobLogger>) -> Result<PathBuf>;
|
||||
async fn restore(&self, restore_file: &Path) -> Result<()>;
|
||||
async fn restore(&self, restore_file: &Path, logger: Arc<JobLogger>) -> Result<()>;
|
||||
}
|
||||
|
||||
pub struct DatabaseFactory;
|
||||
|
||||
@@ -13,7 +13,7 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting Firebird backup for database: {}", cfg.name));
|
||||
logger.log("info", format!("Starting Firebird backup for database: {}", cfg.name));
|
||||
|
||||
let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension));
|
||||
let db_path = format!("{}/{}:{}", cfg.host, cfg.port, cfg.database);
|
||||
|
||||
@@ -42,9 +42,9 @@ impl Database for FirebirdDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
@@ -1,18 +1,21 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting Firebird restore for database {}", cfg.name);
|
||||
logger.log("debug", format!("Starting Firebird restore for database {}", cfg.name));
|
||||
|
||||
let db_path = format!("{}/{}:{}", cfg.host, cfg.port, cfg.database);
|
||||
|
||||
info!("Restore source: {}", restore_file.display());
|
||||
info!("Restore target: {}", db_path);
|
||||
logger.log("info", format!("Restore source: {}", restore_file.display()));
|
||||
logger.log("info", format!("Restore target: {}", db_path));
|
||||
|
||||
let start = Instant::now();
|
||||
let output = Command::new("gbak")
|
||||
.arg("-c")
|
||||
.arg("-v")
|
||||
@@ -26,13 +29,18 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.output()
|
||||
.with_context(|| format!("Failed to run gbak restore for {}", cfg.name))?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let exit_code = output.status.code().unwrap_or(-1);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
error!("Firebird restore failed for {}: {}", cfg.name, stderr);
|
||||
logger.log_command("gbak", Some(stderr.clone()), Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("Firebird restore failed for {}: {}", cfg.name, stderr));
|
||||
anyhow::bail!("Firebird restore failed for {}: {}", cfg.name, stderr);
|
||||
}
|
||||
|
||||
info!("Firebird restore completed for {}", cfg.name);
|
||||
logger.log_command("gbak", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("Firebird restore completed for {}", cfg.name));
|
||||
|
||||
Ok(())
|
||||
})
|
||||
|
||||
@@ -16,11 +16,11 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting backup for database {}", cfg.name));
|
||||
logger.log("info", format!("Starting backup for database {}", cfg.name));
|
||||
|
||||
let version = match futures::executor::block_on(server_version(&cfg)) {
|
||||
Ok(v) => {
|
||||
logger.log("debug", format!("MariaDB version detected: {}", v));
|
||||
logger.log("info", format!("MariaDB version detected: {}", v));
|
||||
v
|
||||
}
|
||||
Err(e) => {
|
||||
|
||||
@@ -49,9 +49,9 @@ impl Database for MariaDBDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
@@ -1,27 +1,30 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::fs::File;
|
||||
use std::io::{Read, Write};
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
let handle = tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting restore for database {}", cfg.name);
|
||||
logger.log("info", format!("Starting restore for database {}", cfg.name));
|
||||
|
||||
let mut sql_content = String::new();
|
||||
let mut file = File::open(&restore_file)
|
||||
.with_context(|| format!("Failed to open restore file {}", restore_file.display()))?;
|
||||
file.read_to_string(&mut sql_content)
|
||||
.with_context(|| format!("Failed to read restore file {}", restore_file.display()))?;
|
||||
|
||||
|
||||
let drop_create_cmd = format!(
|
||||
"DROP DATABASE IF EXISTS `{0}`; CREATE DATABASE `{0}`;",
|
||||
cfg.database
|
||||
);
|
||||
|
||||
let drop_status = Command::new("mariadb")
|
||||
let drop_start = Instant::now();
|
||||
let drop_output = Command::new("mariadb")
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
@@ -31,14 +34,22 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.arg("-e")
|
||||
.arg(&drop_create_cmd)
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
.status()
|
||||
.output()
|
||||
.with_context(|| format!("Failed to drop/recreate database {}", cfg.name))?;
|
||||
|
||||
if !drop_status.success() {
|
||||
error!("Drop/create database failed for {}", cfg.name);
|
||||
let drop_duration_ms = drop_start.elapsed().as_millis() as f64;
|
||||
let drop_exit_code = drop_output.status.code().unwrap_or(-1);
|
||||
let drop_stderr = String::from_utf8_lossy(&drop_output.stderr).to_string();
|
||||
|
||||
if !drop_output.status.success() {
|
||||
logger.log_command("mariadb", Some(drop_stderr.clone()), Some(drop_exit_code), Some(drop_duration_ms));
|
||||
logger.log("error", format!("Drop/create database failed for {}: {}", cfg.name, drop_stderr));
|
||||
anyhow::bail!("Failed to drop/recreate database {}", cfg.name);
|
||||
}
|
||||
info!("Database {} dropped and recreated", cfg.name);
|
||||
logger.log_command("mariadb", if drop_stderr.is_empty() { None } else { Some(drop_stderr) }, Some(0), Some(drop_duration_ms));
|
||||
logger.log("info", format!("Database {} dropped and recreated", cfg.name));
|
||||
|
||||
let start = Instant::now();
|
||||
|
||||
let mut child = Command::new("mariadb")
|
||||
.arg("--host")
|
||||
@@ -65,13 +76,18 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.wait_with_output()
|
||||
.with_context(|| format!("Failed to complete MariaDB restore for {}", cfg.name))?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let exit_code = output.status.code().unwrap_or(-1);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
error!("MariaDB restore failed for {}: {}", cfg.name, stderr);
|
||||
logger.log_command("mariadb", Some(stderr.clone()), Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("MariaDB restore failed for {}: {}", cfg.name, stderr));
|
||||
anyhow::bail!("MariaDB restore failed for {}", cfg.name);
|
||||
}
|
||||
|
||||
info!("Restore finished successfully for database {}", cfg.name);
|
||||
logger.log_command("mariadb", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("Restore finished successfully for database {}", cfg.name));
|
||||
Ok(())
|
||||
});
|
||||
|
||||
|
||||
@@ -6,7 +6,6 @@ use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
use tracing::error;
|
||||
|
||||
pub async fn run(
|
||||
cfg: DatabaseConfig,
|
||||
@@ -15,7 +14,7 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting MongoDB backup for database {}", cfg.name));
|
||||
logger.log("info", format!("Starting MongoDB backup for database {}", cfg.name));
|
||||
|
||||
let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension));
|
||||
let mongodump = select_mongo_path().join("mongodump");
|
||||
@@ -37,7 +36,7 @@ pub async fn run(
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
|
||||
if !output.status.success() {
|
||||
error!("MongoDB backup failed for {}: {}", cfg.name, stderr);
|
||||
logger.log("error", format!("MongoDB backup failed for {}: {}", cfg.name, stderr));
|
||||
logger.log_command("mongodump", Some(stderr.clone()), Some(exit_code), Some(duration_ms));
|
||||
anyhow::bail!("MongoDB backup failed for {}: {}", cfg.name, stderr);
|
||||
}
|
||||
|
||||
@@ -36,9 +36,9 @@ impl Database for MongoDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
@@ -1,17 +1,20 @@
|
||||
use crate::domain::mongodb::connection::{extract_db_name, get_mongo_uri, select_mongo_path};
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting MongoDB restore for database {}", cfg.name);
|
||||
logger.log("debug", format!("Starting MongoDB restore for database {}", cfg.name));
|
||||
|
||||
let mongorestore = select_mongo_path().join("mongorestore");
|
||||
let uri = get_mongo_uri(cfg.clone())?;
|
||||
|
||||
let dry_start = Instant::now();
|
||||
let dry_run = Command::new(&mongorestore)
|
||||
.arg(format!(
|
||||
"--uri={}",
|
||||
@@ -26,14 +29,23 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.arg("--verbose")
|
||||
.output()?;
|
||||
|
||||
let dry_duration_ms = dry_start.elapsed().as_millis() as f64;
|
||||
let dry_exit_code = dry_run.status.code().unwrap_or(-1);
|
||||
let dry_output = String::from_utf8_lossy(&dry_run.stderr);
|
||||
logger.log_command(
|
||||
"mongorestore --dryRun",
|
||||
if dry_output.is_empty() { None } else { Some(dry_output.to_string()) },
|
||||
Some(dry_exit_code),
|
||||
Some(dry_duration_ms),
|
||||
);
|
||||
let source_db = extract_db_name(&dry_output).unwrap_or_else(|| {
|
||||
info!("Could not detect source database from archive, falling back to configured database: {}", cfg.database);
|
||||
logger.log("info", format!("Could not detect source database from archive, falling back to configured database: {}", cfg.database));
|
||||
cfg.database.clone()
|
||||
});
|
||||
|
||||
info!("Using source database in archive: {}", source_db);
|
||||
logger.log("info", format!("Using source database in archive: {}", source_db));
|
||||
|
||||
let start = Instant::now();
|
||||
let output = Command::new(&mongorestore)
|
||||
.arg(format!("--uri={}", uri))
|
||||
.arg(format!("--archive={}", restore_file.display()))
|
||||
@@ -45,13 +57,18 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.output()
|
||||
.with_context(|| format!("Failed to run mongorestore for {}", cfg.name))?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let exit_code = output.status.code().unwrap_or(-1);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
error!("MongoDB restore failed for {}: {}", cfg.name, stderr);
|
||||
logger.log_command("mongorestore", Some(stderr.clone()), Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("MongoDB restore failed for {}: {}", cfg.name, stderr));
|
||||
anyhow::bail!("MongoDB restore failed for: {}", cfg.name);
|
||||
}
|
||||
|
||||
info!("MongoDB restore completed for {}", cfg.name);
|
||||
logger.log_command("mongorestore", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("MongoDB restore completed for {}", cfg.name));
|
||||
Ok(())
|
||||
})
|
||||
.await?
|
||||
|
||||
@@ -13,7 +13,7 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting MSSQL backup for database {}", cfg.name));
|
||||
logger.log("info", format!("Starting MSSQL backup for database {}", cfg.name));
|
||||
|
||||
let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension));
|
||||
let connection_string = format!(
|
||||
|
||||
@@ -41,9 +41,9 @@ impl Database for MssqlDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
+19
-10
@@ -1,26 +1,29 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting MSSQL restore for database {}", cfg.name);
|
||||
logger.log("debug", format!("Starting MSSQL restore for database {}", cfg.name));
|
||||
|
||||
let connection_string = format!(
|
||||
"Server=tcp:{},{};Database={};User Id={};Password={};TrustServerCertificate=True;Encrypt=False",
|
||||
cfg.host, cfg.port, cfg.database, cfg.username, cfg.password
|
||||
);
|
||||
|
||||
info!(
|
||||
logger.log("info", format!(
|
||||
"MSSQL restore: {} → {}:{}/{}",
|
||||
restore_file.display(),
|
||||
cfg.host,
|
||||
cfg.port,
|
||||
cfg.database
|
||||
);
|
||||
));
|
||||
|
||||
let start = Instant::now();
|
||||
let output = Command::new("sqlpackage")
|
||||
.arg("/a:Import")
|
||||
.arg(format!("/tcs:{}", connection_string))
|
||||
@@ -28,17 +31,23 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.output()
|
||||
.with_context(|| format!("Failed to run sqlpackage restore for {}", cfg.name))?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let exit_code = output.status.code().unwrap_or(-1);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
let stdout = String::from_utf8_lossy(&output.stdout).to_string();
|
||||
let combined = format!("{}{}", stdout, stderr);
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
let stdout = String::from_utf8_lossy(&output.stdout);
|
||||
error!(
|
||||
logger.log_command("sqlpackage", if combined.is_empty() { None } else { Some(combined.clone()) }, Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!(
|
||||
"MSSQL restore failed for {} — stderr: {} stdout: {}",
|
||||
cfg.name, stderr, stdout
|
||||
);
|
||||
));
|
||||
anyhow::bail!("MSSQL restore failed for {}: {}", cfg.name, stderr);
|
||||
}
|
||||
|
||||
info!("MSSQL restore completed for {}", cfg.name);
|
||||
logger.log_command("sqlpackage", if combined.is_empty() { None } else { Some(combined) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("MSSQL restore completed for {}", cfg.name));
|
||||
Ok(())
|
||||
})
|
||||
.await?
|
||||
|
||||
@@ -16,7 +16,7 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting backup for database {}", cfg.name));
|
||||
logger.log("info", format!("Starting backup for database {}", cfg.name));
|
||||
|
||||
let _version = match futures::executor::block_on(server_version(&cfg)) {
|
||||
Ok(v) => {
|
||||
|
||||
@@ -49,9 +49,9 @@ impl Database for MySQLDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
+27
-11
@@ -1,14 +1,16 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::fs::File;
|
||||
use std::io::{Read, Write};
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
let handle = tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting restore for database {}", cfg.name);
|
||||
logger.log("info", format!("Starting restore for database {}", cfg.name));
|
||||
|
||||
let mut sql_content = String::new();
|
||||
let mut file = File::open(&restore_file)
|
||||
@@ -21,7 +23,8 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
cfg.database
|
||||
);
|
||||
|
||||
let drop_status = Command::new("mysql")
|
||||
let drop_start = Instant::now();
|
||||
let drop_output = Command::new("mysql")
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
@@ -31,14 +34,22 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.arg("-e")
|
||||
.arg(&drop_create_cmd)
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
.status()
|
||||
.output()
|
||||
.with_context(|| format!("Failed to drop/recreate database {}", cfg.name))?;
|
||||
|
||||
if !drop_status.success() {
|
||||
error!("Drop/create database failed for {}", cfg.name);
|
||||
let drop_duration_ms = drop_start.elapsed().as_millis() as f64;
|
||||
let drop_exit_code = drop_output.status.code().unwrap_or(-1);
|
||||
let drop_stderr = String::from_utf8_lossy(&drop_output.stderr).to_string();
|
||||
|
||||
if !drop_output.status.success() {
|
||||
logger.log_command("mysql", Some(drop_stderr.clone()), Some(drop_exit_code), Some(drop_duration_ms));
|
||||
logger.log("error", format!("Drop/create database failed for {}: {}", cfg.name, drop_stderr));
|
||||
anyhow::bail!("Failed to drop/recreate database {}", cfg.name);
|
||||
}
|
||||
info!("Database {} dropped and recreated", cfg.name);
|
||||
logger.log_command("mysql", if drop_stderr.is_empty() { None } else { Some(drop_stderr) }, Some(0), Some(drop_duration_ms));
|
||||
logger.log("info", format!("Database {} dropped and recreated", cfg.name));
|
||||
|
||||
let start = Instant::now();
|
||||
|
||||
let mut child = Command::new("mysql")
|
||||
.arg("--host")
|
||||
@@ -65,13 +76,18 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.wait_with_output()
|
||||
.with_context(|| format!("Failed to complete mysql restore for {}", cfg.name))?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let exit_code = output.status.code().unwrap_or(-1);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
error!("MySQL restore failed for {}: {}", cfg.name, stderr);
|
||||
logger.log_command("mysql", Some(stderr.clone()), Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("MySQL restore failed for {}: {}", cfg.name, stderr));
|
||||
anyhow::bail!("MySQL restore failed for {}", cfg.name);
|
||||
}
|
||||
|
||||
info!("Restore finished successfully for database {}", cfg.name);
|
||||
logger.log_command("mysql", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("Restore finished successfully for database {}", cfg.name));
|
||||
Ok(())
|
||||
});
|
||||
|
||||
|
||||
@@ -16,7 +16,7 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting backup for database {}", cfg.name));
|
||||
logger.log("info", format!("Starting backup for database {}", cfg.name));
|
||||
|
||||
let version = match futures::executor::block_on(server_version(&cfg)) {
|
||||
Ok(v) => {
|
||||
|
||||
@@ -40,9 +40,9 @@ impl Database for PostgresDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), self.format, file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), self.format, file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
@@ -1,52 +1,54 @@
|
||||
use anyhow::Result;
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
use super::connection::{select_pg_path, server_version, terminate_connections};
|
||||
use super::format::PostgresDumpFormat;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
|
||||
pub async fn run(
|
||||
cfg: DatabaseConfig,
|
||||
format: PostgresDumpFormat,
|
||||
restore_file: PathBuf,
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<()> {
|
||||
tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting restore for database {}", cfg.name);
|
||||
logger.log("info", format!("Starting restore for database {}", cfg.name));
|
||||
|
||||
let version = match futures::executor::block_on(server_version(&cfg)) {
|
||||
Ok(v) => {
|
||||
debug!("Postgres version detected: {}", v);
|
||||
logger.log("debug", format!("Postgres version detected: {}", v));
|
||||
v
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Failed to get server version for {}: {:?}", cfg.name, e);
|
||||
logger.log("error", format!("Failed to get server version for {}: {:?}", cfg.name, e));
|
||||
return Err(e.into());
|
||||
}
|
||||
};
|
||||
|
||||
let pg_restore = select_pg_path(&version).join("pg_restore");
|
||||
|
||||
debug!("Using pg_restore at {:?}", pg_restore);
|
||||
logger.log("debug", format!("Using pg_restore at {:?}", pg_restore));
|
||||
|
||||
if let Err(e) = futures::executor::block_on(terminate_connections(&cfg)) {
|
||||
error!("Failed to terminate connections for {}: {:?}", cfg.name, e);
|
||||
logger.log("error", format!("Failed to terminate connections for {}: {:?}", cfg.name, e));
|
||||
return Err(e.into());
|
||||
}
|
||||
info!("Connections terminated for database {}", cfg.name);
|
||||
logger.log("info", format!("Connections terminated for database {}", cfg.name));
|
||||
|
||||
let url = format!(
|
||||
"postgresql://{}:{}@{}:{}/{}",
|
||||
cfg.username, cfg.password, cfg.host, cfg.port, cfg.database
|
||||
);
|
||||
|
||||
debug!("Restore URL: {}", url);
|
||||
|
||||
match format {
|
||||
PostgresDumpFormat::Fc => {
|
||||
info!("Running FC restore for {}", cfg.name);
|
||||
let status = Command::new(&pg_restore)
|
||||
logger.log("info", format!("Running FC restore for {}", cfg.name));
|
||||
let start = Instant::now();
|
||||
let output = Command::new(&pg_restore)
|
||||
.arg("--no-owner")
|
||||
.arg("--no-privileges")
|
||||
.arg("--clean")
|
||||
@@ -57,38 +59,49 @@ pub async fn run(
|
||||
.arg("-v")
|
||||
.arg(&restore_file)
|
||||
.env("PGPASSWORD", &cfg.password)
|
||||
.status();
|
||||
.output();
|
||||
|
||||
match status {
|
||||
Ok(s) if s.success() => {
|
||||
info!("FC restore completed successfully for {}", cfg.name)
|
||||
}
|
||||
Ok(s) => {
|
||||
error!("FC restore failed with status {:?} for {}", s, cfg.name);
|
||||
anyhow::bail!("Postgres restore failed for {}", cfg.name);
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
|
||||
match output {
|
||||
Ok(o) => {
|
||||
let stderr = String::from_utf8_lossy(&o.stderr).to_string();
|
||||
let stdout = String::from_utf8_lossy(&o.stdout).to_string();
|
||||
let combined = format!("{}{}", stdout, stderr);
|
||||
let exit_code = o.status.code().unwrap_or(-1);
|
||||
|
||||
if o.status.success() {
|
||||
logger.log_command("pg_restore", if combined.is_empty() { None } else { Some(combined) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("FC restore completed successfully for {}", cfg.name))
|
||||
} else {
|
||||
logger.log_command("pg_restore", if combined.is_empty() { None } else { Some(combined) }, Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("FC restore failed with status {:?} for {}", o.status, cfg.name));
|
||||
anyhow::bail!("Postgres restore failed for {}", cfg.name);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Error executing pg_restore for {}: {:?}", cfg.name, e);
|
||||
logger.log_command("pg_restore", Some(e.to_string()), Some(-1), Some(duration_ms));
|
||||
logger.log("error", format!("Error executing pg_restore for {}: {:?}", cfg.name, e));
|
||||
return Err(e.into());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
PostgresDumpFormat::Fd => {
|
||||
info!("Running FD restore for {}", cfg.name);
|
||||
logger.log("info", format!("Running FD restore for {}", cfg.name));
|
||||
|
||||
let tar_gz = match std::fs::File::open(&restore_file) {
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
error!(
|
||||
logger.log("error", format!(
|
||||
"Failed to open restore file {:?} for {}: {:?}",
|
||||
restore_file, cfg.name, e
|
||||
);
|
||||
));
|
||||
return Err(e.into());
|
||||
}
|
||||
};
|
||||
|
||||
info!("tar_gz {:?}", tar_gz);
|
||||
logger.log("info", format!("tar_gz {:?}", tar_gz));
|
||||
|
||||
let dec = flate2::read::GzDecoder::new(tar_gz);
|
||||
let mut archive = tar::Archive::new(dec);
|
||||
@@ -96,31 +109,31 @@ pub async fn run(
|
||||
let tmp_dir = match tempfile::TempDir::new() {
|
||||
Ok(d) => d,
|
||||
Err(e) => {
|
||||
error!(
|
||||
logger.log("error", format!(
|
||||
"Failed to create temporary directory for FD restore of {}: {:?}",
|
||||
cfg.name, e
|
||||
);
|
||||
));
|
||||
return Err(e.into());
|
||||
}
|
||||
};
|
||||
|
||||
if let Err(e) = archive.unpack(tmp_dir.path()) {
|
||||
error!("Failed to unpack FD archive for {}: {:?}", cfg.name, e);
|
||||
logger.log("error", format!("Failed to unpack FD archive for {}: {:?}", cfg.name, e));
|
||||
return Err(e.into());
|
||||
}
|
||||
|
||||
debug!("Listing contents of temp dir: {}", tmp_dir.path().display());
|
||||
logger.log("debug", format!("Listing contents of temp dir: {}", tmp_dir.path().display()));
|
||||
|
||||
for entry in std::fs::read_dir(tmp_dir.path())? {
|
||||
if let Ok(entry) = entry {
|
||||
let path = entry.path();
|
||||
let file_type = entry.file_type()?;
|
||||
debug!(
|
||||
logger.log("debug", format!(
|
||||
" - {} | is_dir: {} | is_file: {}",
|
||||
path.display(),
|
||||
file_type.is_dir(),
|
||||
file_type.is_file()
|
||||
);
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -134,7 +147,8 @@ pub async fn run(
|
||||
.ok_or_else(|| anyhow::anyhow!("Invalid FD archive: toc.dat not found"))?
|
||||
};
|
||||
|
||||
let status = Command::new(&pg_restore)
|
||||
let start = Instant::now();
|
||||
let output = Command::new(&pg_restore)
|
||||
.arg("--no-owner")
|
||||
.arg("--no-privileges")
|
||||
.arg("--clean")
|
||||
@@ -147,25 +161,36 @@ pub async fn run(
|
||||
.arg("4")
|
||||
.arg(dump_dir)
|
||||
.env("PGPASSWORD", &cfg.password)
|
||||
.status();
|
||||
.output();
|
||||
|
||||
match status {
|
||||
Ok(s) if s.success() => {
|
||||
info!("FD restore completed successfully for {}", cfg.name)
|
||||
}
|
||||
Ok(s) => {
|
||||
error!("FD restore failed with status {:?} for {}", s, cfg.name);
|
||||
anyhow::bail!("Postgres FD restore failed for {}", cfg.name);
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
|
||||
match output {
|
||||
Ok(o) => {
|
||||
let stderr = String::from_utf8_lossy(&o.stderr).to_string();
|
||||
let stdout = String::from_utf8_lossy(&o.stdout).to_string();
|
||||
let combined = format!("{}{}", stdout, stderr);
|
||||
let exit_code = o.status.code().unwrap_or(-1);
|
||||
|
||||
if o.status.success() {
|
||||
logger.log_command("pg_restore", if combined.is_empty() { None } else { Some(combined) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("FD restore completed successfully for {}", cfg.name))
|
||||
} else {
|
||||
logger.log_command("pg_restore", if combined.is_empty() { None } else { Some(combined) }, Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("FD restore failed with status {:?} for {}", o.status, cfg.name));
|
||||
anyhow::bail!("Postgres FD restore failed for {}", cfg.name);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Error executing pg_restore for {}: {:?}", cfg.name, e);
|
||||
logger.log_command("pg_restore", Some(e.to_string()), Some(-1), Some(duration_ms));
|
||||
logger.log("error", format!("Error executing pg_restore for {}: {:?}", cfg.name, e));
|
||||
return Err(e.into());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
info!("Restore finished for database {}", cfg.name);
|
||||
logger.log("info", format!("Restore finished for database {}", cfg.name));
|
||||
|
||||
Ok(())
|
||||
})
|
||||
|
||||
@@ -14,7 +14,7 @@ pub async fn run(
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log(
|
||||
"debug",
|
||||
"info",
|
||||
format!("Starting Redis backup for database {}", cfg.name),
|
||||
);
|
||||
|
||||
|
||||
@@ -36,7 +36,7 @@ impl Database for RedisDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, _file: &Path) -> Result<()> {
|
||||
async fn restore(&self, _file: &Path, _logger: Arc<JobLogger>) -> Result<()> {
|
||||
bail!("Restore not supported for Redis databases")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@ pub async fn run(
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log("debug", format!("Starting SQLite backup for database {}", cfg.name));
|
||||
logger.log("info", format!("Starting SQLite backup for database {}", cfg.name));
|
||||
|
||||
let db_path_str = if cfg.path.is_empty() {
|
||||
anyhow::bail!("Database path not configured");
|
||||
|
||||
@@ -35,9 +35,9 @@ impl Database for SqliteDatabase {
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
async fn restore(&self, file: &Path) -> Result<()> {
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf()).await;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
@@ -1,12 +1,14 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::path::PathBuf;
|
||||
use std::process::Command;
|
||||
use tracing::{debug, error, info};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
debug!("Starting SQLite restore for database {}", cfg.name);
|
||||
logger.log("debug", format!("Starting SQLite restore for database {}", cfg.name));
|
||||
|
||||
let db_path_str = if cfg.path.is_empty() {
|
||||
anyhow::bail!("Database path not configured");
|
||||
@@ -25,19 +27,25 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> {
|
||||
.with_context(|| format!("Failed to remove existing DB {}", db_path.display()))?;
|
||||
}
|
||||
|
||||
let start = Instant::now();
|
||||
let output = Command::new("sqlite3")
|
||||
.arg(db_path.as_os_str())
|
||||
.arg(format!(".restore '{}'", restore_file.display()))
|
||||
.output()
|
||||
.with_context(|| format!("Failed to run sqlite3 restore for {}", cfg.name))?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let exit_code = output.status.code().unwrap_or(-1);
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
error!("SQLite restore failed for {}: {}", cfg.name, stderr);
|
||||
logger.log_command("sqlite3", Some(stderr.clone()), Some(exit_code), Some(duration_ms));
|
||||
logger.log("error", format!("SQLite restore failed for {}: {}", cfg.name, stderr));
|
||||
anyhow::bail!("SQLite restore failed for {}", cfg.name);
|
||||
}
|
||||
|
||||
info!("SQLite restore completed for {}", cfg.name);
|
||||
logger.log_command("sqlite3", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("SQLite restore completed for {}", cfg.name));
|
||||
Ok(())
|
||||
})
|
||||
.await?
|
||||
|
||||
@@ -14,7 +14,7 @@ pub async fn run(
|
||||
) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
logger.log(
|
||||
"debug",
|
||||
"info",
|
||||
format!("Starting Valkey backup for database {}", cfg.name),
|
||||
);
|
||||
|
||||
|
||||
@@ -36,7 +36,7 @@ impl Database for ValkeyDatabase {
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, _file: &Path) -> Result<()> {
|
||||
async fn restore(&self, _file: &Path, _logger: Arc<JobLogger>) -> Result<()> {
|
||||
bail!("Restore not supported for Valkey databases")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,12 +3,16 @@ use crate::services::api::{ApiClient, ApiError};
|
||||
use anyhow::Result;
|
||||
use reqwest::Method;
|
||||
use serde::Serialize;
|
||||
use crate::services::backup::logger::JobLogEntry;
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct ResultRestoreRequest {
|
||||
#[serde(rename = "generatedId")]
|
||||
pub generated_id: String,
|
||||
pub status: String,
|
||||
pub logs: Vec<JobLogEntry>,
|
||||
#[serde(rename = "durationMs")]
|
||||
pub duration_ms: f64,
|
||||
}
|
||||
|
||||
impl ApiClient {
|
||||
@@ -17,10 +21,14 @@ impl ApiClient {
|
||||
agent_id: impl Into<String>,
|
||||
generated_id: impl Into<String>,
|
||||
status: impl Into<String>,
|
||||
job_logs: Vec<JobLogEntry>,
|
||||
duration_ms: f64,
|
||||
) -> Result<Option<ResultRestoreResponse>, ApiError> {
|
||||
let body = ResultRestoreRequest {
|
||||
generated_id: generated_id.into(),
|
||||
status: status.into(),
|
||||
logs: job_logs,
|
||||
duration_ms,
|
||||
};
|
||||
|
||||
let agent_id = agent_id.into();
|
||||
|
||||
@@ -5,22 +5,30 @@ use crate::utils::file::decrypt_file_stream_gcm;
|
||||
|
||||
use anyhow::Result;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn prepare_archive(
|
||||
&self,
|
||||
downloaded_file: PathBuf,
|
||||
tmp_path: &Path,
|
||||
logger: Arc<JobLogger>
|
||||
) -> Result<PathBuf> {
|
||||
logger.log("info", "Start preparing backup archive".to_string());
|
||||
|
||||
let filename = downloaded_file
|
||||
.file_name()
|
||||
.unwrap()
|
||||
.to_string_lossy()
|
||||
.to_string();
|
||||
|
||||
logger.log("debug", format!("Archive filename: {}", filename));
|
||||
|
||||
let is_legacy = filename.ends_with(".sql") || filename.ends_with(".dump");
|
||||
|
||||
if is_legacy {
|
||||
logger.log("info", "Legacy archive detected, skipping extraction".to_string());
|
||||
return Ok(downloaded_file);
|
||||
}
|
||||
|
||||
@@ -29,29 +37,50 @@ impl RestoreService {
|
||||
let mut archive = downloaded_file.clone();
|
||||
|
||||
if encrypted {
|
||||
logger.log("info", "Archive is encrypted, decrypting".to_string());
|
||||
|
||||
let new_name = filename.strip_suffix(".enc").unwrap();
|
||||
|
||||
let decrypted = tmp_path.join(new_name);
|
||||
|
||||
decrypt_file_stream_gcm(
|
||||
if let Err(e) = decrypt_file_stream_gcm(
|
||||
downloaded_file,
|
||||
decrypted.clone(),
|
||||
self.ctx.edge_key.master_key_b64.clone(),
|
||||
)
|
||||
.await?;
|
||||
.await
|
||||
{
|
||||
logger.log("error", format!("Failed to decrypt archive: {}", e));
|
||||
return Err(e);
|
||||
}
|
||||
|
||||
logger.log("info", format!("Archive decrypted to {}", decrypted.display()));
|
||||
|
||||
archive = decrypted;
|
||||
}
|
||||
|
||||
let files = decompress_large_tar_gz(archive.as_path(), tmp_path).await?;
|
||||
logger.log("info", format!("Decompressing archive {}", archive.display()));
|
||||
|
||||
let files = match decompress_large_tar_gz(archive.as_path(), tmp_path).await {
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
logger.log("error", format!("Failed to decompress archive: {}", e));
|
||||
return Err(e);
|
||||
}
|
||||
};
|
||||
|
||||
if files.is_empty() {
|
||||
logger.log("error", "Archive is empty after decompression".to_string());
|
||||
anyhow::bail!("archive empty");
|
||||
}
|
||||
|
||||
logger.log("info", format!("Archive prepared, {} file(s) extracted", files.len()));
|
||||
|
||||
if files.len() == 1 {
|
||||
logger.log("debug", format!("Using single extracted file: {}", files[0].display()));
|
||||
Ok(files[0].clone())
|
||||
} else {
|
||||
logger.log("debug", format!("Multiple files extracted, using archive root: {}", archive.display()));
|
||||
Ok(archive)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,15 +3,19 @@ use super::service::RestoreService;
|
||||
use anyhow::Result;
|
||||
use reqwest::{Client, Url};
|
||||
use std::path::{Path, PathBuf};
|
||||
use tracing::info;
|
||||
use std::sync::Arc;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn download_backup(&self, file_url: &str, tmp_path: &Path) -> Result<PathBuf> {
|
||||
pub async fn download_backup(&self, file_url: &str, tmp_path: &Path, logger: Arc<JobLogger>) -> Result<PathBuf> {
|
||||
logger.log("info", "Start downloading backup archive".to_string());
|
||||
|
||||
let client = Client::new();
|
||||
|
||||
let response = client.get(file_url).send().await?;
|
||||
|
||||
if !response.status().is_success() {
|
||||
logger.log("error", "Failed to download".to_string());
|
||||
anyhow::bail!("download failed");
|
||||
}
|
||||
|
||||
@@ -39,8 +43,7 @@ impl RestoreService {
|
||||
|
||||
tokio::fs::write(&path, &bytes).await?;
|
||||
|
||||
info!("Backup downloaded to {}", path.display());
|
||||
|
||||
logger.log("info", format!("Backup downloaded to {}", path.display()));
|
||||
Ok(path)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,23 +1,36 @@
|
||||
use super::service::RestoreService;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::Result;
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
use tempfile::TempDir;
|
||||
use tracing::info;
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn execute_restore(&self, cfg: DatabaseConfig, file_url: String) -> Result<()> {
|
||||
let logger = Arc::new(JobLogger::new());
|
||||
let start = Instant::now();
|
||||
|
||||
logger.log("info", "Database restoration job started".to_string());
|
||||
|
||||
let temp_dir = TempDir::new()?;
|
||||
let tmp_path = temp_dir.path();
|
||||
|
||||
info!("Created temp directory {}", tmp_path.display());
|
||||
logger.log("info", format!("Created temp directory {}", tmp_path.display()));
|
||||
|
||||
let downloaded = self.download_backup(&file_url, tmp_path).await?;
|
||||
let downloaded = self.download_backup(&file_url, tmp_path, Arc::clone(&logger)).await?;
|
||||
|
||||
let backup_file = self.prepare_archive(downloaded, tmp_path).await?;
|
||||
let backup_file = self.prepare_archive(downloaded, tmp_path, Arc::clone(&logger)).await?;
|
||||
|
||||
let result = self.run_restore(cfg, backup_file).await?;
|
||||
let result = self.run_restore(cfg, backup_file, Arc::clone(&logger)).await?;
|
||||
|
||||
self.send_result(result).await;
|
||||
logger.log("info", "Database restore job finished".to_string());
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
let logs = Arc::try_unwrap(logger)
|
||||
.unwrap_or_else(|_| JobLogger::new())
|
||||
.into_entries();
|
||||
self.send_result(result, logs, duration_ms).await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1,31 +1,38 @@
|
||||
use super::models::RestoreResult;
|
||||
use super::service::RestoreService;
|
||||
use crate::services::api::ApiError;
|
||||
use crate::services::api::models::agent::restore::ResultRestoreResponse;
|
||||
|
||||
use tracing::{error, info};
|
||||
use crate::services::backup::logger::JobLogEntry;
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn send_result(&self, result: RestoreResult) {
|
||||
pub async fn send_result(
|
||||
&self,
|
||||
result: RestoreResult,
|
||||
logs: Vec<JobLogEntry>,
|
||||
duration_ms: f64,
|
||||
) -> Result<Option<ResultRestoreResponse>, ApiError> {
|
||||
info!(
|
||||
"[RestoreService] DB: {} | Status: {}",
|
||||
result.generated_id, result.status
|
||||
"[RestoreService] DB: {} | Status: {} | Duration: {}ms",
|
||||
result.generated_id,
|
||||
result.status,
|
||||
duration_ms
|
||||
);
|
||||
|
||||
match self
|
||||
.ctx
|
||||
self.ctx
|
||||
.api
|
||||
.restore_result(
|
||||
self.ctx.edge_key.agent_id.clone(),
|
||||
&result.generated_id,
|
||||
&result.status,
|
||||
logs,
|
||||
duration_ms
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
info!("Restoration result sent successfully");
|
||||
}
|
||||
Err(e) => {
|
||||
.map_err(|e| {
|
||||
error!("Failed to send restoration result: {}", e);
|
||||
}
|
||||
}
|
||||
e.into()
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -4,39 +4,65 @@ use super::service::RestoreService;
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use anyhow::Result;
|
||||
use std::path::PathBuf;
|
||||
use tracing::{error, info};
|
||||
use std::sync::Arc;
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn run_restore(
|
||||
&self,
|
||||
cfg: DatabaseConfig,
|
||||
backup_file: PathBuf,
|
||||
logger: Arc<JobLogger>,
|
||||
) -> Result<RestoreResult> {
|
||||
let generated_id = cfg.generated_id.clone();
|
||||
|
||||
logger.log(
|
||||
"info",
|
||||
format!("Preparing restore for database {}", cfg.name),
|
||||
);
|
||||
|
||||
let db = DatabaseFactory::create_for_restore(cfg.clone(), &backup_file).await;
|
||||
|
||||
logger.log(
|
||||
"debug",
|
||||
format!("Checking reachability for database {}", cfg.name),
|
||||
);
|
||||
|
||||
let reachable = db.ping().await.unwrap_or(false);
|
||||
|
||||
info!("Reachable: {}", reachable);
|
||||
logger.log("info", format!("Reachable: {}", reachable));
|
||||
|
||||
if !reachable {
|
||||
logger.log(
|
||||
"error",
|
||||
format!("Database {} unreachable, aborting restore", cfg.name),
|
||||
);
|
||||
|
||||
return Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
});
|
||||
}
|
||||
|
||||
match db.restore(&backup_file).await {
|
||||
Ok(_) => Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "success".into(),
|
||||
}),
|
||||
match db.restore(&backup_file, Arc::clone(&logger)).await {
|
||||
Ok(_) => {
|
||||
logger.log(
|
||||
"info",
|
||||
format!("Restore completed successfully for database {}", cfg.name),
|
||||
);
|
||||
Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "success".into(),
|
||||
})
|
||||
}
|
||||
|
||||
Err(e) => {
|
||||
error!("Restore failed: {:?}", e);
|
||||
logger.log(
|
||||
"error",
|
||||
format!("Restore failed for database {}: {:?}", cfg.name, e),
|
||||
);
|
||||
|
||||
Ok(RestoreResult {
|
||||
generated_id,
|
||||
|
||||
@@ -95,7 +95,7 @@ async fn firebird_backup_restore_test() {
|
||||
info!("Reachable: {}", reachable);
|
||||
assert!(reachable);
|
||||
|
||||
match db.restore(&backup_file).await {
|
||||
match db.restore(&backup_file, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await {
|
||||
Ok(_) => {
|
||||
info!("Restore succeeded for {}", config.generated_id);
|
||||
assert!(true)
|
||||
|
||||
@@ -81,7 +81,7 @@ async fn mariadb_backup_restore_test() {
|
||||
info!("Reachable: {}", reachable);
|
||||
assert!(reachable);
|
||||
|
||||
match db.restore(&backup_file).await {
|
||||
match db.restore(&backup_file, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await {
|
||||
Ok(_) => {
|
||||
info!("Restore succeeded for {}", config.generated_id);
|
||||
assert!(true)
|
||||
|
||||
@@ -85,7 +85,7 @@ async fn mongodb_backup_restore_test() {
|
||||
info!("Reachable: {}", reachable);
|
||||
assert!(reachable);
|
||||
|
||||
match db.restore(&file_path).await {
|
||||
match db.restore(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await {
|
||||
Ok(_) => {
|
||||
info!("Restore succeeded for {}", config.generated_id);
|
||||
assert!(true)
|
||||
|
||||
@@ -141,7 +141,7 @@ async fn mssql_backup_restore_test() {
|
||||
let restore_config = make_config(host, port, "restoreddb", "5a445eb4-c2c6-4bde-a423-ee1385dcf6d5");
|
||||
let db_restore = DatabaseFactory::create_for_restore(restore_config, &backup_file).await;
|
||||
|
||||
match db_restore.restore(&backup_file).await {
|
||||
match db_restore.restore(&backup_file, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await {
|
||||
Ok(_) => info!("MSSQL restore succeeded"),
|
||||
Err(e) => {
|
||||
error!("MSSQL restore failed: {:?}", e);
|
||||
|
||||
@@ -81,7 +81,7 @@ async fn mysql_backup_restore_test() {
|
||||
info!("Reachable: {}", reachable);
|
||||
assert!(reachable);
|
||||
|
||||
match db.restore(&backup_file).await {
|
||||
match db.restore(&backup_file, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await {
|
||||
Ok(_) => {
|
||||
info!("Restore succeeded for {}", config.generated_id);
|
||||
assert!(true)
|
||||
|
||||
@@ -96,7 +96,7 @@ async fn postgres_backup_restore_test() {
|
||||
|
||||
info!("Running pg_restore: {:?}", backup_file);
|
||||
|
||||
match db.restore(&backup_file).await {
|
||||
match db.restore(&backup_file, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await {
|
||||
Ok(_) => {
|
||||
info!("Restore succeeded for {}", config.generated_id);
|
||||
assert!(true)
|
||||
|
||||
Reference in New Issue
Block a user