diff --git a/.github/workflows/codecov.yml b/.github/workflows/codecov.yml index 07efe04..bb044bc 100644 --- a/.github/workflows/codecov.yml +++ b/.github/workflows/codecov.yml @@ -1,86 +1,3 @@ -#name: Codecov Rust -# -#on: -# push: -# branches: ["main"] -# pull_request: -# branches: ["main"] -# -#env: -# CARGO_TERM_COLOR: always -# -#jobs: -# coverage: -# runs-on: ubuntu-latest -# steps: -# - uses: actions/checkout@v4 -# -# - name: Install Rust toolchain -# uses: dtolnay/rust-toolchain@stable -# with: -# components: llvm-tools-preview -# -# - name: Build test image (with grcov included) -# run: docker compose -f docker-compose.test.yml build agent-test -# -# - name: Run tests in container -# env: -# CARGO_TARGET_DIR: /app/target -# CARGO_INCREMENTAL: 0 -# RUSTFLAGS: "-C instrument-coverage -C link-dead-code" -# LLVM_PROFILE_FILE: "/app/coverage/cargo-test-%p-%m.profraw" -# run: | -# docker compose -f docker-compose.test.yml run \ -# -e CARGO_TARGET_DIR \ -# -e CARGO_INCREMENTAL \ -# -e RUSTFLAGS \ -# -e LLVM_PROFILE_FILE \ -# agent-test bash -c "cargo clean && cargo test --verbose && sync" -# -# - name: Verify profraw files exist -# run: | -# docker compose -f docker-compose.test.yml run agent-test \ -# find /app/coverage -type f -name "*.profraw" | wc -l || true -# -# - name: Generate coverage report inside container -# run: | -# docker compose -f docker-compose.test.yml run agent-test bash -c " -# rustup component add llvm-tools && -# grcov /app/coverage \ -# --binary-path /app/target/debug \ -# -s /app \ -# --llvm \ -# -t lcov \ -# --branch \ -# --ignore-not-existing \ -# --ignore '/app/target/*' \ -# --ignore '/*' \ -# -o /app/lcov.info -# " -# -# - name: Copy lcov.info from container to host -# run: | -# docker compose -f docker-compose.test.yml cp agent-test:/app/lcov.info ./lcov.info -# -# - name: Remove container -# run: | -# docker rm agent-test-run -# -# - name: Show basic coverage report info (debug) -# run: | -# echo "lcov.info size:" $(wc -c ./lcov.info | awk '{print $1}') -# head -n 30 ./lcov.info || true -# -# - name: Upload coverage to Codecov -# uses: codecov/codecov-action@v5 -# with: -# files: ./lcov.info -# flags: unittests -# name: rust-unit-coverage -# verbose: true -# fail_ci_if_error: true -# env: -# CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} name: Codecov Rust on: @@ -99,15 +16,11 @@ env: jobs: coverage: runs-on: ubuntu-latest + steps: - uses: actions/checkout@v4 - - name: Install Rust toolchain - uses: dtolnay/rust-toolchain@stable - with: - components: llvm-tools-preview - - - name: Build test image (with grcov included) + - name: Build test image run: docker compose -f docker-compose.test.yml build agent-test - name: Start agent-test container @@ -115,28 +28,32 @@ jobs: - name: Run tests inside container run: | - docker compose -f docker-compose.test.yml exec \ + docker compose -f docker-compose.test.yml exec -T \ -e CARGO_TARGET_DIR \ -e CARGO_INCREMENTAL \ -e RUSTFLAGS \ -e LLVM_PROFILE_FILE \ agent-test bash -c " - cargo clean && + mkdir -p /app/coverage && + rm -rf /app/target/* /app/coverage/* && cargo test --verbose && sync " - name: Verify profraw files exist run: | - docker compose -f docker-compose.test.yml exec \ - -e LLVM_PROFILE_FILE \ - agent-test find /app/coverage -type f -name "*.profraw" | wc -l || true + docker compose -f docker-compose.test.yml exec -T \ + agent-test bash -c ' + count=$(find /app/coverage -type f -name "*.profraw" | wc -l) + echo "profraw files: $count" + test "$count" -gt 0 + ' - name: Generate coverage report inside container run: | - docker compose -f docker-compose.test.yml exec \ + docker compose -f docker-compose.test.yml exec -T \ agent-test bash -c " - rustup component add llvm-tools && + rustup component add llvm-tools-preview && grcov /app/coverage \ --binary-path /app/target/debug \ -s /app \ @@ -145,17 +62,17 @@ jobs: --branch \ --ignore-not-existing \ --ignore '/app/target/*' \ - --ignore '/*' \ -o /app/lcov.info " - name: Copy lcov.info from container to host run: docker compose -f docker-compose.test.yml cp agent-test:/app/lcov.info ./lcov.info - - name: Show basic coverage report info (debug) + - name: Show basic coverage report info run: | echo "lcov.info size:" $(wc -c ./lcov.info | awk '{print $1}') head -n 30 ./lcov.info || true + test -s ./lcov.info - name: Upload coverage to Codecov uses: codecov/codecov-action@v5 @@ -165,8 +82,9 @@ jobs: name: rust-unit-coverage verbose: true fail_ci_if_error: true - env: - CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} + token: ${{ secrets.CODECOV_TOKEN }} + use_pypi: true - name: Stop and remove container - run: docker compose -f docker-compose.test.yml down \ No newline at end of file + if: always() + run: docker compose -f docker-compose.test.yml down --volumes \ No newline at end of file diff --git a/CITATION.cff b/CITATION.cff index e7a7092..4bbfc3f 100644 --- a/CITATION.cff +++ b/CITATION.cff @@ -27,5 +27,5 @@ keywords: - self-hosted - portabase license: Apache-2.0 -version: 1.11.1 +version: 1.12.1 date-released: '2026-02-24' diff --git a/Cargo.lock b/Cargo.lock index c7b0809..3c4b240 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3083,7 +3083,7 @@ dependencies = [ [[package]] name = "portabase-agent" -version = "1.11.1" +version = "1.12.1" dependencies = [ "aes", "aes-gcm", diff --git a/Cargo.toml b/Cargo.toml index fa3a8d0..6a36a77 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "portabase-agent" -version = "1.11.1" +version = "1.12.1" edition = "2024" [dependencies] diff --git a/docker-compose.yml b/docker-compose.yml index d3016a3..189b569 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -12,14 +12,14 @@ services: - ./databases.json:/config/config.json #- ./databases.toml:/config/config.toml #- /var/run/docker.sock:/var/run/docker.sock - #- cargo-target:/app/target +# - cargo-target:/app/target - databases_sqlite-data:/sqlite-data/workspace/data - ./scripts/sqlite/test-db:/sqlite-data-2/workspace/data environment: APP_ENV: development LOG: debug TZ: "Europe/Paris" - EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiNzNlZmJhNjMtNTkzMy00Mzk3LWI0ZmMtMjlmNTViNmI5YzA4IiwibWFzdGVyS2V5QjY0IjoiQlhWM1hvbEM2NTZTVjdkTmdjV1BHUWxrKytycExJNmxHRGk3Q1BCNWllbz0ifQ==" + EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiMGZiNDYyMmUtMTMxNS00MzMxLTlkMTMtZWMzMjAyZjZiNTIwIiwibWFzdGVyS2V5QjY0IjoiMUh0djdtWCtYVkJxL0IzUEV2WDlZZjlQeUdVZW5oRHlXemo5THRqNW90WT0ifQ==" #CHUNK_SIZE_MB: "1" #POOLING: 1 #DATABASES_CONFIG_FILE: "config.toml" @@ -28,10 +28,18 @@ services: networks: - portabase + cpus: "1.50" + + mem_limit: 4g + memswap_limit: 4g + + pids_limit: 512 + + volumes: cargo-registry: cargo-git: - #cargo-target: +# cargo-target: databases_sqlite-data: external: true diff --git a/justfile b/justfile index 45b7b7c..8a98edf 100644 --- a/justfile +++ b/justfile @@ -58,4 +58,14 @@ seed-all: just seed-sqlite just seed-mongo just seed-firebird - just seed-mssql \ No newline at end of file + just seed-mssql + +test: + echo "Build test image (with grcov included)" + docker compose -f docker-compose.test.yml build agent-test + echo "Start agent-test container" + docker compose -f docker-compose.test.yml up -d agent-test + echo "Run tests inside container" + docker compose -f docker-compose.test.yml exec -e CARGO_INCREMENTAL -e RUSTFLAGS -e LLVM_PROFILE_FILE agent-test bash -c "cargo test --verbose && sync" + echo "Down volumes tests databases" + docker compose -f docker-compose.test.yml down --volumes diff --git a/src/domain/factory.rs b/src/domain/factory.rs index e4907ff..df5651d 100644 --- a/src/domain/factory.rs +++ b/src/domain/factory.rs @@ -5,20 +5,21 @@ use crate::domain::postgres::{detect_format_from_file, detect_format_from_size}; use crate::domain::redis::database::RedisDatabase; use crate::domain::sqlite::database::SqliteDatabase; use crate::domain::valkey::database::ValkeyDatabase; +use crate::domain::firebird::database::FirebirdDatabase; +use crate::domain::mariadb::database::MariaDBDatabase; +use crate::domain::mssql::database::MssqlDatabase; +use crate::services::backup::logger::JobLogger; use crate::services::config::{DatabaseConfig, DbType}; use anyhow::Result; use std::path::{Path, PathBuf}; use std::sync::Arc; -use crate::domain::firebird::database::FirebirdDatabase; -use crate::domain::mariadb::database::MariaDBDatabase; -use crate::domain::mssql::database::MssqlDatabase; #[async_trait::async_trait] pub trait Database: Send + Sync { fn file_extension(&self) -> &'static str; async fn ping(&self) -> Result; - async fn backup(&self, backup_dir: &Path) -> Result; - async fn restore(&self, restore_file: &Path) -> Result<()>; + async fn backup(&self, backup_dir: &Path, logger: Arc) -> Result; + async fn restore(&self, restore_file: &Path, logger: Arc) -> Result<()>; } pub struct DatabaseFactory; diff --git a/src/domain/firebird/backup.rs b/src/domain/firebird/backup.rs index b1d90bf..959bf9e 100644 --- a/src/domain/firebird/backup.rs +++ b/src/domain/firebird/backup.rs @@ -1,47 +1,53 @@ +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, backup_dir: PathBuf, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("Starting 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); - let db_path = format!( - "{}/{}:{}", - cfg.host, - cfg.port, - cfg.database - ); - - info!("Firebird database target: {}", db_path); - info!("Backup file: {}", file_path.display()); - + logger.log("info", format!("Firebird target: {} → {}", db_path, file_path.display())); + + let start = Instant::now(); let output = Command::new("gbak") .arg("-b") .arg("-v") .arg("-user").arg(&cfg.username) .arg("-password").arg(&cfg.password) - .arg(db_path) + .arg(&db_path) .arg(&file_path) .output() .with_context(|| format!("Failed to run gbak 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_output = if stderr.is_empty() && stdout.is_empty() { + None + } else { + Some(format!("{}{}", stdout, stderr).trim().to_string()) + }; if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - error!("Firebird backup failed: {}", stderr); + logger.log_command("gbak", combined_output, Some(exit_code), Some(duration_ms)); anyhow::bail!("Firebird backup failed for {}: {}", cfg.name, stderr); } - info!("Firebird backup completed: {}", file_path.display()); - + logger.log_command("gbak", combined_output, Some(0), Some(duration_ms)); + logger.log("info", format!("Firebird backup completed for {}", cfg.name)); Ok(file_path) }) .await? diff --git a/src/domain/firebird/database.rs b/src/domain/firebird/database.rs index b2af6da..93eb8ae 100644 --- a/src/domain/firebird/database.rs +++ b/src/domain/firebird/database.rs @@ -1,10 +1,12 @@ use super::{backup, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; use anyhow::Result; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; pub struct FirebirdDatabase { cfg: DatabaseConfig, @@ -27,21 +29,22 @@ impl Database for FirebirdDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; let res = backup::run( self.cfg.clone(), dir.to_path_buf(), self.file_extension(), + logger, ) .await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/firebird/restore.rs b/src/domain/firebird/restore.rs index d6ac1e7..b6a01dc 100644 --- a/src/domain/firebird/restore.rs +++ b/src/domain/firebird/restore.rs @@ -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) -> 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(()) }) diff --git a/src/domain/mariadb/backup.rs b/src/domain/mariadb/backup.rs index 1a2bfef..a4b59e6 100644 --- a/src/domain/mariadb/backup.rs +++ b/src/domain/mariadb/backup.rs @@ -1,39 +1,41 @@ use crate::domain::mariadb::connection::{select_mariadb_path, server_version}; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use anyhow::{Context, Result}; use std::collections::HashMap; 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, backup_dir: PathBuf, env: HashMap, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("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) => { - debug!("Mariadb version detected: {}", v); + logger.log("info", format!("MariaDB 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: {}", e)); return Err(e.into()); } }; - info!("Mariadb version found: {}", version); - let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension)); + let _mariadb_dump = select_mariadb_path(&version).join("mariadb-dump"); - let mariadb_dump = select_mariadb_path(&version).join("mariadb-dump"); - info!("Mariadb dump found: {}", mariadb_dump.display()); - + logger.log("debug", format!("Using mariadb-dump at {}", _mariadb_dump.display())); + logger.log("info", format!("Running mariadb-dump for {}", cfg.name)); + let start = Instant::now(); let output = Command::new("mariadb-dump") .arg("--host").arg(&cfg.host) .arg("--port").arg(cfg.port.to_string()) @@ -45,8 +47,9 @@ pub async fn run( .arg("--quick") .arg("--skip-lock-tables") .arg("--no-create-db") - .arg("--skip-add-drop-table") + .arg("--skip-add-drop-table") .arg("--compress") + .arg("--verbose") .arg("--max-allowed-packet=512M") .arg("--net-buffer-length=16K") .arg("--default-character-set=utf8mb4") @@ -55,12 +58,18 @@ pub async fn run( .envs(env) .output() .with_context(|| format!("Failed to run mariadb-dump for {}", cfg.name))?; + let duration_ms = start.elapsed().as_millis() as f64; + let exit_code = output.status.code().unwrap_or(-1); if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); + let stderr = String::from_utf8_lossy(&output.stderr).to_string(); + logger.log_command("mariadb-dump", Some(stderr.clone()), Some(exit_code), Some(duration_ms)); anyhow::bail!("Mariadb backup failed for {}: {}", cfg.name, stderr); } + let stderr = String::from_utf8_lossy(&output.stderr).to_string(); + logger.log_command("mariadb-dump", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms)); + logger.log("info", format!("mariadb-dump completed for {}", cfg.name)); Ok(file_path) }) .await? diff --git a/src/domain/mariadb/database.rs b/src/domain/mariadb/database.rs index 37d04fc..8ae7788 100644 --- a/src/domain/mariadb/database.rs +++ b/src/domain/mariadb/database.rs @@ -1,11 +1,13 @@ use super::{backup, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; use anyhow::Result; use async_trait::async_trait; use std::collections::HashMap; use std::path::{Path, PathBuf}; +use std::sync::Arc; pub struct MariaDBDatabase { cfg: DatabaseConfig, @@ -33,22 +35,23 @@ impl Database for MariaDBDatabase { ping::run(self.cfg.clone(), self.build_env().clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; let res = backup::run( self.cfg.clone(), dir.to_path_buf(), self.build_env().clone(), self.file_extension(), + logger, ) .await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/mariadb/restore.rs b/src/domain/mariadb/restore.rs index 5357cbf..d6fc555 100644 --- a/src/domain/mariadb/restore.rs +++ b/src/domain/mariadb/restore.rs @@ -1,27 +1,26 @@ +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) -> 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 +30,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") @@ -55,23 +62,26 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> { .with_context(|| format!("Failed to start MariaDB restore for {}", cfg.name))?; let mut stdin = child.stdin.take().context("Failed to open child stdin")?; - stdin - .write_all(sql_content.as_bytes()) - .context("Failed to write SQL content to MariaDB stdin")?; - stdin.flush()?; + std::io::copy(&mut file, &mut stdin) + .context("Failed to stream SQL content to MariaDB stdin")?; drop(stdin); let output = child .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(()) }); diff --git a/src/domain/mongodb/backup.rs b/src/domain/mongodb/backup.rs index d4724b3..022131e 100644 --- a/src/domain/mongodb/backup.rs +++ b/src/domain/mongodb/backup.rs @@ -1,35 +1,48 @@ use crate::domain::mongodb::connection::{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, backup_dir: PathBuf, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("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"); let uri = get_mongo_uri(cfg.clone())?; + logger.log("info", format!("Running mongodump for {}", cfg.name)); + + let start = Instant::now(); let output = Command::new(mongodump) .arg(format!("--uri={}", uri)) .arg(format!("--archive={}", file_path.display())) .arg("--gzip") + .arg("--verbose") .output() .context("MongoDB backup failed")?; + 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 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); } - info!("MongoDB backup completed for {}", cfg.name); + + logger.log_command("mongodump", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms)); + logger.log("info", format!("MongoDB backup completed for {}", cfg.name)); Ok(file_path) }) .await? diff --git a/src/domain/mongodb/database.rs b/src/domain/mongodb/database.rs index e4c40b8..0fe9ee2 100644 --- a/src/domain/mongodb/database.rs +++ b/src/domain/mongodb/database.rs @@ -1,9 +1,11 @@ use anyhow::Result; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; use super::{backup, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; @@ -27,16 +29,16 @@ impl Database for MongoDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; - let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension()).await; + let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension(), logger).await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/mongodb/restore.rs b/src/domain/mongodb/restore.rs index 0a8012d..627aa4d 100644 --- a/src/domain/mongodb/restore.rs +++ b/src/domain/mongodb/restore.rs @@ -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) -> 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,12 +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); - let source_db = extract_db_name(&dry_output) - .context("Could not detect source database name from archive")?; + 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(|| { + logger.log("info", format!("Could not detect source database from archive, falling back to configured database: {}", cfg.database)); + cfg.database.clone() + }); - info!("Detected 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())) @@ -43,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? diff --git a/src/domain/mssql/backup.rs b/src/domain/mssql/backup.rs index 5edb941..be564e1 100644 --- a/src/domain/mssql/backup.rs +++ b/src/domain/mssql/backup.rs @@ -1,16 +1,19 @@ +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, backup_dir: PathBuf, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("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!( @@ -18,29 +21,35 @@ pub async fn run( cfg.host, cfg.port, cfg.database, cfg.username, cfg.password ); - info!( - "MSSQL backup: {}:{}/{} → {}", - cfg.host, - cfg.port, - cfg.database, - file_path.display() - ); + logger.log("info", format!("MSSQL backup: {}:{}/{} → {}", cfg.host, cfg.port, cfg.database, file_path.display())); + let start = Instant::now(); let output = Command::new("sqlpackage") .arg("/a:Export") .arg(format!("/scs:{}", connection_string)) .arg(format!("/tf:{}", file_path.display())) .output() .with_context(|| format!("Failed to run sqlpackage 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(); if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - let stdout = String::from_utf8_lossy(&output.stdout); - error!("MSSQL backup failed — stderr: {} stdout: {}", stderr, stdout); + logger.log("error", format!("MSSQL backup failed — stderr: {} stdout: {}", stderr, stdout)); + let out = format!("stderr: {} stdout: {}", stderr, stdout); + logger.log_command("sqlpackage", Some(out), Some(exit_code), Some(duration_ms)); anyhow::bail!("MSSQL backup failed for {}: {}", cfg.name, stderr); } - info!("MSSQL backup completed: {}", file_path.display()); + let combined = if stdout.is_empty() && stderr.is_empty() { + None + } else { + Some(format!("{}{}", stdout, stderr).trim().to_string()) + }; + logger.log_command("sqlpackage", combined, Some(0), Some(duration_ms)); + logger.log("info", format!("MSSQL backup completed for {}", cfg.name)); Ok(file_path) }) .await? diff --git a/src/domain/mssql/database.rs b/src/domain/mssql/database.rs index 673cb95..fda0652 100644 --- a/src/domain/mssql/database.rs +++ b/src/domain/mssql/database.rs @@ -1,10 +1,12 @@ use super::{backup, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; use anyhow::Result; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; pub struct MssqlDatabase { cfg: DatabaseConfig, @@ -26,21 +28,22 @@ impl Database for MssqlDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; let res = backup::run( self.cfg.clone(), dir.to_path_buf(), self.file_extension(), + logger, ) .await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/mssql/restore.rs b/src/domain/mssql/restore.rs index 181313d..af56d6c 100644 --- a/src/domain/mssql/restore.rs +++ b/src/domain/mssql/restore.rs @@ -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) -> 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? diff --git a/src/domain/mysql/backup.rs b/src/domain/mysql/backup.rs index 9cc1071..ebadcc2 100644 --- a/src/domain/mysql/backup.rs +++ b/src/domain/mysql/backup.rs @@ -1,65 +1,71 @@ use crate::domain::mysql::connection::server_version; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use anyhow::{Context, Result}; use std::collections::HashMap; 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, backup_dir: PathBuf, env: HashMap, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("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)) { + let _version = match futures::executor::block_on(server_version(&cfg)) { Ok(v) => { - debug!("Mysql version detected: {}", v); + logger.log("debug", format!("MySQL 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: {}", e)); return Err(e.into()); } }; - info!("Mysql version found: {}", version); - let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension)); + logger.log("info", format!("Running mysqldump for {}", cfg.name)); + + let start = Instant::now(); let output = Command::new("mysqldump") - .arg("--host") - .arg(&cfg.host) - .arg("--port") - .arg(cfg.port.to_string()) - .arg("--user") - .arg(&cfg.username) + .arg("--host").arg(&cfg.host) + .arg("--port").arg(cfg.port.to_string()) + .arg("--user").arg(&cfg.username) .arg("--routines") .arg("--events") .arg("--triggers") + .arg("--verbose") .arg("--single-transaction") .arg("--quick") .arg("--skip-lock-tables") .arg("--skip-add-drop-table") - .arg("--no-create-db") // IMPORTANT + .arg("--no-create-db") .arg("--default-character-set=utf8mb4") - .arg(&cfg.database) // IMPORTANT: NOT --databases - .arg("-r") - .arg(&file_path) + .arg(&cfg.database) + .arg("-r").arg(&file_path) .envs(env) .output() .with_context(|| format!("Failed to run mysqldump for {}", cfg.name))?; + let duration_ms = start.elapsed().as_millis() as f64; + let exit_code = output.status.code().unwrap_or(-1); + + let _stdout = String::from_utf8_lossy(&output.stdout).to_string(); + let stderr = String::from_utf8_lossy(&output.stderr).to_string(); if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - info!("mysqldump stderr: {}", stderr); + logger.log_command("mysqldump", Some(stderr.clone()), Some(exit_code), Some(duration_ms)); anyhow::bail!("MySQL backup failed for {}: {}", cfg.name, stderr); } - info!("Output {}", String::from_utf8_lossy(&output.stdout)); + logger.log_command("mysqldump", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms)); + logger.log("info", format!("mysqldump completed for {}", cfg.name)); Ok(file_path) }) diff --git a/src/domain/mysql/database.rs b/src/domain/mysql/database.rs index 9c39bd6..12a34ae 100644 --- a/src/domain/mysql/database.rs +++ b/src/domain/mysql/database.rs @@ -1,11 +1,13 @@ use super::{backup, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; use anyhow::Result; use async_trait::async_trait; use std::collections::HashMap; use std::path::{Path, PathBuf}; +use std::sync::Arc; pub struct MySQLDatabase { cfg: DatabaseConfig, @@ -33,22 +35,23 @@ impl Database for MySQLDatabase { ping::run(self.cfg.clone(), self.build_env().clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; let res = backup::run( self.cfg.clone(), dir.to_path_buf(), self.build_env().clone(), self.file_extension(), + logger, ) .await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/mysql/restore.rs b/src/domain/mysql/restore.rs index 8f569aa..955b4c6 100644 --- a/src/domain/mysql/restore.rs +++ b/src/domain/mysql/restore.rs @@ -1,27 +1,26 @@ +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) -> 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("mysql") + let drop_start = Instant::now(); + let drop_output = Command::new("mysql") .arg("--host") .arg(&cfg.host) .arg("--port") @@ -31,14 +30,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") @@ -55,23 +62,26 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf) -> Result<()> { .with_context(|| format!("Failed to start mysql restore for {}", cfg.name))?; let mut stdin = child.stdin.take().context("Failed to open child stdin")?; - stdin - .write_all(sql_content.as_bytes()) - .context("Failed to write SQL content to mysql stdin")?; - stdin.flush()?; + std::io::copy(&mut file, &mut stdin) + .context("Failed to stream SQL content to mysql stdin")?; drop(stdin); let output = child .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(()) }); diff --git a/src/domain/postgres/backup.rs b/src/domain/postgres/backup.rs index 9ec9eb4..61f24f3 100644 --- a/src/domain/postgres/backup.rs +++ b/src/domain/postgres/backup.rs @@ -1,82 +1,89 @@ 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}; use super::format::PostgresDumpFormat; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; pub async fn run( cfg: DatabaseConfig, format: PostgresDumpFormat, backup_dir: PathBuf, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("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) => { - 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: {}", e)); return Err(e.into()); } }; let pg_dump = select_pg_path(&version).join("pg_dump"); - - debug!("Using pg_dump at {:?}", pg_dump); + logger.log("debug", format!("Using pg_dump at {:?}", pg_dump)); match format { PostgresDumpFormat::Fc => { - info!("Running FC backup for {}", cfg.name); + logger.log("info", format!("Running FC backup for {}", cfg.name)); + let file_path = backup_dir.join(format!("{}.dump", cfg.generated_id)); let url = format!( "postgresql://{}:{}@{}:{}/{}", cfg.username, cfg.password, cfg.host, cfg.port, cfg.database ); - let status = Command::new(&pg_dump) - .arg("--dbname") - .arg(&url) + let start = Instant::now(); + let output = Command::new(&pg_dump) + .arg("--dbname").arg(&url) .arg("-Fc") - .arg("-f") - .arg(&file_path) + .arg("-f").arg(&file_path) .arg("-v") .arg("--compress=3") - .status(); + .output(); + let duration_ms = start.elapsed().as_millis() as f64; - match status { - Ok(s) if s.success() => info!( - "FC backup completed successfully for {} at {:?}", - cfg.name, file_path - ), - Ok(s) => { - error!("FC backup failed with status {:?} for {}", s, cfg.name); - anyhow::bail!("Postgres backup failed for {}", cfg.name); + match output { + Ok(o) => { + let stderr = String::from_utf8_lossy(&o.stderr).to_string(); + let exit_code = o.status.code().unwrap_or(-1); + + if o.status.success() { + logger.log("info", format!("FC backup completed successfully for {} at {:?}", cfg.name, file_path)); + logger.log_command("pg_dump", if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms)); + } else { + logger.log("error", format!("FC backup failed with status {:?} for {}", o.status, cfg.name)); + logger.log_command("pg_dump", Some(stderr), Some(exit_code), Some(duration_ms)); + anyhow::bail!("Postgres backup failed for {}", cfg.name); + } } Err(e) => { - error!("Error executing pg_dump for {}: {:?}", cfg.name, e); + logger.log("error", format!("Error executing pg_dump for {}: {:?}", cfg.name, e)); + logger.log_command("pg_dump", Some(e.to_string()), Some(-1), Some(duration_ms)); return Err(e.into()); } } - info!("Backup finished for database {}", cfg.name); + logger.log("info", format!("Backup finished for database {}", cfg.name)); Ok(file_path) } PostgresDumpFormat::Fd => { - info!("Running FD backup for {}", cfg.name); + logger.log("info", format!("Running FD backup for {}", cfg.name)); + let dump_dir = backup_dir.join(format!("{}_dir", cfg.generated_id)); let tar_file = backup_dir.join(format!("{}.tar.gz", cfg.generated_id)); if let Err(e) = std::fs::create_dir_all(&dump_dir) { - error!( - "Failed to create dump directory {:?} for {}: {:?}", - dump_dir, cfg.name, e - ); + logger.log("error", format!("Failed to create dump directory {:?} for {}: {:?}", dump_dir, cfg.name, e)); return Err(e.into()); } @@ -84,59 +91,58 @@ pub async fn run( "postgresql://{}:{}@{}:{}/{}", cfg.username, cfg.password, cfg.host, cfg.port, cfg.database ); + let cmd_label = format!("pg_dump -Fd {}", url); - let status = Command::new(&pg_dump) - .arg("--dbname") - .arg(&url) + let start = Instant::now(); + let output = Command::new(&pg_dump) + .arg("--dbname").arg(&url) .arg("-Fd") - .arg("-j") - .arg("4") - .arg("-f") - .arg(&dump_dir) + .arg("-j").arg("4") + .arg("-f").arg(&dump_dir) .arg("-v") - .status(); + .output(); + let duration_ms = start.elapsed().as_millis() as f64; - match status { - Ok(s) if s.success() => { - info!("FD backup pg_dump completed successfully for {}", cfg.name) - } - Ok(s) => { - error!( - "FD backup pg_dump failed with status {:?} for {}", - s, cfg.name - ); - anyhow::bail!("Postgres FD backup failed for {}", cfg.name); + match output { + Ok(o) => { + let stderr = String::from_utf8_lossy(&o.stderr).to_string(); + let exit_code = o.status.code().unwrap_or(-1); + if o.status.success() { + logger.log("info", format!("FD backup pg_dump completed successfully for {}", cfg.name)); + logger.log_command(cmd_label, if stderr.is_empty() { None } else { Some(stderr) }, Some(0), Some(duration_ms)); + } else { + logger.log("error", format!("FD backup pg_dump failed with status {:?} for {}", o.status, cfg.name)); + logger.log_command(cmd_label, Some(stderr), Some(exit_code), Some(duration_ms)); + anyhow::bail!("Postgres FD backup failed for {}", cfg.name); + } } Err(e) => { - error!("Error executing pg_dump for {}: {:?}", cfg.name, e); + logger.log("error", format!("Error executing pg_dump for {}: {:?}", cfg.name, e)); + logger.log_command(cmd_label, Some(e.to_string()), Some(-1), Some(duration_ms)); return Err(e.into()); } } match std::fs::File::create(&tar_file) { Ok(tar_gz) => { - let enc = - flate2::write::GzEncoder::new(tar_gz, flate2::Compression::default()); + let enc = flate2::write::GzEncoder::new(tar_gz, flate2::Compression::default()); let mut tar = tar::Builder::new(enc); if let Err(e) = tar.append_dir_all(".", &dump_dir) { - error!("Failed to append dump_dir to tar for {}: {:?}", cfg.name, e); + logger.log("error", format!("Failed to append dump_dir to tar for {}: {:?}", cfg.name, e)); return Err(e.into()); } if let Err(e) = tar.finish() { - error!("Failed to finish tar archive for {}: {:?}", cfg.name, e); + logger.log("error", format!("Failed to finish tar archive for {}: {:?}", cfg.name, e)); return Err(e.into()); } - info!("FD backup archive created at {:?}", tar_file); + logger.log("info", format!("FD backup archive created at {:?}", tar_file)); } Err(e) => { - error!( - "Failed to create tar.gz file {:?} for {}: {:?}", - tar_file, cfg.name, e - ); + logger.log("error", format!("Failed to create tar.gz file {:?} for {}: {:?}", tar_file, cfg.name, e)); return Err(e.into()); } } - info!("Backup finished for database {}", cfg.name); + logger.log("info", format!("Backup finished for database {}", cfg.name)); Ok(tar_file) } } diff --git a/src/domain/postgres/database.rs b/src/domain/postgres/database.rs index 12aa7ce..205de95 100644 --- a/src/domain/postgres/database.rs +++ b/src/domain/postgres/database.rs @@ -1,9 +1,11 @@ use anyhow::Result; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; use super::{backup, format::PostgresDumpFormat, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; @@ -31,16 +33,16 @@ impl Database for PostgresDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; - let res = backup::run(self.cfg.clone(), self.format, dir.to_path_buf()).await; + let res = backup::run(self.cfg.clone(), self.format, dir.to_path_buf(), logger).await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/postgres/restore.rs b/src/domain/postgres/restore.rs index f722556..2eca26a 100644 --- a/src/domain/postgres/restore.rs +++ b/src/domain/postgres/restore.rs @@ -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, ) -> 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(()) }) diff --git a/src/domain/redis/backup.rs b/src/domain/redis/backup.rs index 9fbe995..c63b1b0 100644 --- a/src/domain/redis/backup.rs +++ b/src/domain/redis/backup.rs @@ -1,21 +1,26 @@ +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, backup_dir: PathBuf, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("Starting Redis backup for database {}", cfg.name); + logger.log( + "info", + format!("Starting Redis backup for database {}", cfg.name), + ); let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension)); let mut cmd = Command::new("redis-cli"); - cmd.arg("-h") .arg(&cfg.host) .arg("-p") @@ -24,40 +29,66 @@ pub async fn run( if !cfg.username.is_empty() { cmd.arg("--user").arg(&cfg.username); } - if !cfg.password.is_empty() { cmd.arg("-a").arg(&cfg.password); } - cmd.arg("--rdb").arg(&file_path); - debug!("Command Backup: {:?}", cmd); + logger.log("info", format!("Running redis-cli --rdb for {}", cfg.name)); + let start = Instant::now(); let output = cmd.output().context("Redis backup command failed")?; + 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); let stdout = String::from_utf8_lossy(&output.stdout); if !output.status.success() { if stderr.contains("NOAUTH") { - error!( - "Redis backup failed for {}: Authentication required (NOAUTH)", - cfg.name + logger.log( + "error", + format!( + "Redis backup failed for {}: Authentication required (NOAUTH)", + cfg.name + ), + ); + logger.log_command( + "redis-cli", + Some("Authentication required (NOAUTH)".into()), + Some(exit_code), + Some(duration_ms), ); anyhow::bail!( "Redis backup failed for {}: Authentication required", cfg.name ); } else { - error!("Redis backup failed for {}: {}", cfg.name, stderr); + logger.log( + "error", + format!("Redis backup failed for {}: {}", cfg.name, stderr), + ); + logger.log_command( + "redis-cli", + Some(stderr.to_string()), + Some(exit_code), + Some(duration_ms), + ); anyhow::bail!("Redis backup failed for {}: {}", cfg.name, stderr); } } - info!( - "Redis backup completed for {}. Output: {}", - cfg.name, stdout + logger.log_command( + "redis-cli", + if stdout.is_empty() { + None + } else { + Some(stdout.to_string()) + }, + Some(0), + Some(duration_ms), ); + logger.log("info", format!("Redis backup completed for {}", cfg.name)); Ok(file_path) }) diff --git a/src/domain/redis/database.rs b/src/domain/redis/database.rs index 33e67e2..91a2cb7 100644 --- a/src/domain/redis/database.rs +++ b/src/domain/redis/database.rs @@ -1,9 +1,11 @@ use anyhow::{Result, bail}; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; use crate::domain::factory::Database; use crate::domain::redis::{backup, ping}; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; @@ -27,14 +29,14 @@ impl Database for RedisDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; - let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension()).await; + let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension(), logger).await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, _file: &Path) -> Result<()> { + async fn restore(&self, _file: &Path, _logger: Arc) -> Result<()> { bail!("Restore not supported for Redis databases") } } diff --git a/src/domain/sqlite/backup.rs b/src/domain/sqlite/backup.rs index 2ab9cd1..9b8f12c 100644 --- a/src/domain/sqlite/backup.rs +++ b/src/domain/sqlite/backup.rs @@ -1,16 +1,19 @@ +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, backup_dir: PathBuf, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("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"); @@ -19,27 +22,35 @@ pub async fn run( }; let db_path = PathBuf::from(db_path_str); - info!("database path: {}", db_path.display()); + logger.log("info", format!("Database path: {}", db_path.display())); if !db_path.exists() { + logger.log("error", format!("SQLite database file not found: {}", db_path.display())); anyhow::bail!("SQLite database file not found: {}", db_path.display()); } let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension)); + logger.log("info", format!("Running sqlite3 backup for {}", cfg.name)); + + let start = Instant::now(); let output = Command::new("sqlite3") .arg(db_path.as_os_str()) .arg(format!(".backup '{}'", file_path.display())) .output() .context("SQLite backup command failed to start")?; - info!("Backup successful: {:?}", output); + let duration_ms = start.elapsed().as_millis() as f64; + let exit_code = output.status.code().unwrap_or(-1); + if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - error!("SQLite backup failed for {}: {}", cfg.name, stderr); + let stderr = String::from_utf8_lossy(&output.stderr).to_string(); + logger.log("error", format!("SQLite backup failed for {}: {}", cfg.name, stderr)); + logger.log_command("sqlite3", Some(stderr.clone()), Some(exit_code), Some(duration_ms)); anyhow::bail!("SQLite backup failed for {}: {}", cfg.name, stderr); } - info!("SQLite backup completed for {}", cfg.name); + logger.log_command("sqlite3", None, Some(0), Some(duration_ms)); + logger.log("info", format!("SQLite backup completed for {}", cfg.name)); Ok(file_path) }) .await? diff --git a/src/domain/sqlite/database.rs b/src/domain/sqlite/database.rs index 2579061..1e0f2d1 100644 --- a/src/domain/sqlite/database.rs +++ b/src/domain/sqlite/database.rs @@ -1,9 +1,11 @@ use anyhow::Result; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; use super::{backup, ping, restore}; use crate::domain::factory::Database; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; @@ -27,15 +29,15 @@ impl Database for SqliteDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; - let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension()).await; + let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension(), logger).await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, file: &Path) -> Result<()> { + async fn restore(&self, file: &Path, logger: Arc) -> 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 } diff --git a/src/domain/sqlite/restore.rs b/src/domain/sqlite/restore.rs index a386a5d..47cfccc 100644 --- a/src/domain/sqlite/restore.rs +++ b/src/domain/sqlite/restore.rs @@ -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) -> 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? diff --git a/src/domain/valkey/backup.rs b/src/domain/valkey/backup.rs index d064b19..56fe05b 100644 --- a/src/domain/valkey/backup.rs +++ b/src/domain/valkey/backup.rs @@ -1,21 +1,26 @@ +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, backup_dir: PathBuf, file_extension: &'static str, + logger: Arc, ) -> Result { tokio::task::spawn_blocking(move || -> Result { - debug!("Starting Valkey backup for database {}", cfg.name); + logger.log( + "info", + format!("Starting Valkey backup for database {}", cfg.name), + ); let file_path = backup_dir.join(format!("{}{}", cfg.generated_id, file_extension)); let mut cmd = Command::new("valkey-cli"); - cmd.arg("-h") .arg(&cfg.host) .arg("-p") @@ -24,40 +29,66 @@ pub async fn run( if !cfg.username.is_empty() { cmd.arg("--user").arg(&cfg.username); } - if !cfg.password.is_empty() { cmd.arg("-a").arg(&cfg.password); } - cmd.arg("--rdb").arg(&file_path); - debug!("Command Backup: {:?}", cmd); + logger.log("info", format!("Running valkey-cli --rdb for {}", cfg.name)); + let start = Instant::now(); let output = cmd.output().context("Valkey backup command failed")?; + 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); let stdout = String::from_utf8_lossy(&output.stdout); if !output.status.success() { if stderr.contains("NOAUTH") { - error!( - "Valkey backup failed for {}: Authentication required (NOAUTH)", - cfg.name + logger.log( + "error", + format!( + "Valkey backup failed for {}: Authentication required (NOAUTH)", + cfg.name + ), + ); + logger.log_command( + "valkey-cli", + Some("Authentication required (NOAUTH)".into()), + Some(exit_code), + Some(duration_ms), ); anyhow::bail!( "Valkey backup failed for {}: Authentication required", cfg.name ); } else { - error!("Valkey backup failed for {}: {}", cfg.name, stderr); + logger.log( + "error", + format!("Valkey backup failed for {}: {}", cfg.name, stderr), + ); + logger.log_command( + "valkey-cli", + Some(stderr.to_string()), + Some(exit_code), + Some(duration_ms), + ); anyhow::bail!("Valkey backup failed for {}: {}", cfg.name, stderr); } } - info!( - "Valkey backup completed for {}. Output: {}", - cfg.name, stdout + logger.log_command( + "valkey-cli", + if stdout.is_empty() { + None + } else { + Some(stdout.to_string()) + }, + Some(0), + Some(duration_ms), ); + logger.log("info", format!("Valkey backup completed for {}", cfg.name)); Ok(file_path) }) diff --git a/src/domain/valkey/database.rs b/src/domain/valkey/database.rs index 8696179..ec9b6a7 100644 --- a/src/domain/valkey/database.rs +++ b/src/domain/valkey/database.rs @@ -1,9 +1,11 @@ use anyhow::{Result, bail}; use async_trait::async_trait; use std::path::{Path, PathBuf}; +use std::sync::Arc; use crate::domain::factory::Database; use crate::domain::valkey::{backup, ping}; +use crate::services::backup::logger::JobLogger; use crate::services::config::DatabaseConfig; use crate::utils::locks::{DbOpLock, FileLock}; @@ -27,14 +29,14 @@ impl Database for ValkeyDatabase { ping::run(self.cfg.clone()).await } - async fn backup(&self, dir: &Path) -> Result { + async fn backup(&self, dir: &Path, logger: Arc) -> Result { FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?; - let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension()).await; + let res = backup::run(self.cfg.clone(), dir.to_path_buf(), self.file_extension(), logger).await; FileLock::release(&self.cfg.generated_id).await?; res } - async fn restore(&self, _file: &Path) -> Result<()> { + async fn restore(&self, _file: &Path, _logger: Arc) -> Result<()> { bail!("Restore not supported for Valkey databases") } } diff --git a/src/services/api/endpoints/agent/backup/mod.rs b/src/services/api/endpoints/agent/backup/mod.rs index 5a489e4..f47764c 100644 --- a/src/services/api/endpoints/agent/backup/mod.rs +++ b/src/services/api/endpoints/agent/backup/mod.rs @@ -2,6 +2,7 @@ pub mod upload; use crate::services::api::models::agent::backup::BackupResponse; use crate::services::api::{ApiClient, ApiError}; +use crate::services::backup::logger::JobLogEntry; use anyhow::Result; use reqwest::Method; use serde::Serialize; @@ -21,6 +22,9 @@ pub struct BackupUpdateRequest { pub size: Option, #[serde(rename = "generatedId")] pub generated_id: String, + pub logs: Vec, + #[serde(rename = "durationMs")] + pub duration_ms: f64, } impl ApiClient { @@ -49,12 +53,16 @@ impl ApiClient { status: impl Into, file_size: impl Into>, generated_id: impl Into, + job_logs: Vec, + duration_ms: f64, ) -> Result, ApiError> { let body = BackupUpdateRequest { backup_id: backup_id.into(), status: status.into(), size: file_size.into(), generated_id: generated_id.into(), + logs: job_logs, + duration_ms: duration_ms, }; let agent_id = agent_id.into(); diff --git a/src/services/api/endpoints/agent/restore/mod.rs b/src/services/api/endpoints/agent/restore/mod.rs index 82d53f7..553a9de 100644 --- a/src/services/api/endpoints/agent/restore/mod.rs +++ b/src/services/api/endpoints/agent/restore/mod.rs @@ -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, + #[serde(rename = "durationMs")] + pub duration_ms: f64, } impl ApiClient { @@ -17,10 +21,14 @@ impl ApiClient { agent_id: impl Into, generated_id: impl Into, status: impl Into, + job_logs: Vec, + duration_ms: f64, ) -> Result, ApiError> { let body = ResultRestoreRequest { generated_id: generated_id.into(), status: status.into(), + logs: job_logs, + duration_ms, }; let agent_id = agent_id.into(); diff --git a/src/services/backup/compressor.rs b/src/services/backup/compressor.rs index ad36a09..ecb8ebf 100644 --- a/src/services/backup/compressor.rs +++ b/src/services/backup/compressor.rs @@ -1,14 +1,17 @@ +use super::logger::JobLogger; use super::service::BackupService; use crate::utils::compress::compress_to_tar_gz_large; use anyhow::Result; use std::path::PathBuf; +use std::sync::Arc; impl BackupService { - pub async fn compress_backup(&self, backup_file: Option) -> Result { + pub async fn compress_backup(&self, backup_file: Option, logger: Arc) -> Result { let file = backup_file.ok_or_else(|| anyhow::anyhow!("No backup file generated"))?; - let compression = compress_to_tar_gz_large(&file).await?; - + logger.log("info", "Start compressing archive".to_string()); + let compression = compress_to_tar_gz_large(&file, logger).await?; + Ok(compression.compressed_path) } } diff --git a/src/services/backup/executor.rs b/src/services/backup/executor.rs index 4b88c86..018eed9 100644 --- a/src/services/backup/executor.rs +++ b/src/services/backup/executor.rs @@ -1,3 +1,4 @@ +use super::logger::JobLogger; use super::service::BackupService; use crate::services::api::models::agent::status::DatabaseStorage; use crate::services::config::DatabaseConfig; @@ -5,6 +6,8 @@ use crate::utils::common::BackupMethod; use crate::utils::locks::FileLock; use anyhow::Result; +use std::sync::Arc; +use std::time::Instant; use tempfile::TempDir; impl BackupService { @@ -16,9 +19,13 @@ impl BackupService { storages: Vec, encrypt: bool, ) -> Result<()> { + let logger = Arc::new(JobLogger::new()); + if FileLock::is_locked(&generated_id).await? { anyhow::bail!("backup already running"); } + let start = Instant::now(); + logger.log("info", "Database backup job started".to_string()); let backup = self.create_backup_record(&generated_id, &method).await?; let backup_id = backup.backup.id; @@ -26,21 +33,27 @@ impl BackupService { let temp_dir = TempDir::new()?; let tmp_path = temp_dir.path(); - let mut result = Self::run(db_cfg, tmp_path).await?; + let mut result = Self::run(db_cfg, tmp_path, Arc::clone(&logger)).await?; if result.status == "failed" { - self.send_result(result, vec![], &backup_id).await?; + 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, vec![], &backup_id, logs, duration_ms).await?; return Ok(()); } - let compressed = self.compress_backup(result.backup_file.take()).await?; + let compressed = self.compress_backup(result.backup_file.take(), Arc::clone(&logger)).await?; result.backup_file = Some(compressed); let uploads = self - .upload(result.clone(), method, storages, encrypt, &backup_id) + .upload(result.clone(), method, storages, encrypt, &backup_id, Arc::clone(&logger)) .await?; - self.send_result(result, uploads, &backup_id).await?; + logger.log("info", "Database backup 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, uploads, &backup_id, logs, duration_ms).await?; Ok(()) } diff --git a/src/services/backup/logger.rs b/src/services/backup/logger.rs new file mode 100644 index 0000000..0bb0182 --- /dev/null +++ b/src/services/backup/logger.rs @@ -0,0 +1,137 @@ +use chrono::Utc; +use serde::Serialize; +use std::sync::Mutex; +use tracing::{event, Level}; + +#[derive(Serialize, Clone, Debug)] +pub struct JobLogEntry { + pub timestamp: String, + #[serde(rename = "type")] + pub entry_type: &'static str, + pub level: &'static str, + pub message: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub command: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub output: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub exit_code: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub duration_ms: Option, +} + +#[derive(Default, Debug)] +pub struct JobLogger { + entries: Mutex>, +} + +impl JobLogger { + pub fn new() -> Self { + Self::default() + } + + fn trace_log(level: &'static str, message: &str) { + match level { + "trace" => event!(Level::TRACE, "{message}"), + "debug" => event!(Level::DEBUG, "{message}"), + "info" => event!(Level::INFO, "{message}"), + "warn" | "warning" => event!(Level::WARN, "{message}"), + "error" => event!(Level::ERROR, "{message}"), + _ => event!(Level::INFO, "{message}"), + } + } + + + #[allow(dead_code)] + pub fn log(&self, level: &'static str, message: impl Into) { + let message = message.into(); + + Self::trace_log(level, &message); + + let mut entries = self.entries.lock().unwrap(); + entries.push(JobLogEntry { + timestamp: Utc::now().to_rfc3339(), + entry_type: "log", + level, + message, + command: None, + output: None, + exit_code: None, + duration_ms: None, + }); + } + + pub fn log_command( + &self, + command: impl Into, + output: Option, + exit_code: Option, + duration_ms: Option, + ) { + let level = match exit_code { + Some(0) | None => "debug", + _ => "error", + }; + + + let cmd = command.into(); + let message = format!("Executed: {}", cmd); + + match level { + "debug" => { + tracing::debug!( + command = %cmd, + exit_code = ?exit_code, + duration_ms = ?duration_ms, + "Command executed" + ); + } + "error" => { + tracing::error!( + command = %cmd, + exit_code = ?exit_code, + duration_ms = ?duration_ms, + "Command failed" + ); + } + _ => { + tracing::info!( + command = %cmd, + exit_code = ?exit_code, + duration_ms = ?duration_ms, + "Command executed" + ); + } + } + + if let Some(output_text) = output.as_deref() { + for line in output_text.lines() { + if line.trim().is_empty() { + continue; + } + + match level { + "debug" => tracing::debug!("{line}"), + "error" => tracing::error!("{line}"), + _ => tracing::info!("{line}"), + } + } + } + + let mut entries = self.entries.lock().unwrap(); + entries.push(JobLogEntry { + timestamp: Utc::now().to_rfc3339(), + entry_type: "command", + level, + message, + command: Some(cmd), + output, + exit_code, + duration_ms, + }); + } + + pub fn into_entries(self) -> Vec { + self.entries.into_inner().unwrap() + } +} diff --git a/src/services/backup/mod.rs b/src/services/backup/mod.rs index 98f1aec..b761556 100644 --- a/src/services/backup/mod.rs +++ b/src/services/backup/mod.rs @@ -1,4 +1,5 @@ pub mod compressor; +pub mod logger; pub mod dispatcher; pub mod executor; pub mod helpers; diff --git a/src/services/backup/result.rs b/src/services/backup/result.rs index 2de5588..e005948 100644 --- a/src/services/backup/result.rs +++ b/src/services/backup/result.rs @@ -1,3 +1,4 @@ +use super::logger::JobLogEntry; use super::models::{BackupResult, UploadResult}; use super::service::BackupService; use crate::services::api::ApiError; @@ -11,6 +12,8 @@ impl BackupService { result: BackupResult, upload_results: Vec, backup_id: &String, + logs: Vec, + duration_ms: f64, ) -> Result, ApiError> { let status = if upload_results.iter().any(|r| r.success) { "success" @@ -32,6 +35,8 @@ impl BackupService { status, file_size, &result.generated_id, + logs, + duration_ms ) .await .map_err(|e| { diff --git a/src/services/backup/runner.rs b/src/services/backup/runner.rs index 73ff871..52aa5d5 100644 --- a/src/services/backup/runner.rs +++ b/src/services/backup/runner.rs @@ -1,3 +1,4 @@ +use super::logger::JobLogger; use super::models::BackupResult; use super::service::BackupService; @@ -6,10 +7,11 @@ use crate::services::config::DatabaseConfig; use anyhow::Result; use std::path::Path; -use tracing::{error, info}; +use std::sync::Arc; +use tracing::error; impl BackupService { - pub async fn run(cfg: DatabaseConfig, tmp_path: &Path) -> Result { + pub async fn run(cfg: DatabaseConfig, tmp_path: &Path, logger: Arc) -> Result { let db = DatabaseFactory::create_for_backup(cfg.clone()).await; let generated_id = cfg.generated_id.clone(); @@ -19,13 +21,15 @@ impl BackupService { Ok(v) => v, Err(e) => { error!("Ping failed: {}", e); + logger.log("error", format!("Ping failed: {}", e)); return Err(e.into()); } }; - info!("Reachable: {}", reachable); + logger.log("info", format!("Database reachable: {}", reachable)); if !reachable { + logger.log("error", "Database unreachable, backup aborted"); return Ok(BackupResult { generated_id, db_type, @@ -35,7 +39,7 @@ impl BackupService { }); } - match db.backup(tmp_path).await { + match db.backup(tmp_path, Arc::clone(&logger)).await { Ok(file) => Ok(BackupResult { generated_id, db_type, @@ -44,16 +48,19 @@ impl BackupService { code: None, }), - Err(e) if e.to_string() == "backup_already_in_progress" => Ok(BackupResult { - generated_id, - db_type, - status: "failed".into(), - backup_file: None, - code: Some("backup_already_in_progress".into()), - }), + Err(e) if e.to_string() == "backup_already_in_progress" => { + logger.log("warn", "Backup already in progress"); + Ok(BackupResult { + generated_id, + db_type, + status: "failed".into(), + backup_file: None, + code: Some("backup_already_in_progress".into()), + }) + } Err(e) => { - error!("Backup failed for {}: {:?}", generated_id, e); + logger.log("error", format!("Backup failed: {}", e)); Ok(BackupResult { generated_id, db_type, diff --git a/src/services/backup/uploader.rs b/src/services/backup/uploader.rs index 35ecdf5..116bf8b 100644 --- a/src/services/backup/uploader.rs +++ b/src/services/backup/uploader.rs @@ -1,13 +1,13 @@ +use super::logger::JobLogger; use super::models::{BackupResult, UploadResult}; use super::service::BackupService; - use crate::services::api::models::agent::status::DatabaseStorage; use crate::services::storage; use crate::utils::common::BackupMethod; - use anyhow::{Result, bail}; use futures::future::join_all; -use tracing::{error, info}; +use std::sync::Arc; +use tracing::info; impl BackupService { pub async fn upload( @@ -17,6 +17,7 @@ impl BackupService { storages: Vec, encrypt: bool, backup_id: &String, + logger: Arc, ) -> Result> { if result.code.as_deref() == Some("backup_already_in_progress") { info!("Skipping send: backup already in progress"); @@ -29,15 +30,13 @@ impl BackupService { let ctx_clone = ctx.clone(); let result_clone = result.clone(); let provider = storage::get_provider(&storage); + let logger_clone = Arc::clone(&logger); let storage_id = storage.id.clone(); let generated_id = result_clone.generated_id.clone(); async move { - info!( - "Uploading storage -> {:?} for {:?}", - storage.provider, storage_id - ); + logger_clone.log("info", format!("Uploading storage {:?} (id: {})", storage.provider, storage_id)); /* INIT STEP @@ -54,7 +53,7 @@ impl BackupService { { Ok(v) => v, Err(e) => { - error!("backup_upload_init failed: {}", e); + logger_clone.log("error", format!("Upload init failed: {}", e)); return UploadResult { storage_id, @@ -69,6 +68,7 @@ impl BackupService { let backup_storage_id = match init { Some(v) => v.backup_storage.id, None => { + logger_clone.log("error", "Upload init returned empty response"); return UploadResult { storage_id, success: false, @@ -83,7 +83,7 @@ impl BackupService { PROVIDER CHECK */ let Some(provider) = provider else { - error!("Skipping storage due to missing provider"); + logger_clone.log("error", format!("Missing provider for storage {}", storage_id)); return UploadResult { storage_id, @@ -107,20 +107,23 @@ impl BackupService { ) .await; - let status = if upload_result.success { - "success" - } else { - "failed" - }; + let status = if upload_result.success { "success" } else { "failed" }; if status != "success" { + logger_clone.log("error", format!( + "Upload failed for storage {}: {}", + storage_id, + upload_result.error.as_deref().unwrap_or("unknown error") + )); return upload_result; } - info!( - "Storage {} uploaded to remote path {:?}", - storage_id, upload_result.remote_file_path - ); + logger_clone.log("info", format!( + "Storage {} uploaded to {:?} ({} bytes)", + storage_id, + upload_result.remote_file_path.clone().unwrap().to_string(), + upload_result.total_size.unwrap_or(0) + )); /* METADATA VALIDATION @@ -129,6 +132,7 @@ impl BackupService { match (&upload_result.remote_file_path, upload_result.total_size) { (Some(path), Some(size)) => (path.clone(), size), _ => { + logger_clone.log("error", format!("Missing remote_file_path or total_size for storage {}", storage_id)); return UploadResult { storage_id, success: false, @@ -158,10 +162,7 @@ impl BackupService { Ok(_) => upload_result, Err(err) => { - error!( - "backup_upload_status failed (storage_id={}): {}", - storage_id, err - ); + logger_clone.log("error", format!("Upload status update failed for {}: {}", storage_id, err)); UploadResult { storage_id, diff --git a/src/services/restore/archive.rs b/src/services/restore/archive.rs index be87159..237f8de 100644 --- a/src/services/restore/archive.rs +++ b/src/services/restore/archive.rs @@ -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 ) -> Result { + 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) } } diff --git a/src/services/restore/downloader.rs b/src/services/restore/downloader.rs index 0b6cf41..ef10394 100644 --- a/src/services/restore/downloader.rs +++ b/src/services/restore/downloader.rs @@ -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 { + pub async fn download_backup(&self, file_url: &str, tmp_path: &Path, logger: Arc) -> Result { + 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) } } diff --git a/src/services/restore/executor.rs b/src/services/restore/executor.rs index b790623..7f1cff1 100644 --- a/src/services/restore/executor.rs +++ b/src/services/restore/executor.rs @@ -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(()) } diff --git a/src/services/restore/result.rs b/src/services/restore/result.rs index d3f61c1..c4f92a0 100644 --- a/src/services/restore/result.rs +++ b/src/services/restore/result.rs @@ -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, + duration_ms: f64, + ) -> Result, 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() + }) } -} +} \ No newline at end of file diff --git a/src/services/restore/runner.rs b/src/services/restore/runner.rs index 3806865..756f72f 100644 --- a/src/services/restore/runner.rs +++ b/src/services/restore/runner.rs @@ -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, ) -> Result { 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, diff --git a/src/tests/domain/firebird.rs b/src/tests/domain/firebird.rs index a465081..228dfdf 100644 --- a/src/tests/domain/firebird.rs +++ b/src/tests/domain/firebird.rs @@ -7,18 +7,12 @@ use std::time::Duration; use tempfile::TempDir; use testcontainers::runners::AsyncRunner; use testcontainers::{ContainerAsync, GenericImage, ImageExt}; -use testcontainers::core::{AccessMode, IntoContainerPort}; +use testcontainers::core::IntoContainerPort; use tracing::{error, info}; -use testcontainers::core::{Mount}; async fn create_config() -> (ContainerAsync, DatabaseConfig) { - - let mount = Mount::volume_mount("firebird-test-data", "/var/lib/firebird/data") - .with_access_mode(AccessMode::ReadWrite); - let container = GenericImage::new("firebirdsql/firebird", "latest") .with_exposed_port(3050.tcp()) - .with_mount(mount) .with_env_var("FIREBIRD_ROOT_PASSWORD", "fake_root_password") .with_env_var("FIREBIRD_USER", "alice") .with_env_var("FIREBIRD_PASSWORD", "fake_password") @@ -73,13 +67,13 @@ async fn firebird_backup_restore_test() { let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); tokio::time::sleep(Duration::from_secs(10)).await; assert!(file_path.is_file()); - let compression = compress_to_tar_gz_large(&file_path).await.unwrap(); + let compression = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(compression.compressed_path.is_file()); let files = decompress_large_tar_gz( @@ -101,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) @@ -111,4 +105,4 @@ async fn firebird_backup_restore_test() { assert!(false) } } -} \ No newline at end of file +} diff --git a/src/tests/domain/mariadb.rs b/src/tests/domain/mariadb.rs index 6538151..836d211 100644 --- a/src/tests/domain/mariadb.rs +++ b/src/tests/domain/mariadb.rs @@ -58,11 +58,11 @@ async fn mariadb_backup_restore_test() { let backup_path = temp_dir.path(); let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file()); - let compression = compress_to_tar_gz_large(&file_path).await.unwrap(); + let compression = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(compression.compressed_path.is_file()); let files = decompress_large_tar_gz(compression.compressed_path.as_path(), temp_dir.path()) @@ -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) diff --git a/src/tests/domain/mongodb.rs b/src/tests/domain/mongodb.rs index 1b86ae2..e596435 100644 --- a/src/tests/domain/mongodb.rs +++ b/src/tests/domain/mongodb.rs @@ -75,7 +75,7 @@ async fn mongodb_backup_restore_test() { let backup_path = temp_dir.path(); let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); info!("Backup path: {:?}", file_path); assert!(file_path.is_file()); @@ -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) diff --git a/src/tests/domain/mssql.rs b/src/tests/domain/mssql.rs index 61781be..520d854 100644 --- a/src/tests/domain/mssql.rs +++ b/src/tests/domain/mssql.rs @@ -43,7 +43,7 @@ async fn create_user_database(host: &str, port: u16, db_name: &str) { client.simple_query(sql.as_str()).await.unwrap(); } -fn make_config(host: String, port: u16, database: &str) -> DatabaseConfig { +fn make_config(host: String, port: u16, database: &str, generated_id: &str) -> DatabaseConfig { DatabaseConfig { name: "Test MSSQL".to_string(), database: database.to_string(), @@ -52,7 +52,7 @@ fn make_config(host: String, port: u16, database: &str) -> DatabaseConfig { password: SA_PASSWORD.to_string(), port, host, - generated_id: "5a445eb4-c2c6-4bde-a423-ee1385dcf6d3".to_string(), + generated_id: generated_id.to_string(), path: "".to_string(), } } @@ -66,7 +66,7 @@ async fn mssql_ping_test() { let host = container.get_host().await.unwrap().to_string(); let port = container.get_host_port_ipv4(1433).await.unwrap(); - let config = make_config(host, port, "master"); + let config = make_config(host, port, "master", "5a445eb4-c2c6-4bde-a423-ee1385dcf6d3"); let db = DatabaseFactory::create_for_backup(config).await; let reachable = db.ping().await.unwrap_or(false); @@ -86,11 +86,11 @@ async fn mssql_backup_test() { create_user_database(&host, port, "backupdb").await; - let config = make_config(host, port, "backupdb"); + let config = make_config(host, port, "backupdb", "5a445eb4-c2c6-4bde-a423-ee1385dcf6d4"); let temp_dir = TempDir::new().unwrap(); let db = DatabaseFactory::create_for_backup(config).await; - let file_path = db.backup(temp_dir.path()).await.unwrap(); + let file_path = db.backup(temp_dir.path(), std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file(), "backup file should exist"); assert!( @@ -115,14 +115,14 @@ async fn mssql_backup_restore_test() { create_user_database(&host, port, "sourcedb").await; - let backup_config = make_config(host.clone(), port, "sourcedb"); + let backup_config = make_config(host.clone(), port, "sourcedb", "5a445eb4-c2c6-4bde-a423-ee1385dcf6d5"); let temp_dir = TempDir::new().unwrap(); let db = DatabaseFactory::create_for_backup(backup_config).await; - let file_path = db.backup(temp_dir.path()).await.unwrap(); + let file_path = db.backup(temp_dir.path(), std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file()); - let compression = compress_to_tar_gz_large(&file_path).await.unwrap(); + let compression = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(compression.compressed_path.is_file()); let files = decompress_large_tar_gz( @@ -138,10 +138,10 @@ async fn mssql_backup_restore_test() { panic!("Unexpected number of files after decompression: {}", files.len()); }; - let restore_config = make_config(host, port, "restoreddb"); + 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); diff --git a/src/tests/domain/mysql.rs b/src/tests/domain/mysql.rs index 20b8861..f912d64 100644 --- a/src/tests/domain/mysql.rs +++ b/src/tests/domain/mysql.rs @@ -58,11 +58,11 @@ async fn mysql_backup_restore_test() { let backup_path = temp_dir.path(); let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file()); - let compression = compress_to_tar_gz_large(&file_path).await.unwrap(); + let compression = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(compression.compressed_path.is_file()); let files = decompress_large_tar_gz(compression.compressed_path.as_path(), temp_dir.path()) @@ -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) diff --git a/src/tests/domain/postgres.rs b/src/tests/domain/postgres.rs index bac4f6d..029ef84 100644 --- a/src/tests/domain/postgres.rs +++ b/src/tests/domain/postgres.rs @@ -66,11 +66,11 @@ async fn postgres_backup_restore_test() { let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file()); - let compression = compress_to_tar_gz_large(&file_path).await.unwrap(); + let compression = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(compression.compressed_path.is_file()); @@ -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) diff --git a/src/tests/domain/redis.rs b/src/tests/domain/redis.rs index 32e3ebd..ef56dac 100644 --- a/src/tests/domain/redis.rs +++ b/src/tests/domain/redis.rs @@ -56,7 +56,7 @@ async fn redis_backup_test() { let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file()); } diff --git a/src/tests/domain/valkey.rs b/src/tests/domain/valkey.rs index f018541..42c6df6 100644 --- a/src/tests/domain/valkey.rs +++ b/src/tests/domain/valkey.rs @@ -55,7 +55,7 @@ async fn valkey_backup_test() { let db = DatabaseFactory::create_for_backup(config.clone()).await; - let file_path = db.backup(backup_path).await.unwrap(); + let file_path = db.backup(backup_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await.unwrap(); assert!(file_path.is_file()); } diff --git a/src/tests/utils/compress_tests.rs b/src/tests/utils/compress_tests.rs index b735d95..7e96c69 100644 --- a/src/tests/utils/compress_tests.rs +++ b/src/tests/utils/compress_tests.rs @@ -9,7 +9,7 @@ async fn compress_creates_tar_gz() -> Result<()> { let file_path = tmp.path().join("test.txt"); write(&file_path, b"hello world").await?; - let result = compress_to_tar_gz_large(&file_path).await?; + let result = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await?; assert!(result.compressed_path.exists()); assert_eq!(result.compressed_path.extension().unwrap(), "gz"); @@ -22,8 +22,7 @@ async fn compress_skips_existing_tar_gz() -> Result<()> { let file_path = tmp.path().join("already.tar.gz"); write(&file_path, b"compressed").await?; - let result = compress_to_tar_gz_large(&file_path).await?; - // Should return same path without creating a new file + let result = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await?; assert_eq!(result.compressed_path, file_path); Ok(()) @@ -35,7 +34,7 @@ async fn decompress_restores_file() -> Result<()> { let file_path = tmp.path().join("file.txt"); write(&file_path, b"data for decompress").await?; - let compress_result = compress_to_tar_gz_large(&file_path).await?; + let compress_result = compress_to_tar_gz_large(&file_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await?; let output_dir = tmp.path().join("out"); tokio::fs::create_dir_all(&output_dir).await?; @@ -58,8 +57,8 @@ async fn decompress_multiple_files() -> Result<()> { write(&file2, b"file2").await?; // Compress both files individually (for simplicity in this test) - let compress1 = compress_to_tar_gz_large(&file1).await?; - let compress2 = compress_to_tar_gz_large(&file2).await?; + let compress1 = compress_to_tar_gz_large(&file1, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await?; + let compress2 = compress_to_tar_gz_large(&file2, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await?; let output_dir = tmp.path().join("out_multi"); tokio::fs::create_dir_all(&output_dir).await?; diff --git a/src/utils/compress.rs b/src/utils/compress.rs index 64bb945..4a84a07 100644 --- a/src/utils/compress.rs +++ b/src/utils/compress.rs @@ -3,6 +3,7 @@ use async_compression::tokio::bufread::GzipDecoder; use async_compression::tokio::write::GzipEncoder as AsyncGzipEncoder; use futures::StreamExt; use std::path::{Path, PathBuf}; +use std::sync::Arc; use tokio::fs::File; use tokio::fs::create_dir_all; use tokio::io::AsyncWriteExt; @@ -11,19 +12,21 @@ use tokio_tar::Archive; use tokio_tar::Builder as TokioTarBuilder; use tracing::info; +use crate::services::backup::logger::JobLogger; + #[allow(dead_code)] pub struct CompressionResult { pub compressed_path: PathBuf, } -pub async fn compress_to_tar_gz_large(file: &PathBuf) -> Result { +pub async fn compress_to_tar_gz_large(file: &PathBuf, logger: Arc) -> Result { if file .file_name() .and_then(|n| n.to_str()) .map(|n| n.ends_with(".tar.gz")) .unwrap_or(false) { - info!("File {:?} is already a tar.gz, skipping compression", file); + logger.log("info", format!("File {:?} is already a tar.gz, skipping compression", file)); return Ok(CompressionResult { compressed_path: file.clone(), }); @@ -31,34 +34,56 @@ pub async fn compress_to_tar_gz_large(file: &PathBuf) -> Result